diff --git a/PROJECT_MEMORY.md b/PROJECT_MEMORY.md index 5c32fcc..5011015 100644 --- a/PROJECT_MEMORY.md +++ b/PROJECT_MEMORY.md @@ -102,7 +102,14 @@ Diese Datei hält den verifizierten technischen Arbeitsstand fest. Sie enthält - Das Hauptfenster ist unter Windows nun ausdrücklich Besitzer des App-Lebenszyklus: Nach einem natürlichen, vom Renderer bestätigten Schließen wird der kontrollierte Shutdown auch bei verbleibenden Providerfenstern gestartet. Der normale `beforeunload`-Pfad bleibt erhalten, damit noch nicht übertragene Linksammler-Änderungen synchron gesichert werden können. - `second-instance` und `activate` stellen ein vorhandenes Hauptfenster wieder her oder erzeugen ein fehlendes neu. Ein Start während eines bereits laufenden Shutdowns plant genau einen Relaunch nach dem Prozessende ein. - Der Shutdown besitzt einen äußeren 10-Sekunden-Wächter. Nach erfolgreichem Controller-Shutdown wird weiterhin normal mit `app.quit()` beendet; blockieren Renderer- oder Providerfenster diesen letzten Schritt, erzwingt ein separater 2-Sekunden-Wächter den Exit. `shutdownDaemon()` läuft im Graceful-Pfad erst nach `controller.shutdown()`. -- Die Änderung ist nur auf dem Arbeitsbranch implementiert und getestet. Sie wurde weder veröffentlicht noch auf dem Server installiert; dafür ist eine neue ausdrückliche Release-/Deployment-Freigabe erforderlich. +- Der Debug-Server serialisiert Start, Stop und Restart. Gleichzeitige Restarts können keinen unverwalteten Listener mehr hinterlassen, wartende Restarts können einen Stop nicht mehr überholen und ein Stop während des Socket-Starts löst die Lifecycle-Queue zuverlässig auf. +- Settings- und Session-Recovery unterscheiden jetzt gültige leere Zustände von JSON-gültigen, aber unbrauchbaren Strukturen. Eine absichtlich leere Session bleibt maßgeblich; beschädigte Settings-/Session-Primärdateien fallen auf ein intaktes Backup zurück. Sparse historische Settings bleiben kompatibel, und eine fehlende Settings-Primärdatei wird aus dem Backup repariert. +- Das sichere Crash-Journal für ausstehende Hoster-Metadaten-Umbenennungen bleibt über einen Neustart erhalten. Nur absolute echte Unterpfade des Paketausgabeordners werden übernommen; fremde, relative oder mit dem Ausgabeordner identische Pfade werden verworfen. +- Asynchrone Settings- und Session-Saves lösen ihre Promises erst nach dem zugehörigen dauerhaften Schreibvorgang auf und propagieren Fehler. Coalescing, Drain-Ende, Importbarriere und Replay warten auch bei Fehlern auf alle gestarteten Writes; verlorene Wakeups, hängende Barrieren, still verworfene Saves und falsch erfolgreiche Waiter wurden mit Race-Regressionen abgedeckt. +- Der Updater beendet Backpressure-, Shutdown- und Idle-Wartezustände auch bei Schreibfehlern zuverlässig. Er wartet vor dem Rename auf den tatsächlichen Dateistream-Close, propagiert späte Close-Fehler, hält die Fehlerabsicherung bis `close` aktiv und blockiert bei einem nie auflösenden Response-`cancel()` nicht mehr dauerhaft. +- Die übergroßen Updater-Test-Fixtures wurden durch gleichwertige 128-KiB-MZ-Payloads ersetzt, damit Integritäts- und Fallback-Tests unter paralleler Last nicht an Vitests 5-Sekunden-Grenze flaken. +- Die Entpacklogik und ihre Produktionspfade wurden nicht verändert. +- Die Änderungen sind nur auf dem Arbeitsbranch implementiert und getestet. Sie wurden weder veröffentlicht noch auf dem Server installiert; dafür ist eine neue ausdrückliche Release-/Deployment-Freigabe erforderlich. ## Start-, Build- und Testbefehle @@ -152,6 +159,10 @@ npm exec -- tsc --noEmit - Vollständiger Client-Lauf ohne den privilegierten Symlink-Metadatentest: 138 Testdateien erfolgreich, 1 JVM-Testdatei übersprungen; 2.609 Tests erfolgreich und 4 übersprungen. Ein unter paralleler Build-Last einmal auffälliger, fachfremder Ordner-Aufräumtest war anschließend dreimal isoliert und im vollständigen Wiederholungslauf erfolgreich. - `public-release-metadata.test.ts` separat: 22 von 24 erfolgreich. Die zwei übrigen Fälle scheitern weiterhin ausschließlich beim Anlegen ihrer Symlink-Fixtures mit Windows-`EPERM`, bevor die Produktassertion ausgeführt wird. - Nach dem Lifecycle-Fix erneut erfolgreich: TypeScript, vollständiger Main-/Renderer-Build, Self-Check und 16 von 16 Backup-API-Tests. Die bekannte Vite-Warnung zum rund 574 KiB großen Renderer-Chunk bleibt bestehen. +- Gezielte finale Regressionen: Storage und Session-Restart 141 von 141, Updater 40 von 40 und Debug-Server 18 von 18 erfolgreich. +- Finaler vollständiger Client-Lauf ohne den privilegierten Symlink-Metadatentest: 138 Testdateien erfolgreich, 1 JVM-Testdatei übersprungen; 2.640 Tests erfolgreich und 4 übersprungen. +- Der erste vollständige Lauf dieser Änderung hatte bei 2.639 erfolgreichen Tests genau einen 5-Sekunden-Timeout in einer Updater-Testfixture mit der 88,5-MiB-Node-EXE. Nach dem Ersatz aller reinen MZ-Fixtures durch 128-KiB-Payloads war der vollständige Wiederholungslauf grün. +- Nach allen finalen Korrekturen erneut erfolgreich: TypeScript, vollständiger Main-/Renderer-Build, Self-Check und 16 von 16 Backup-API-Tests. Die bekannte Vite-Warnung zum rund 574 KiB großen Renderer-Chunk bleibt bestehen. ## Bekannte Probleme und Risiken @@ -165,11 +176,11 @@ npm exec -- tsc --noEmit - Der Renderer-Bundle-Chunk liegt über Vites 500-KiB-Warnschwelle. - Die Shutdown-Wächter können einen vollständig synchron blockierten Main-Thread nicht präemptieren. Im äußersten Fall beträgt die kombinierte Grenze knapp 12 Sekunden: bis zu 10 Sekunden Controller-Shutdown plus 2 Sekunden Quit-Bestätigung. - Beim erzwungenen Exit können letzte asynchrone Logzeilen oder Benachrichtigungen fehlen; die wesentlichen Queue-, Settings-, Statistik- und Collector-Daten werden zuvor synchron gesichert. Laufende externe Extraktions-/Remux-Prozesse bleiben ein separater späterer Härtungspunkt. -- Der Start-/Beenden-Fix ist noch nicht veröffentlicht oder auf dem Server installiert. Bis zu einer ausdrücklich freigegebenen Auslieferung läuft dort weiterhin die bisherige Version. +- Die aktuellen Lifecycle-, Debug-Server-, Persistenz- und Updater-Fixes sind noch nicht veröffentlicht oder auf dem Server installiert. Bis zu einer ausdrücklich freigegebenen Auslieferung läuft dort weiterhin die bisherige Version. ## Nächste sinnvolle Schritte -1. Für eine Auslieferung des Start-/Beenden-Fixes zuerst Stand, Ziel, Tests und Rollback-Weg nennen und Saschas ausdrückliches Release-/Deployment-Go abwarten. +1. Für eine Auslieferung der aktuellen Fehlerkorrekturen zuerst Stand, Ziel, Tests und Rollback-Weg nennen und Saschas ausdrückliches Release-/Deployment-Go abwarten. 2. Nach einer freigegebenen Serverinstallation den Ablauf Real-Debrid-Web-Download, Hauptfenster schließen und unmittelbar neu starten am echten Zielsystem verifizieren; vor jedem Eingriff Prozessbaum und Logtail sichern. 3. Sicherheitsabhängigkeiten in einem separaten Upgrade-Branch aktualisieren, Electron-/Vite-/Vitest-Major-Wechsel einzeln testen und danach den vollständigen Windows-Paketpfad prüfen. 4. Eine nichtdestruktive Strategie zur Bereinigung der divergierenden `main`-Branches abstimmen; kein Force-Push ohne ausdrückliche Freigabe. diff --git a/src/main/debug-server.ts b/src/main/debug-server.ts index e04a8b1..bebe0c4 100644 --- a/src/main/debug-server.ts +++ b/src/main/debug-server.ts @@ -76,6 +76,9 @@ let runtimeBaseDir = ""; let allowlist: string[] = []; let requestLimits = new Map(); let notificationStatusProvider: (() => NotificationSupportPayload) | null = null; +let serverLifecycleQueue: Promise = Promise.resolve(); +let serverLifecycleGeneration = 0; +let debugServerEnabled = false; export interface DebugServerRuntimeStatus { running: boolean; @@ -1171,13 +1174,19 @@ function openServerSocket(): Promise { } const srv = http.createServer(handleRequest); let settled = false; - const settle = (): void => { - if (!settled) { - settled = true; - resolve(); - } - }; - srv.on("error", (err: NodeJS.ErrnoException) => { + const settle = (): void => { + if (!settled) { + settled = true; + resolve(); + } + }; + srv.once("close", () => { + if (server === srv) { + server = null; + } + settle(); + }); + srv.on("error", (err: NodeJS.ErrnoException) => { if (err.code === "EADDRINUSE") { logger.warn(`Debug-Server: Port ${bindPort} belegt (EADDRINUSE) - Server nicht gestartet`); } else { @@ -1193,9 +1202,49 @@ function openServerSocket(): Promise { settle(); }); server = srv; - }); -} - + }); +} + +function enqueueServerLifecycle(operation: () => Promise): Promise { + const pending = serverLifecycleQueue.then(operation, operation); + serverLifecycleQueue = pending.catch(() => undefined); + return pending; +} + +function closeServerSocket(): Promise { + const current = server; + if (!current) { + return Promise.resolve(); + } + server = null; + return new Promise((resolve) => { + let settled = false; + let timeout: NodeJS.Timeout | null = null; + const done = (): void => { + if (settled) { + return; + } + settled = true; + if (timeout) { + clearTimeout(timeout); + } + resolve(); + }; + try { + current.close(() => done()); + } catch { + done(); + } + try { + current.closeAllConnections?.(); + } catch { + } + if (!settled) { + timeout = setTimeout(done, 1500); + } + }); +} + export function startDebugServer( mgr: DownloadManager, baseDir: string, @@ -1204,32 +1253,34 @@ export function startDebugServer( runtimeBaseDir = baseDir; manager = mgr; notificationStatusProvider = readCurrentNotificationStatus || null; - void openServerSocket(); + debugServerEnabled = true; + const generation = ++serverLifecycleGeneration; + void enqueueServerLifecycle(async () => { + if (!debugServerEnabled || generation !== serverLifecycleGeneration) { + return; + } + await closeServerSocket(); + if (!debugServerEnabled || generation !== serverLifecycleGeneration) { + return; + } + await openServerSocket(); + }); +} + +export async function restartDebugServer(): Promise { + const generation = serverLifecycleGeneration; + await enqueueServerLifecycle(async () => { + if (!debugServerEnabled || generation !== serverLifecycleGeneration) { + return; + } + await closeServerSocket(); + if (!debugServerEnabled || generation !== serverLifecycleGeneration) { + return; + } + await openServerSocket(); + }); + return getDebugServerRuntimeStatus(); } - -export async function restartDebugServer(): Promise { - const old = server; - if (old) { - server = null; - await new Promise((resolve) => { - let settled = false; - const done = (): void => { - if (!settled) { - settled = true; - resolve(); - } - }; - old.close(() => done()); - try { - old.closeAllConnections?.(); - } catch { - } - setTimeout(done, 1500); - }); - } - await openServerSocket(); - return getDebugServerRuntimeStatus(); -} export function getDebugServerRuntimeStatus(): DebugServerRuntimeStatus { return { @@ -1273,15 +1324,21 @@ export function clearDebugToken(): void { } export function stopDebugServer(): void { + debugServerEnabled = false; + serverLifecycleGeneration += 1; notificationStatusProvider = null; - if (server) { - server.close(); - try { - server.closeAllConnections?.(); - } catch { - } - server = null; - requestLimits = new Map(); + const current = server; + server = null; + requestLimits = new Map(); + if (current) { + try { + current.close(); + } catch { + } + try { + current.closeAllConnections?.(); + } catch { + } logger.info("Debug-Server gestoppt"); } } diff --git a/src/main/storage.ts b/src/main/storage.ts index a6124cb..dab96e5 100644 --- a/src/main/storage.ts +++ b/src/main/storage.ts @@ -75,16 +75,30 @@ function isPathInsideDir(filePath: string, dirPath: string): boolean { } } -function normalizeSessionTargetPath(value: unknown, packageOutputDir: string): string { - const targetPath = asText(value); - if (!targetPath || !packageOutputDir || !path.isAbsolute(targetPath)) { +function normalizeSessionTargetPath(value: unknown, packageOutputDir: string): string { + const targetPath = asText(value); + if (!targetPath || !packageOutputDir || !path.isAbsolute(targetPath)) { return ""; } if (!isPathInsideDir(targetPath, packageOutputDir)) { return ""; } - return path.resolve(targetPath); -} + return path.resolve(targetPath); +} + +function normalizeSessionMetadataRenameTargetPath(value: unknown, packageOutputDir: string): string { + if (!path.isAbsolute(packageOutputDir)) { + return ""; + } + const targetPath = normalizeSessionTargetPath(value, packageOutputDir); + if (!targetPath) { + return ""; + } + const resolvedOutputDir = path.resolve(packageOutputDir); + const normalizedTargetPath = process.platform === "win32" ? targetPath.toLowerCase() : targetPath; + const normalizedOutputDir = process.platform === "win32" ? resolvedOutputDir.toLowerCase() : resolvedOutputDir; + return normalizedTargetPath === normalizedOutputDir ? "" : targetPath; +} function clampNumber(value: unknown, fallback: number, min: number, max: number): number { const num = Number(value); @@ -907,9 +921,16 @@ function migrateLegacyMegaEnableFlags(parsed: AppSettings): AppSettings { return parsed; } const preferApi = parsed.megaDebridPreferApi !== undefined ? Boolean(parsed.megaDebridPreferApi) : true; - return { ...parsed, megaDebridApiEnabled: preferApi, megaDebridWebEnabled: !preferApi }; -} - + return { ...parsed, megaDebridApiEnabled: preferApi, megaDebridWebEnabled: !preferApi }; +} + +const PERSISTED_SETTINGS_KEYS = new Set(Object.keys(defaultSettings())); + +function isPersistedSettingsEnvelope(value: unknown): boolean { + const parsed = asRecord(value); + return Boolean(parsed && Object.keys(parsed).some((key) => PERSISTED_SETTINGS_KEYS.has(key))); +} + interface LoadedSettingsFile { settings: AppSettings; needsCredentialRewrite: boolean; @@ -917,7 +938,12 @@ interface LoadedSettingsFile { function readSettingsFile(filePath: string): LoadedSettingsFile | null { try { - const parsed = JSON.parse(fs.readFileSync(filePath, "utf8")) as AppSettings; + const raw = JSON.parse(fs.readFileSync(filePath, "utf8")) as unknown; + if (!isPersistedSettingsEnvelope(raw)) { + logger.error(`Settings-Datei beschädigt (Struktur ungültig): ${filePath}`); + return null; + } + const parsed = raw as AppSettings; const needsCredentialRewrite = needsPersistedSettingsRewrite(parsed); const restored = restorePersistedSettings(parsed); const migratedLanguage = (restored as Partial).language === undefined ? "de" : restored.language; @@ -1001,10 +1027,11 @@ export function normalizeLoadedSession(raw: unknown): SessionState { speedBps: clampNumber(item.speedBps, 0, 0, 10_000_000_000), downloadedBytes: clampNumber(item.downloadedBytes, 0, 0, 10_000_000_000_000), totalBytes: item.totalBytes == null ? null : clampNumber(item.totalBytes, 0, 0, 10_000_000_000_000), - progressPercent: clampNumber(item.progressPercent, 0, 0, 100), - fileName: asText(item.fileName) || "download.bin", - targetPath: asText(item.targetPath), - resumable: item.resumable === undefined ? true : Boolean(item.resumable), + progressPercent: clampNumber(item.progressPercent, 0, 0, 100), + fileName: asText(item.fileName) || "download.bin", + targetPath: asText(item.targetPath), + metadataRenameTargetPath: asText(item.metadataRenameTargetPath) || undefined, + resumable: item.resumable === undefined ? true : Boolean(item.resumable), attempts: clampNumber(item.attempts, 0, 0, 10_000), lastError: asText(item.lastError), fullStatus: asText(item.fullStatus), @@ -1083,21 +1110,35 @@ export function normalizeLoadedSession(raw: unknown): SessionState { logger.warn(`normalizeLoadedSession: ${orphanedItemCount} verwaiste Items entfernt (fehlende Pakete)`); } - let droppedUnsafeTargetPathCount = 0; - for (const item of Object.values(itemsById)) { - const pkg = packagesById[item.packageId]; - if (!pkg) { + let droppedUnsafeTargetPathCount = 0; + let droppedUnsafeMetadataRenameTargetPathCount = 0; + for (const item of Object.values(itemsById)) { + const pkg = packagesById[item.packageId]; + if (!pkg) { continue; } const safeTargetPath = normalizeSessionTargetPath(item.targetPath, pkg.outputDir); if (!safeTargetPath && asText(item.targetPath)) { - droppedUnsafeTargetPathCount += 1; - } - item.targetPath = safeTargetPath; - } - if (droppedUnsafeTargetPathCount > 0) { - logger.warn(`normalizeLoadedSession: ${droppedUnsafeTargetPathCount} unsichere targetPath-Eintraege verworfen`); - } + droppedUnsafeTargetPathCount += 1; + } + item.targetPath = safeTargetPath; + const metadataRenameTargetPath = asText(item.metadataRenameTargetPath); + const safeMetadataRenameTargetPath = normalizeSessionMetadataRenameTargetPath(metadataRenameTargetPath, pkg.outputDir); + if (metadataRenameTargetPath && !safeMetadataRenameTargetPath) { + droppedUnsafeMetadataRenameTargetPathCount += 1; + } + if (safeMetadataRenameTargetPath) { + item.metadataRenameTargetPath = safeMetadataRenameTargetPath; + } else { + delete item.metadataRenameTargetPath; + } + } + if (droppedUnsafeTargetPathCount > 0) { + logger.warn(`normalizeLoadedSession: ${droppedUnsafeTargetPathCount} unsichere targetPath-Eintraege verworfen`); + } + if (droppedUnsafeMetadataRenameTargetPathCount > 0) { + logger.warn(`normalizeLoadedSession: ${droppedUnsafeMetadataRenameTargetPathCount} unsichere metadataRenameTargetPath-Einträge verworfen`); + } for (const pkg of Object.values(packagesById)) { pkg.itemIds = pkg.itemIds.filter((itemId) => { @@ -1142,14 +1183,13 @@ export function normalizeLoadedSession(raw: unknown): SessionState { } export function loadSettings(paths: StoragePaths): AppSettings { - ensureBaseDir(paths.baseDir); - if (!fs.existsSync(paths.configFile)) { - return defaultSettings(); - } - const loaded = readSettingsFile(paths.configFile); + ensureBaseDir(paths.baseDir); + const primaryExists = fs.existsSync(paths.configFile); + const backupFile = `${paths.configFile}.bak`; + const backupExists = fs.existsSync(backupFile); + const loaded = primaryExists ? readSettingsFile(paths.configFile) : null; if (loaded) { - const backupFile = `${paths.configFile}.bak`; - const backupNeedsCredentialRewrite = fs.existsSync(backupFile) + const backupNeedsCredentialRewrite = backupExists ? needsSettingsFileCredentialRewrite(backupFile) : false; if (loaded.needsCredentialRewrite || backupNeedsCredentialRewrite) { @@ -1158,17 +1198,20 @@ export function loadSettings(paths: StoragePaths): AppSettings { return loaded.settings; } - const backupFile = `${paths.configFile}.bak`; - const backupLoaded = fs.existsSync(backupFile) ? readSettingsFile(backupFile) : null; + const backupLoaded = backupExists ? readSettingsFile(backupFile) : null; if (backupLoaded) { - logger.warn("Konfiguration defekt, Backup-Datei wird verwendet"); + logger.warn(primaryExists + ? "Konfiguration defekt, Backup-Datei wird verwendet" + : "Konfiguration fehlt, Backup-Datei wird verwendet"); rewriteProtectedSettings(paths, backupLoaded.settings); return backupLoaded.settings; - } - - logger.error("Konfiguration konnte nicht geladen werden (auch Backup fehlgeschlagen)"); - return defaultSettings(); -} + } + + if (primaryExists || backupExists) { + logger.error("Konfiguration konnte nicht geladen werden (auch Backup fehlgeschlagen)"); + } + return defaultSettings(); +} function syncRenameWithExdevFallback(tempPath: string, targetPath: string): void { try { @@ -1282,6 +1325,16 @@ interface LoadedSessionFile { wasRunning: boolean; } +function isPersistedSessionEnvelope(value: unknown): boolean { + const parsed = asRecord(value); + return Boolean( + parsed + && Array.isArray(parsed.packageOrder) + && asRecord(parsed.packages) + && asRecord(parsed.items) + ); +} + function readSessionFile(filePath: string): LoadedSessionFile | null { let raw: string | null = null; const maxAttempts = 5; @@ -1307,9 +1360,13 @@ function readSessionFile(filePath: string): LoadedSessionFile | null { } if (raw === null) { return null; - } + } try { const parsed = JSON.parse(raw) as unknown; + if (!isPersistedSessionEnvelope(parsed)) { + logger.error(`Session-Datei beschädigt (Struktur ungültig): ${filePath}`); + return null; + } const normalized = normalizeLoadedSession(parsed); const wasRunning = normalized.running; const session = normalizeLoadedSessionTransientFields(normalized); @@ -1331,8 +1388,19 @@ export function saveSettings(paths: StoragePaths, settings: AppSettings): void { writeSettingsFileAtomically(paths.configFile, payloads.primary); } -let asyncSettingsSaveRunning = false; -let asyncSettingsSaveQueued: { paths: StoragePaths; settings: AppSettings; generation: number } | null = null; +interface AsyncSaveWaiter { + resolve: () => void; + reject: (error: unknown) => void; +} + +interface QueuedSettingsSave { + paths: StoragePaths; + settings: AppSettings; + generation: number; + waiters: AsyncSaveWaiter[]; +} + +let asyncSettingsSaveQueued: QueuedSettingsSave | null = null; let syncSettingsSaveGeneration = 0; let activeSettingsSave: Promise | null = null; @@ -1364,27 +1432,52 @@ async function writeSettingsPayload(paths: StoragePaths, settings: AppSettings, } } -async function saveSettingsPayloadAsync(paths: StoragePaths, settings: AppSettings, generation: number): Promise { - if (asyncSettingsSaveRunning) { - asyncSettingsSaveQueued = { paths, settings, generation }; +function startSettingsSaveDrain(): void { + if (activeSettingsSave) { return; } - asyncSettingsSaveRunning = true; - const operation = writeSettingsPayload(paths, settings, generation).catch((error) => { - logger.error(`Async Settings-Save fehlgeschlagen: ${String(error)}`); - }).finally(() => { - asyncSettingsSaveRunning = false; + let operation: Promise; + operation = (async () => { + while (asyncSettingsSaveQueued) { + const current = asyncSettingsSaveQueued; + asyncSettingsSaveQueued = null; + try { + await writeSettingsPayload(current.paths, current.settings, current.generation); + current.waiters.forEach((waiter) => waiter.resolve()); + } catch (error) { + current.waiters.forEach((waiter) => waiter.reject(error)); + const pending = asyncSettingsSaveQueued as QueuedSettingsSave | null; + asyncSettingsSaveQueued = null; + pending?.waiters.forEach((waiter) => waiter.reject(error)); + logger.error(`Async Settings-Save fehlgeschlagen: ${String(error)}`); + throw error; + } + } + })().finally(() => { if (activeSettingsSave === operation) { activeSettingsSave = null; } if (asyncSettingsSaveQueued) { - const queued = asyncSettingsSaveQueued; - asyncSettingsSaveQueued = null; - void saveSettingsPayloadAsync(queued.paths, queued.settings, queued.generation); + startSettingsSaveDrain(); } }); activeSettingsSave = operation; - await operation; + void operation.catch(() => {}); +} + +function saveSettingsPayloadAsync(paths: StoragePaths, settings: AppSettings, generation: number): Promise { + return new Promise((resolve, reject) => { + const waiter = { resolve, reject }; + if (asyncSettingsSaveQueued) { + asyncSettingsSaveQueued.paths = paths; + asyncSettingsSaveQueued.settings = settings; + asyncSettingsSaveQueued.generation = generation; + asyncSettingsSaveQueued.waiters.push(waiter); + } else { + asyncSettingsSaveQueued = { paths, settings, generation, waiters: [waiter] }; + } + startSettingsSaveDrain(); + }); } export async function saveSettingsAsync(paths: StoragePaths, settings: AppSettings): Promise { @@ -1446,25 +1539,8 @@ export function loadSessionWithStatus(paths: StoragePaths): SessionLoadResult { const primary = primaryExists ? readSessionFile(paths.sessionFile) : null; if (primary) { - const primaryPkgCount = Object.keys(primary.session.packages).length; - if (primaryPkgCount === 0 && backupExists) { - const backup = readSessionFile(backupFile); - if (backup) { - const backupPkgCount = Object.keys(backup.session.packages).length; - if (backupPkgCount > 0) { - logger.warn(`Session-Datei ist leer (0 Pakete), aber Backup hat ${backupPkgCount} Pakete — verwende Backup`); - try { - const payload = JSON.stringify({ ...backup.session, updatedAt: Date.now() }, safeJsonReplacer); - fs.writeFileSync(syncTempFile, payload, "utf8"); - syncRenameWithExdevFallback(syncTempFile, paths.sessionFile); - } catch { - } - return { session: backup.session, status: "recovered-backup", wasRunning: backup.wasRunning }; - } - } - } return { session: primary.session, status: "ok", wasRunning: primary.wasRunning }; - } + } const backup = backupExists ? readSessionFile(backupFile) : null; if (backup) { @@ -1532,8 +1608,14 @@ export function saveSession(paths: StoragePaths, session: SessionState): void { } } -let asyncSaveRunning = false; -let asyncSaveQueued: { paths: StoragePaths; payload: string; generation: number } | null = null; +interface QueuedSessionSave { + paths: StoragePaths; + payload: string; + generation: number; + waiters: AsyncSaveWaiter[]; +} + +let asyncSaveQueued: QueuedSessionSave | null = null; let syncSaveGeneration = 0; let activeSessionSave: Promise | null = null; @@ -1568,27 +1650,52 @@ async function writeSessionPayload(paths: StoragePaths, payload: string, generat } } -async function saveSessionPayloadAsync(paths: StoragePaths, payload: string, generation: number): Promise { - if (asyncSaveRunning) { - asyncSaveQueued = { paths, payload, generation }; +function startSessionSaveDrain(): void { + if (activeSessionSave) { return; } - asyncSaveRunning = true; - const operation = writeSessionPayload(paths, payload, generation).catch((error) => { - logger.error(`Async Session-Save fehlgeschlagen: ${String(error)}`); - }).finally(() => { - asyncSaveRunning = false; + let operation: Promise; + operation = (async () => { + while (asyncSaveQueued) { + const current = asyncSaveQueued; + asyncSaveQueued = null; + try { + await writeSessionPayload(current.paths, current.payload, current.generation); + current.waiters.forEach((waiter) => waiter.resolve()); + } catch (error) { + current.waiters.forEach((waiter) => waiter.reject(error)); + const pending = asyncSaveQueued as QueuedSessionSave | null; + asyncSaveQueued = null; + pending?.waiters.forEach((waiter) => waiter.reject(error)); + logger.error(`Async Session-Save fehlgeschlagen: ${String(error)}`); + throw error; + } + } + })().finally(() => { if (activeSessionSave === operation) { activeSessionSave = null; } if (asyncSaveQueued) { - const queued = asyncSaveQueued; - asyncSaveQueued = null; - void saveSessionPayloadAsync(queued.paths, queued.payload, queued.generation); + startSessionSaveDrain(); } }); activeSessionSave = operation; - await operation; + void operation.catch(() => {}); +} + +function saveSessionPayloadAsync(paths: StoragePaths, payload: string, generation: number): Promise { + return new Promise((resolve, reject) => { + const waiter = { resolve, reject }; + if (asyncSaveQueued) { + asyncSaveQueued.paths = paths; + asyncSaveQueued.payload = payload; + asyncSaveQueued.generation = generation; + asyncSaveQueued.waiters.push(waiter); + } else { + asyncSaveQueued = { paths, payload, generation, waiters: [waiter] }; + } + startSessionSaveDrain(); + }); } interface BlockedSettingsSave { @@ -1616,6 +1723,47 @@ export interface PersistenceBarrier { let activePersistenceBarrier: PersistenceBarrierState | null = null; +interface DetachedAsyncSaves { + settings: QueuedSettingsSave | null; + session: QueuedSessionSave | null; +} + +function detachPendingAsyncSaves(): DetachedAsyncSaves { + const detached = { + settings: asyncSettingsSaveQueued, + session: asyncSaveQueued + }; + asyncSettingsSaveQueued = null; + asyncSaveQueued = null; + return detached; +} + +function rejectDetachedAsyncSaves(detached: DetachedAsyncSaves, error: unknown): void { + detached.settings?.waiters.forEach((waiter) => waiter.reject(error)); + detached.session?.waiters.forEach((waiter) => waiter.reject(error)); +} + +function requeueDetachedAsyncSaves(detached: DetachedAsyncSaves): void { + if (detached.settings) { + detached.settings.generation = syncSettingsSaveGeneration; + if (asyncSettingsSaveQueued) { + asyncSettingsSaveQueued.waiters.unshift(...detached.settings.waiters); + } else { + asyncSettingsSaveQueued = detached.settings; + } + startSettingsSaveDrain(); + } + if (detached.session) { + detached.session.generation = syncSaveGeneration; + if (asyncSaveQueued) { + asyncSaveQueued.waiters.unshift(...detached.session.waiters); + } else { + asyncSaveQueued = detached.session; + } + startSessionSaveDrain(); + } +} + function blockSettingsSave(state: PersistenceBarrierState, paths: StoragePaths, settings: AppSettings): Promise { return new Promise((resolve, reject) => { const waiter = { resolve, reject }; @@ -1643,28 +1791,61 @@ function blockSessionSave(state: PersistenceBarrierState, paths: StoragePaths, p } async function replayBlockedSaves(state: PersistenceBarrierState): Promise { + let firstError: unknown; + let failed = false; while (state.blockedSettings || state.blockedSession) { const blockedSettings = state.blockedSettings; const blockedSession = state.blockedSession; state.blockedSettings = null; state.blockedSession = null; - try { - await Promise.all([ - blockedSettings - ? saveSettingsPayloadAsync(blockedSettings.paths, blockedSettings.settings, syncSettingsSaveGeneration) - : Promise.resolve(), - blockedSession - ? saveSessionPayloadAsync(blockedSession.paths, blockedSession.payload, syncSaveGeneration) - : Promise.resolve() - ]); - blockedSettings?.waiters.forEach((waiter) => waiter.resolve()); - blockedSession?.waiters.forEach((waiter) => waiter.resolve()); - } catch (error) { - blockedSettings?.waiters.forEach((waiter) => waiter.reject(error)); - blockedSession?.waiters.forEach((waiter) => waiter.reject(error)); - throw error; + const [settingsResult, sessionResult] = await Promise.allSettled([ + blockedSettings + ? saveSettingsPayloadAsync(blockedSettings.paths, blockedSettings.settings, syncSettingsSaveGeneration) + : Promise.resolve(), + blockedSession + ? saveSessionPayloadAsync(blockedSession.paths, blockedSession.payload, syncSaveGeneration) + : Promise.resolve() + ]); + if (blockedSettings) { + if (settingsResult.status === "fulfilled") { + blockedSettings.waiters.forEach((waiter) => waiter.resolve()); + } else { + blockedSettings.waiters.forEach((waiter) => waiter.reject(settingsResult.reason)); + if (!failed) { + firstError = settingsResult.reason; + failed = true; + } + } + } + if (blockedSession) { + if (sessionResult.status === "fulfilled") { + blockedSession.waiters.forEach((waiter) => waiter.resolve()); + } else { + blockedSession.waiters.forEach((waiter) => waiter.reject(sessionResult.reason)); + if (!failed) { + firstError = sessionResult.reason; + failed = true; + } + } } } + if (failed) { + throw firstError; + } +} + +function rejectBlockedSaves(state: PersistenceBarrierState, error: unknown): void { + state.blockedSettings?.waiters.forEach((waiter) => waiter.reject(error)); + state.blockedSession?.waiters.forEach((waiter) => waiter.reject(error)); + state.blockedSettings = null; + state.blockedSession = null; +} + +function resolveBlockedSaves(state: PersistenceBarrierState): void { + state.blockedSettings?.waiters.forEach((waiter) => waiter.resolve()); + state.blockedSession?.waiters.forEach((waiter) => waiter.resolve()); + state.blockedSettings = null; + state.blockedSession = null; } export async function acquirePersistenceBarrier(): Promise { @@ -1680,19 +1861,46 @@ export async function acquirePersistenceBarrier(): Promise { resolveReleased }; activePersistenceBarrier = state; - cancelPendingAsyncSaves(); - await Promise.all([activeSettingsSave ?? Promise.resolve(), activeSessionSave ?? Promise.resolve()]); + const detachedSaves = detachPendingAsyncSaves(); + const activeWriteResults = await Promise.allSettled([ + activeSettingsSave ?? Promise.resolve(), + activeSessionSave ?? Promise.resolve() + ]); + const activeWriteFailure = activeWriteResults.find((result) => result.status === "rejected"); + if (activeWriteFailure?.status === "rejected") { + rejectBlockedSaves(state, activeWriteFailure.reason); + if (activePersistenceBarrier === state) { + activePersistenceBarrier = null; + } + requeueDetachedAsyncSaves(detachedSaves); + state.resolveReleased(); + throw activeWriteFailure.reason; + } + rejectDetachedAsyncSaves(detachedSaves, new Error("Persistenzbarriere hat ausstehenden Save verworfen")); let finished = false; return { release: async ({ replayBlocked }) => { if (finished) return; finished = true; + let replayError: unknown; + let replayFailed = false; try { if (replayBlocked) { - await replayBlockedSaves(state); + do { + try { + await replayBlockedSaves(state); + } catch (error) { + if (!replayFailed) { + replayError = error; + replayFailed = true; + } + } + } while (state.blockedSettings || state.blockedSession); } else { - state.blockedSettings?.waiters.forEach((waiter) => waiter.resolve()); - state.blockedSession?.waiters.forEach((waiter) => waiter.resolve()); + resolveBlockedSaves(state); + } + if (replayFailed) { + throw replayError; } } finally { if (activePersistenceBarrier === state) { @@ -1704,12 +1912,13 @@ export async function acquirePersistenceBarrier(): Promise { }; } -export function cancelPendingAsyncSaves(): void { - asyncSaveQueued = null; - asyncSettingsSaveQueued = null; - syncSaveGeneration += 1; - syncSettingsSaveGeneration += 1; -} +export function cancelPendingAsyncSaves(): void { + syncSettingsSaveGeneration += 1; + syncSaveGeneration += 1; + const detached = detachPendingAsyncSaves(); + detached.session?.waiters.forEach((waiter) => waiter.resolve()); + detached.settings?.waiters.forEach((waiter) => waiter.resolve()); +} export async function saveSessionAsync(paths: StoragePaths, session: SessionState): Promise { const payload = JSON.stringify({ ...session, updatedAt: Date.now() }, safeJsonReplacer); diff --git a/src/main/update.ts b/src/main/update.ts index 178e783..ef2ddc5 100644 --- a/src/main/update.ts +++ b/src/main/update.ts @@ -698,14 +698,138 @@ async function downloadFile( reportProgress(true); - await fs.promises.mkdir(path.dirname(targetPath), { recursive: true }); - const tempPath = `${targetPath}.tmp`; - const writeStream = fs.createWriteStream(tempPath); - const reader = response.body.getReader(); - - const idleMs = getBodyIdleTimeout(); - let idleTimer: NodeJS.Timeout | null = null; - let idleTimedOut = false; + await fs.promises.mkdir(path.dirname(targetPath), { recursive: true }); + const tempPath = `${targetPath}.tmp`; + const writeStream = fs.createWriteStream(tempPath); + const reader = response.body.getReader(); + let writeError: unknown = null; + const idleMs = getBodyIdleTimeout(); + const idleAbortController = new AbortController(); + let idleTimer: NodeJS.Timeout | null = null; + let idleTimedOut = false; + + const idleTimeoutError = (): Error => new Error(`Update Download Body Timeout nach ${Math.ceil(idleMs / 1000)}s`); + + const onWriteError = (error: unknown): void => { + writeError = error; + void reader.cancel().catch(() => undefined); + }; + const onWriteClose = (): void => { + writeStream.removeListener("error", onWriteError); + }; + writeStream.on("error", onWriteError); + writeStream.once("close", onWriteClose); + + const waitForDrain = async (): Promise => { + if (writeError !== null) throw writeError; + if (idleTimedOut) throw idleTimeoutError(); + if (shutdown?.aborted) throw new Error("aborted:update_shutdown"); + + await new Promise((resolve, reject) => { + let settled = false; + const cleanup = (): void => { + writeStream.removeListener("drain", onDrain); + writeStream.removeListener("error", onError); + idleAbortController.signal.removeEventListener("abort", onIdleTimeout); + shutdown?.removeEventListener("abort", onAbort); + }; + const settle = (error?: unknown): void => { + if (settled) return; + settled = true; + cleanup(); + if (error !== undefined) { + reject(error); + } else { + resolve(); + } + }; + const onDrain = (): void => settle(); + const onError = (error: unknown): void => settle(error); + const onIdleTimeout = (): void => settle(idleTimeoutError()); + const onAbort = (): void => settle(new Error("aborted:update_shutdown")); + + writeStream.once("drain", onDrain); + writeStream.once("error", onError); + idleAbortController.signal.addEventListener("abort", onIdleTimeout, { once: true }); + shutdown?.addEventListener("abort", onAbort, { once: true }); + + if (writeError !== null) { + settle(writeError); + } else if (idleTimedOut) { + onIdleTimeout(); + } else if (shutdown?.aborted) { + onAbort(); + } + }); + }; + + const finishWriting = async (): Promise => { + if (writeError !== null) throw writeError; + if (idleTimedOut) throw idleTimeoutError(); + if (shutdown?.aborted) throw new Error("aborted:update_shutdown"); + + await new Promise((resolve, reject) => { + let settled = false; + let endSucceeded = false; + let closed = false; + const cleanup = (): void => { + writeStream.removeListener("error", onError); + writeStream.removeListener("close", onClose); + idleAbortController.signal.removeEventListener("abort", onIdleTimeout); + shutdown?.removeEventListener("abort", onAbort); + }; + const settle = (error?: unknown): void => { + if (settled) return; + settled = true; + cleanup(); + if (error !== undefined) { + reject(error); + } else { + resolve(); + } + }; + const onError = (error: unknown): void => settle(error); + const settleIfComplete = (): void => { + if (endSucceeded && closed) settle(); + }; + const onClose = (): void => { + closed = true; + if (writeError !== null) { + settle(writeError); + return; + } + settleIfComplete(); + }; + const onIdleTimeout = (): void => settle(idleTimeoutError()); + const onAbort = (): void => settle(new Error("aborted:update_shutdown")); + + writeStream.once("error", onError); + writeStream.once("close", onClose); + idleAbortController.signal.addEventListener("abort", onIdleTimeout, { once: true }); + shutdown?.addEventListener("abort", onAbort, { once: true }); + + if (writeError !== null) { + settle(writeError); + } else if (idleTimedOut) { + onIdleTimeout(); + } else if (shutdown?.aborted) { + onAbort(); + } else { + try { + writeStream.end((error?: unknown) => { + if (error !== undefined && error !== null) { + settle(error); + } else { + endSucceeded = true; + settleIfComplete(); + } + }); + } catch (error) { + settle(error); + } + } + }); + }; const clearIdle = (): void => { if (idleTimer) { @@ -717,51 +841,53 @@ async function downloadFile( const resetIdle = (): void => { clearIdle(); if (idleMs > 0) { - idleTimer = setTimeout(() => { - idleTimedOut = true; - reader.cancel().catch(() => undefined); - }, idleMs); + idleTimer = setTimeout(() => { + idleTimedOut = true; + idleAbortController.abort(); + void reader.cancel().catch(() => undefined); + }, idleMs); } }; try { - resetIdle(); - for (;;) { - if (shutdown?.aborted) { - await reader.cancel().catch(() => undefined); - throw new Error("aborted:update_shutdown"); - } - const { done, value } = await reader.read(); - if (done) break; - - const buf = Buffer.from(value.buffer, value.byteOffset, value.byteLength); - if (!writeStream.write(buf)) { - await new Promise((resolve) => writeStream.once("drain", resolve)); - } - downloadedBytes += buf.byteLength; - resetIdle(); - reportProgress(false); - } - } catch (error) { - writeStream.destroy(); - await fs.promises.rm(tempPath, { force: true }).catch(() => {}); - if (idleTimedOut) { - throw new Error(`Update Download Body Timeout nach ${Math.ceil(idleMs / 1000)}s`); - } - throw error; - } finally { - clearIdle(); - } - - await new Promise((resolve, reject) => { - writeStream.end(() => resolve()); - writeStream.on("error", reject); - }); - - if (idleTimedOut) { - await fs.promises.rm(tempPath, { force: true }).catch(() => {}); - throw new Error(`Update Download Body Timeout nach ${Math.ceil(idleMs / 1000)}s`); - } + resetIdle(); + for (;;) { + if (shutdown?.aborted) { + void reader.cancel().catch(() => undefined); + throw new Error("aborted:update_shutdown"); + } + const { done, value } = await reader.read(); + if (writeError !== null) throw writeError; + if (done) break; + + const buf = Buffer.from(value.buffer, value.byteOffset, value.byteLength); + if (!writeStream.write(buf)) { + await waitForDrain(); + } + if (writeError !== null) throw writeError; + downloadedBytes += buf.byteLength; + resetIdle(); + reportProgress(false); + } + + await finishWriting(); + if (writeError !== null) throw writeError; + } catch (error) { + writeStream.destroy(); + void reader.cancel().catch(() => undefined); + await fs.promises.rm(tempPath, { force: true }).catch(() => {}); + if (idleTimedOut) { + throw idleTimeoutError(); + } + throw error; + } finally { + clearIdle(); + } + + if (idleTimedOut) { + await fs.promises.rm(tempPath, { force: true }).catch(() => {}); + throw idleTimeoutError(); + } if (totalBytes && downloadedBytes !== totalBytes) { await fs.promises.rm(tempPath, { force: true }).catch(() => {}); diff --git a/tests/debug-server.test.ts b/tests/debug-server.test.ts index fbda6b3..83a9a98 100644 --- a/tests/debug-server.test.ts +++ b/tests/debug-server.test.ts @@ -44,7 +44,7 @@ import { defaultSettings } from "../src/main/constants"; import { configureCredentialProtector } from "../src/main/credential-protection"; import { createRendererState } from "../src/main/renderer-state"; import { getAuditLogPath, initAuditLog, logAuditEvent, shutdownAuditLog } from "../src/main/audit-log"; -import { startDebugServer, stopDebugServer } from "../src/main/debug-server"; +import { getDebugServerRuntimeStatus, restartDebugServer, startDebugServer, stopDebugServer, writeDebugServerConfig } from "../src/main/debug-server"; import { ensureItemLog, initItemLogs, shutdownItemLogs } from "../src/main/item-log"; import { configureLogger, getLogFilePath, logger } from "../src/main/logger"; import { ensurePackageLog, initPackageLogs, shutdownPackageLogs } from "../src/main/package-log"; @@ -370,6 +370,60 @@ afterEach(() => { }); describe("debug-server", () => { + it("keeps exactly one reachable server across concurrent restarts", async () => { + const fixture = await createFixture(); + writeDebugServerConfig({ allowlist: ["203.0.113.5"] }); + + await Promise.all([restartDebugServer(), restartDebugServer()]); + + expect(getDebugServerRuntimeStatus()).toEqual(expect.objectContaining({ + running: true, + allowlistCount: 1 + })); + const response = await authedFetch(`${fixture.baseUrl}/health`, fixture.token); + expect(response.ok).toBe(true); + }); + + it("does not reopen after stop races with queued restarts", async () => { + const fixture = await createFixture(); + const restarts = Promise.all([restartDebugServer(), restartDebugServer()]); + + stopDebugServer(); + await restarts; + + expect(getDebugServerRuntimeStatus().running).toBe(false); + await expect(fetch(`${fixture.baseUrl}/health`, { + signal: AbortSignal.timeout(1000) + })).rejects.toThrow(); + }); + + it("settles the lifecycle queue when stopped during socket startup", async () => { + const baseDir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-debug-start-stop-")); + tempDirs.push(baseDir); + const port = await getFreePort(); + fs.writeFileSync(path.join(baseDir, "debug_token.txt"), "debug-secret", "utf8"); + fs.writeFileSync(path.join(baseDir, "debug_port.txt"), String(port), "utf8"); + fs.writeFileSync(path.join(baseDir, "debug_host.txt"), "127.0.0.1", "utf8"); + fs.writeFileSync(path.join(baseDir, "debug_allowlist.txt"), "", "utf8"); + + startDebugServer({} as DownloadManager, baseDir); + await Promise.resolve(); + await Promise.resolve(); + stopDebugServer(); + + const outcome = await Promise.race([ + restartDebugServer().then(() => "settled"), + new Promise((resolve) => setTimeout(() => resolve("timeout"), 500)) + ]); + + expect(outcome).toBe("settled"); + expect(getDebugServerRuntimeStatus().running).toBe(false); + + startDebugServer({} as DownloadManager, baseDir); + await waitForReady(`http://127.0.0.1:${port}/health`); + expect(getDebugServerRuntimeStatus().running).toBe(true); + }); + it("serves the exact safe notification DTO in its endpoint and diagnostics", async () => { const fixture = await createFixture(); const response = await authedFetch(`${fixture.baseUrl}/notifications`, fixture.token); diff --git a/tests/session-restart-loss.test.ts b/tests/session-restart-loss.test.ts index fbefcd9..33a414c 100644 --- a/tests/session-restart-loss.test.ts +++ b/tests/session-restart-loss.test.ts @@ -185,17 +185,17 @@ describe("session restart loss", () => { expect(Object.keys(loaded.items)).toEqual([]); }); - it("recovers from the backup when the primary exists but is empty", () => { + it("keeps a valid empty primary authoritative over an older populated backup", () => { const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-loss-")); tempDirs.push(dir); const paths = createStoragePaths(dir); fs.writeFileSync(paths.sessionFile, JSON.stringify(emptySession()), "utf8"); - fs.writeFileSync(`${paths.sessionFile}.bak`, JSON.stringify(sessionWith(["A", "B"])), "utf8"); - - const loaded = loadSession(paths); - expect(Object.keys(loaded.packages).sort()).toEqual(["A", "B"]); - }); + fs.writeFileSync(`${paths.sessionFile}.bak`, JSON.stringify(sessionWith(["A", "B"])), "utf8"); + + const loaded = loadSession(paths); + expect(Object.keys(loaded.packages)).toEqual([]); + }); it("does not let an in-flight/queued async settings save clobber a newer synchronous saveSettings", async () => { const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-settings-race-")); diff --git a/tests/storage.test.ts b/tests/storage.test.ts index 8212855..3b4f4dc 100644 --- a/tests/storage.test.ts +++ b/tests/storage.test.ts @@ -1907,7 +1907,7 @@ describe("settings storage", () => { expect(loaded.packageOrder).toEqual(empty.packageOrder); }); - it("loads backup session when primary session is corrupted", () => { + it("loads backup session when primary session is corrupted", () => { const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-")); tempDirs.push(dir); const paths = createStoragePaths(dir); @@ -1958,8 +1958,68 @@ describe("settings storage", () => { expect(loaded.items["item-backup"]?.fileName).toBe("backup-file.rar"); const restoredPrimary = JSON.parse(fs.readFileSync(paths.sessionFile, "utf8")) as { packages?: Record }; - expect(restoredPrimary.packages && "pkg-backup" in restoredPrimary.packages).toBe(true); - }); + expect(restoredPrimary.packages && "pkg-backup" in restoredPrimary.packages).toBe(true); + }); + + it("keeps a valid intentionally empty primary session authoritative over its populated backup", () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-empty-primary-")); + tempDirs.push(dir); + const paths = createStoragePaths(dir); + const populated = normalizeLoadedSession({ + ...emptySession(), + packageOrder: ["pkg-old"], + packages: { + "pkg-old": { + id: "pkg-old", + name: "Old Package", + outputDir: path.join(dir, "out"), + extractDir: path.join(dir, "extract"), + status: "queued", + itemIds: [] + } + } + }); + + saveSession(paths, populated); + saveSession(paths, emptySession()); + const loaded = loadSessionWithStatus(paths); + + expect(loaded.status).toBe("ok"); + expect(loaded.session.packageOrder).toEqual([]); + expect(loaded.session.packages).toEqual({}); + }); + + it.each([ + ["empty object", {}], + ["null", null], + ["array", []] + ])("recovers the backup when the primary contains a JSON-valid invalid %s envelope", (_label, invalidPrimary) => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-invalid-session-envelope-")); + tempDirs.push(dir); + const paths = createStoragePaths(dir); + const backup = normalizeLoadedSession({ + ...emptySession(), + packageOrder: ["pkg-backup"], + packages: { + "pkg-backup": { + id: "pkg-backup", + name: "Backup Package", + outputDir: path.join(dir, "out"), + extractDir: path.join(dir, "extract"), + status: "queued", + itemIds: [] + } + } + }); + fs.writeFileSync(`${paths.sessionFile}.bak`, JSON.stringify(backup), "utf8"); + fs.writeFileSync(paths.sessionFile, JSON.stringify(invalidPrimary), "utf8"); + + const loaded = loadSessionWithStatus(paths); + + expect(loaded.status).toBe("recovered-backup"); + expect(loaded.session.packageOrder).toEqual(["pkg-backup"]); + expect(Object.keys(loaded.session.packages)).toEqual(["pkg-backup"]); + }); it("returns defaults when config file contains invalid JSON", () => { const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-")); @@ -1977,7 +2037,7 @@ describe("settings storage", () => { expect(loaded.cleanupMode).toBe(defaults.cleanupMode); }); - it("loads backup config when primary config is corrupted", () => { + it("loads backup config when primary config is corrupted", () => { const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-")); tempDirs.push(dir); const paths = createStoragePaths(dir); @@ -1992,8 +2052,70 @@ describe("settings storage", () => { const loaded = loadSettings(paths); expect(loaded.outputDir).toBe(backupSettings.outputDir); - expect(loaded.packageName).toBe("from-backup"); - }); + expect(loaded.packageName).toBe("from-backup"); + }); + + it.each([ + ["empty object", {}], + ["null", null], + ["array", []], + ["unrecognized object", { unrelated: true }] + ])("loads and repairs the backup when the primary config contains a JSON-valid invalid %s envelope", (_label, invalidPrimary) => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-invalid-settings-envelope-")); + tempDirs.push(dir); + const paths = createStoragePaths(dir); + const backupSettings = { + ...defaultSettings(), + outputDir: path.join(dir, "backup-output"), + packageName: "from-valid-backup" + }; + fs.writeFileSync(`${paths.configFile}.bak`, JSON.stringify(backupSettings), "utf8"); + fs.writeFileSync(paths.configFile, JSON.stringify(invalidPrimary), "utf8"); + + const loaded = loadSettings(paths); + + expect(loaded.outputDir).toBe(backupSettings.outputDir); + expect(loaded.packageName).toBe("from-valid-backup"); + expect(loadSettings(paths).packageName).toBe("from-valid-backup"); + }); + + it("keeps a sparse legacy primary config authoritative over its backup", () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-sparse-settings-")); + tempDirs.push(dir); + const paths = createStoragePaths(dir); + const legacyOutputDir = path.join(dir, "legacy-output"); + fs.writeFileSync(`${paths.configFile}.bak`, JSON.stringify({ + ...defaultSettings(), + outputDir: path.join(dir, "backup-output"), + packageName: "from-backup" + }), "utf8"); + fs.writeFileSync(paths.configFile, JSON.stringify({ outputDir: legacyOutputDir }), "utf8"); + + const loaded = loadSettings(paths); + + expect(loaded.outputDir).toBe(legacyOutputDir); + expect(loaded.packageName).toBe(""); + }); + + it("loads and repairs backup config when the primary config is missing", () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-missing-config-")); + tempDirs.push(dir); + const paths = createStoragePaths(dir); + const settings = { + ...defaultSettings(), + outputDir: path.join(dir, "kept-output"), + packageName: "from-existing-backup" + }; + saveSettings(paths, settings); + fs.rmSync(paths.configFile); + + const loaded = loadSettings(paths); + + expect(loaded.outputDir).toBe(settings.outputDir); + expect(loaded.packageName).toBe("from-existing-backup"); + expect(loadSettings(paths).packageName).toBe("from-existing-backup"); + expect(fs.existsSync(paths.configFile)).toBe(true); + }); it("sanitizes malformed persisted session structures", () => { const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-")); @@ -2035,7 +2157,7 @@ describe("settings storage", () => { expect(loaded.packageOrder).toEqual(["pkg-valid"]); }); - it("drops unsafe session ids and target paths outside the package output directory", () => { + it("drops unsafe session ids and target paths outside the package output directory", () => { const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-")); tempDirs.push(dir); const paths = createStoragePaths(dir); @@ -2101,9 +2223,105 @@ describe("settings storage", () => { expect(Object.keys(loaded.items).sort()).toEqual(["item-outside", "item-safe"]); expect(loaded.packageOrder).toEqual(["pkg-safe"]); expect(path.resolve(loaded.items["item-safe"]?.targetPath || "")).toBe(path.resolve(safeTargetPath)); - expect(loaded.items["item-outside"]?.targetPath).toBe(""); - }); - + expect(loaded.items["item-outside"]?.targetPath).toBe(""); + }); + + it("preserves a metadata rename journal inside the package output directory across save and load", () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-")); + tempDirs.push(dir); + const paths = createStoragePaths(dir); + const outputDir = path.join(dir, "downloads", "journal"); + const currentPath = path.join(outputDir, "opaque.bin"); + const renameTargetPath = path.join(outputDir, "resolved-name.rar"); + const session = emptySession(); + session.packageOrder = ["pkg-journal"]; + session.packages["pkg-journal"] = { + id: "pkg-journal", + name: "Journal Package", + outputDir, + extractDir: path.join(dir, "extract", "journal"), + status: "completed", + itemIds: ["item-journal"], + cancelled: false, + enabled: true, + createdAt: 1, + updatedAt: 1 + }; + session.items["item-journal"] = { + id: "item-journal", + packageId: "pkg-journal", + url: "https://ddownload.com/journal123", + provider: null, + status: "completed", + retries: 0, + speedBps: 0, + downloadedBytes: 1024, + totalBytes: 1024, + progressPercent: 100, + fileName: "resolved-name.rar", + targetPath: currentPath, + metadataRenameTargetPath: renameTargetPath, + resumable: true, + attempts: 1, + lastError: "", + fullStatus: "Fertig", + createdAt: 1, + updatedAt: 1 + }; + + saveSession(paths, session); + + const loaded = loadSession(paths); + expect(loaded.items["item-journal"]?.metadataRenameTargetPath).toBe(path.resolve(renameTargetPath)); + }); + + it("drops a metadata rename journal outside the package output directory", () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-")); + tempDirs.push(dir); + const paths = createStoragePaths(dir); + const outputDir = path.join(dir, "downloads", "journal"); + const session = emptySession(); + session.packageOrder = ["pkg-journal"]; + session.packages["pkg-journal"] = { + id: "pkg-journal", + name: "Journal Package", + outputDir, + extractDir: path.join(dir, "extract", "journal"), + status: "completed", + itemIds: ["item-journal"], + cancelled: false, + enabled: true, + createdAt: 1, + updatedAt: 1 + }; + session.items["item-journal"] = { + id: "item-journal", + packageId: "pkg-journal", + url: "https://ddownload.com/journal123", + provider: null, + status: "completed", + retries: 0, + speedBps: 0, + downloadedBytes: 1024, + totalBytes: 1024, + progressPercent: 100, + fileName: "resolved-name.rar", + targetPath: path.join(outputDir, "opaque.bin"), + metadataRenameTargetPath: path.join(dir, "outside", "resolved-name.rar"), + resumable: true, + attempts: 1, + lastError: "", + fullStatus: "Fertig", + createdAt: 1, + updatedAt: 1 + }; + + saveSession(paths, session); + + const loaded = loadSession(paths); + expect(loaded.items["item-journal"]?.metadataRenameTargetPath).toBeUndefined(); + }); + it("captures async session save payload before later mutations", async () => { const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-")); tempDirs.push(dir); @@ -2120,6 +2338,158 @@ describe("settings storage", () => { expect(persisted.summaryText).toBe("before-mutation"); }); + it("keeps a queued async settings save pending until its payload is persisted", async () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-settings-queue-")); + tempDirs.push(dir); + const paths = createStoragePaths(dir); + saveSettings(paths, defaultSettings()); + const writeFile = fs.promises.writeFile.bind(fs.promises); + let releaseFirstWrite = () => {}; + let markFirstWriteStarted = () => {}; + const firstWriteStarted = new Promise((resolve) => { markFirstWriteStarted = resolve; }); + const firstWriteRelease = new Promise((resolve) => { releaseFirstWrite = resolve; }); + let settingsTempWrites = 0; + vi.spyOn(fs.promises, "writeFile").mockImplementation(async (filePath, data, options) => { + await writeFile(filePath, data, options); + if (String(filePath) === `${paths.configFile}.settings.tmp` && settingsTempWrites++ === 0) { + markFirstWriteStarted(); + await firstWriteRelease; + } + }); + + const first = saveSettingsAsync(paths, { ...defaultSettings(), packageName: "first" }); + await firstWriteStarted; + let secondSettled = false; + const second = saveSettingsAsync(paths, { ...defaultSettings(), packageName: "second" }) + .then(() => { secondSettled = true; }); + await new Promise((resolve) => setImmediate(resolve)); + const settledWhileBlocked = secondSettled; + + releaseFirstWrite(); + await Promise.all([first, second]); + await vi.waitFor(() => expect(loadSettings(paths).packageName).toBe("second")); + + expect(settledWhileBlocked).toBe(false); + }); + + it("keeps a queued async session save pending until its payload is persisted", async () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-session-queue-")); + tempDirs.push(dir); + const paths = createStoragePaths(dir); + const openFile = fs.promises.open.bind(fs.promises); + let releaseFirstWrite = () => {}; + let markFirstWriteStarted = () => {}; + const firstWriteStarted = new Promise((resolve) => { markFirstWriteStarted = resolve; }); + const firstWriteRelease = new Promise((resolve) => { releaseFirstWrite = resolve; }); + let sessionTempWrites = 0; + vi.spyOn(fs.promises, "open").mockImplementation(async (filePath, flags, mode) => { + const handle = await openFile(filePath, flags, mode); + if (String(filePath) === `${paths.sessionFile}.async.tmp`) { + const writeFile = handle.writeFile.bind(handle); + handle.writeFile = async (...args: Parameters) => { + const result = await writeFile(...args); + if (sessionTempWrites++ === 0) { + markFirstWriteStarted(); + await firstWriteRelease; + } + return result; + }; + } + return handle; + }); + + const first = saveSessionAsync(paths, { ...emptySession(), summaryText: "first" }); + await firstWriteStarted; + let secondSettled = false; + const second = saveSessionAsync(paths, { ...emptySession(), summaryText: "second" }) + .then(() => { secondSettled = true; }); + await new Promise((resolve) => setImmediate(resolve)); + const settledWhileBlocked = secondSettled; + + releaseFirstWrite(); + await Promise.all([first, second]); + await vi.waitFor(() => expect(loadSession(paths).summaryText).toBe("second")); + + expect(settledWhileBlocked).toBe(false); + }); + + it("drains a settings save queued in the completion tail", async () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-settings-tail-")); + tempDirs.push(dir); + const paths = createStoragePaths(dir); + const renameFile = fs.renameSync.bind(fs); + let secondSave: Promise | null = null; + let markSecondSaveStarted = () => {}; + const secondSaveStarted = new Promise((resolve) => { markSecondSaveStarted = resolve; }); + let scheduled = false; + vi.spyOn(fs, "renameSync").mockImplementation((source, destination) => { + const result = renameFile(source, destination); + if (!scheduled && String(destination) === paths.configFile) { + scheduled = true; + queueMicrotask(() => queueMicrotask(() => { + secondSave = saveSettingsAsync(paths, { ...defaultSettings(), packageName: "second" }); + markSecondSaveStarted(); + })); + } + return result; + }); + + const firstSave = saveSettingsAsync(paths, { ...defaultSettings(), packageName: "first" }); + await secondSaveStarted; + await firstSave; + expect(secondSave).not.toBeNull(); + await secondSave; + + expect(loadSettings(paths).packageName).toBe("second"); + }); + + it("drains a session save queued in the completion tail", async () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-session-tail-")); + tempDirs.push(dir); + const paths = createStoragePaths(dir); + const renameFile = fs.renameSync.bind(fs); + let secondSave: Promise | null = null; + let markSecondSaveStarted = () => {}; + const secondSaveStarted = new Promise((resolve) => { markSecondSaveStarted = resolve; }); + let scheduled = false; + vi.spyOn(fs, "renameSync").mockImplementation((source, destination) => { + const result = renameFile(source, destination); + if (!scheduled && String(destination) === paths.sessionFile) { + scheduled = true; + queueMicrotask(() => queueMicrotask(() => { + secondSave = saveSessionAsync(paths, { ...emptySession(), summaryText: "second" }); + markSecondSaveStarted(); + })); + } + return result; + }); + + const firstSave = saveSessionAsync(paths, { ...emptySession(), summaryText: "first" }); + await secondSaveStarted; + await firstSave; + expect(secondSave).not.toBeNull(); + await secondSave; + + expect(loadSession(paths).summaryText).toBe("second"); + }); + + it("rejects async settings and session saves when persistence fails", async () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-async-errors-")); + tempDirs.push(dir); + const paths = createStoragePaths(dir); + const settingsFailure = vi.spyOn(fs.promises, "writeFile") + .mockRejectedValueOnce(new Error("settings write failed")); + + await expect(saveSettingsAsync(paths, defaultSettings())).rejects.toThrow("settings write failed"); + settingsFailure.mockRestore(); + + const sessionFailure = vi.spyOn(fs.promises, "open") + .mockRejectedValueOnce(new Error("session write failed")); + + await expect(saveSessionAsync(paths, emptySession())).rejects.toThrow("session write failed"); + sessionFailure.mockRestore(); + }); + it("keeps newer synchronous settings and session saves when older async commits resume", async () => { const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-sync-wins-")); tempDirs.push(dir); @@ -2260,6 +2630,8 @@ describe("settings storage", () => { releaseSessionWrite(); const barrier = await barrierPromise; await Promise.all([oldSettingsSave, oldSessionSave]); + expect(loadSettings(paths).packageName).toBe("old-settings"); + expect(loadSession(paths).summaryText).toBe("old-session"); const blockedSettingsSave = saveSettingsAsync(paths, oldSettings); const blockedSessionSave = saveSessionAsync(paths, oldSession); @@ -2295,6 +2667,229 @@ describe("settings storage", () => { await nextBarrier.release({ replayBlocked: false }); }); + it("cleans up a barrier acquisition after an active write fails", async () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-barrier-acquire-failure-")); + tempDirs.push(dir); + const paths = createStoragePaths(dir); + let markWriteStarted = () => {}; + const writeStarted = new Promise((resolve) => { markWriteStarted = resolve; }); + let rejectWrite = (_error: Error) => {}; + const writeRelease = new Promise((_resolve, reject) => { rejectWrite = reject; }); + const writeFile = fs.promises.writeFile.bind(fs.promises); + vi.spyOn(fs.promises, "writeFile").mockImplementation((filePath, data, options) => { + if (String(filePath) === `${paths.configFile}.settings.tmp`) { + markWriteStarted(); + return writeRelease; + } + return writeFile(filePath, data, options); + }); + + const activeSave = saveSettingsAsync(paths, { ...defaultSettings(), packageName: "failing" }) + .then(() => "resolved", (error: unknown) => `rejected:${String((error as Error).message || error)}`); + await writeStarted; + const barrierAttempt = acquirePersistenceBarrier(); + const blockedSave = saveSessionAsync(paths, { ...emptySession(), summaryText: "blocked" }) + .then(() => "resolved", (error: unknown) => `rejected:${String((error as Error).message || error)}`); + rejectWrite(new Error("active write failed")); + + await expect(barrierAttempt).rejects.toThrow("active write failed"); + expect(await activeSave).toBe("rejected:active write failed"); + expect(await Promise.race([ + blockedSave, + new Promise((resolve) => setTimeout(() => resolve("timeout"), 100)) + ])).toBe("rejected:active write failed"); + vi.restoreAllMocks(); + + await saveSettingsAsync(paths, { ...defaultSettings(), packageName: "after-failure" }); + expect(loadSettings(paths).packageName).toBe("after-failure"); + const nextBarrier = await acquirePersistenceBarrier(); + await nextBarrier.release({ replayBlocked: false }); + }); + + it("requeues a pending settings save when barrier acquisition fails", async () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-barrier-settings-requeue-")); + tempDirs.push(dir); + const paths = createStoragePaths(dir); + const writeFile = fs.promises.writeFile.bind(fs.promises); + let markWriteStarted = () => {}; + const writeStarted = new Promise((resolve) => { markWriteStarted = resolve; }); + let rejectWrite = (_error: Error) => {}; + const writeRelease = new Promise((_resolve, reject) => { rejectWrite = reject; }); + let settingsWrites = 0; + vi.spyOn(fs.promises, "writeFile").mockImplementation((filePath, data, options) => { + if (String(filePath) === `${paths.configFile}.settings.tmp` && settingsWrites++ === 0) { + markWriteStarted(); + return writeRelease; + } + return writeFile(filePath, data, options); + }); + + const activeSave = saveSettingsAsync(paths, { ...defaultSettings(), packageName: "active" }) + .then(() => "resolved", (error: unknown) => `rejected:${String((error as Error).message || error)}`); + await writeStarted; + let queuedSettled = false; + const queuedSave = saveSettingsAsync(paths, { ...defaultSettings(), packageName: "queued-latest" }) + .then(() => "resolved", (error: unknown) => `rejected:${String((error as Error).message || error)}`) + .finally(() => { queuedSettled = true; }); + const barrierAttempt = acquirePersistenceBarrier(); + await new Promise((resolve) => setImmediate(resolve)); + const settledBeforeFailure = queuedSettled; + rejectWrite(new Error("active settings write failed")); + + await expect(barrierAttempt).rejects.toThrow("active settings write failed"); + expect(await activeSave).toBe("rejected:active settings write failed"); + expect(await queuedSave).toBe("resolved"); + expect(settledBeforeFailure).toBe(false); + expect(loadSettings(paths).packageName).toBe("queued-latest"); + }); + + it("requeues a pending session save when barrier acquisition fails", async () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-barrier-session-requeue-")); + tempDirs.push(dir); + const paths = createStoragePaths(dir); + const openFile = fs.promises.open.bind(fs.promises); + let markOpenStarted = () => {}; + const openStarted = new Promise((resolve) => { markOpenStarted = resolve; }); + let rejectOpen = (_error: Error) => {}; + const openRelease = new Promise((_resolve, reject) => { rejectOpen = reject; }); + let sessionOpens = 0; + vi.spyOn(fs.promises, "open").mockImplementation((filePath, flags, mode) => { + if (String(filePath) === `${paths.sessionFile}.async.tmp` && sessionOpens++ === 0) { + markOpenStarted(); + return openRelease; + } + return openFile(filePath, flags, mode); + }); + + const activeSave = saveSessionAsync(paths, { ...emptySession(), summaryText: "active" }) + .then(() => "resolved", (error: unknown) => `rejected:${String((error as Error).message || error)}`); + await openStarted; + let queuedSettled = false; + const queuedSave = saveSessionAsync(paths, { ...emptySession(), summaryText: "queued-latest" }) + .then(() => "resolved", (error: unknown) => `rejected:${String((error as Error).message || error)}`) + .finally(() => { queuedSettled = true; }); + const barrierAttempt = acquirePersistenceBarrier(); + await new Promise((resolve) => setImmediate(resolve)); + const settledBeforeFailure = queuedSettled; + rejectOpen(new Error("active session open failed")); + + await expect(barrierAttempt).rejects.toThrow("active session open failed"); + expect(await activeSave).toBe("rejected:active session open failed"); + expect(await queuedSave).toBe("resolved"); + expect(settledBeforeFailure).toBe(false); + expect(loadSession(paths).summaryText).toBe("queued-latest"); + }); + + it("rejects pre-barrier queued saves when a successful import barrier discards them", async () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-barrier-explicit-discard-")); + tempDirs.push(dir); + const paths = createStoragePaths(dir); + saveSettings(paths, defaultSettings()); + saveSession(paths, emptySession()); + const writeFile = fs.promises.writeFile.bind(fs.promises); + const openFile = fs.promises.open.bind(fs.promises); + let releaseSettingsWrite = () => {}; + let releaseSessionWrite = () => {}; + let markSettingsWriteStarted = () => {}; + let markSessionWriteStarted = () => {}; + const settingsWriteStarted = new Promise((resolve) => { markSettingsWriteStarted = resolve; }); + const sessionWriteStarted = new Promise((resolve) => { markSessionWriteStarted = resolve; }); + const settingsWriteRelease = new Promise((resolve) => { releaseSettingsWrite = resolve; }); + const sessionWriteRelease = new Promise((resolve) => { releaseSessionWrite = resolve; }); + vi.spyOn(fs.promises, "writeFile").mockImplementation(async (filePath, data, options) => { + await writeFile(filePath, data, options); + if (String(filePath) === `${paths.configFile}.settings.tmp`) { + markSettingsWriteStarted(); + await settingsWriteRelease; + } + }); + vi.spyOn(fs.promises, "open").mockImplementation(async (filePath, flags, mode) => { + const handle = await openFile(filePath, flags, mode); + if (String(filePath) === `${paths.sessionFile}.async.tmp`) { + const originalWrite = handle.writeFile.bind(handle); + handle.writeFile = async (...args: Parameters) => { + const result = await originalWrite(...args); + markSessionWriteStarted(); + await sessionWriteRelease; + return result; + }; + } + return handle; + }); + + const activeSettings = saveSettingsAsync(paths, { ...defaultSettings(), packageName: "active" }); + const activeSession = saveSessionAsync(paths, { ...emptySession(), summaryText: "active" }); + await Promise.all([settingsWriteStarted, sessionWriteStarted]); + const queuedSettings = saveSettingsAsync(paths, { ...defaultSettings(), packageName: "discarded" }) + .then(() => "resolved", (error: unknown) => `rejected:${String((error as Error).message || error)}`); + const queuedSession = saveSessionAsync(paths, { ...emptySession(), summaryText: "discarded" }) + .then(() => "resolved", (error: unknown) => `rejected:${String((error as Error).message || error)}`); + const barrierAttempt = acquirePersistenceBarrier(); + releaseSettingsWrite(); + releaseSessionWrite(); + const barrier = await barrierAttempt; + await Promise.all([activeSettings, activeSession]); + + expect(await queuedSettings).toMatch(/^rejected:.*barrier/i); + expect(await queuedSession).toMatch(/^rejected:.*barrier/i); + saveSettings(paths, { ...defaultSettings(), packageName: "imported" }); + saveSession(paths, { ...emptySession(), summaryText: "imported" }); + await barrier.release({ replayBlocked: false }); + expect(loadSettings(paths).packageName).toBe("imported"); + expect(loadSession(paths).summaryText).toBe("imported"); + }); + + it("waits for every started replay write before releasing a failed barrier", async () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-barrier-replay-failure-")); + tempDirs.push(dir); + const paths = createStoragePaths(dir); + const openFile = fs.promises.open.bind(fs.promises); + let markSessionWriteStarted = () => {}; + const sessionWriteStarted = new Promise((resolve) => { markSessionWriteStarted = resolve; }); + let releaseSessionWrite = () => {}; + const sessionWriteRelease = new Promise((resolve) => { releaseSessionWrite = resolve; }); + const writeFile = fs.promises.writeFile.bind(fs.promises); + vi.spyOn(fs.promises, "writeFile").mockImplementation((filePath, data, options) => { + if (String(filePath) === `${paths.configFile}.settings.tmp`) { + return Promise.reject(new Error("replay settings failed")); + } + return writeFile(filePath, data, options); + }); + vi.spyOn(fs.promises, "open").mockImplementation(async (filePath, flags, mode) => { + const handle = await openFile(filePath, flags, mode); + if (String(filePath) === `${paths.sessionFile}.async.tmp`) { + const writeFile = handle.writeFile.bind(handle); + handle.writeFile = async (...args: Parameters) => { + const result = await writeFile(...args); + markSessionWriteStarted(); + await sessionWriteRelease; + return result; + }; + } + return handle; + }); + + const barrier = await acquirePersistenceBarrier(); + const settingsSave = saveSettingsAsync(paths, { ...defaultSettings(), packageName: "replayed-settings" }) + .then(() => "resolved", (error: unknown) => `rejected:${String((error as Error).message || error)}`); + const sessionSave = saveSessionAsync(paths, { ...emptySession(), summaryText: "replayed-session" }) + .then(() => "resolved", (error: unknown) => `rejected:${String((error as Error).message || error)}`); + let releaseSettled = false; + const releaseResult = barrier.release({ replayBlocked: true }) + .then(() => "resolved", (error: unknown) => `rejected:${String((error as Error).message || error)}`) + .finally(() => { releaseSettled = true; }); + await sessionWriteStarted; + await new Promise((resolve) => setImmediate(resolve)); + + expect(releaseSettled).toBe(false); + releaseSessionWrite(); + expect(await releaseResult).toBe("rejected:replay settings failed"); + expect(await settingsSave).toBe("rejected:replay settings failed"); + expect(await sessionSave).toBe("resolved"); + const nextBarrier = await acquirePersistenceBarrier(); + await nextBarrier.release({ replayBlocked: false }); + }); + it("creates session backup before sync and async session overwrites", async () => { const dir = fs.mkdtempSync(path.join(os.tmpdir(), "rd-store-")); tempDirs.push(dir); diff --git a/tests/update.test.ts b/tests/update.test.ts index 51c75cb..0ca7e1c 100644 --- a/tests/update.test.ts +++ b/tests/update.test.ts @@ -1,6 +1,7 @@ -import fs from "node:fs"; -import crypto from "node:crypto"; -import { afterEach, describe, expect, it, vi } from "vitest"; +import fs from "node:fs"; +import crypto from "node:crypto"; +import { EventEmitter } from "node:events"; +import { afterEach, describe, expect, it, vi } from "vitest"; const { spawnMock, unrefMock, onceMock } = vi.hoisted(() => { const unref = vi.fn(); @@ -27,7 +28,7 @@ vi.mock("node:child_process", () => ({ spawn: spawnMock })); -import { buildInstallerLaunchArgs, checkGitHubUpdate, installLatestUpdate, isRemoteNewer, normalizeUpdateRepo, parseVersionParts } from "../src/main/update"; +import { abortActiveUpdateDownload, buildInstallerLaunchArgs, checkGitHubUpdate, installLatestUpdate, isRemoteNewer, normalizeUpdateRepo, parseVersionParts } from "../src/main/update"; import { APP_VERSION } from "../src/main/constants"; import { UpdateCheckResult, UpdateInstallProgress } from "../src/shared/types"; @@ -37,12 +38,19 @@ function sha256Hex(buffer: Buffer): string { return crypto.createHash("sha256").update(buffer).digest("hex"); } -function sha512Hex(buffer: Buffer): string { - return crypto.createHash("sha512").update(buffer).digest("hex"); -} - -afterEach(() => { - globalThis.fetch = originalFetch; +function sha512Hex(buffer: Buffer): string { + return crypto.createHash("sha512").update(buffer).digest("hex"); +} + +function createExecutablePayload() { + const payload = Buffer.alloc(128 * 1024); + payload.write("MZ"); + return payload; +} + +afterEach(() => { + abortActiveUpdateDownload(); + globalThis.fetch = originalFetch; spawnMock.mockClear(); unrefMock.mockClear(); onceMock.mockClear(); @@ -50,6 +58,416 @@ afterEach(() => { }); describe("update", () => { + it("settles a backpressured update write when the target stream errors", async () => { + class FailingWriteStream extends EventEmitter { + private markWriteStarted!: () => void; + readonly writeStarted = new Promise((resolve) => { + this.markWriteStarted = resolve; + }); + + write(): boolean { + this.markWriteStarted(); + return false; + } + + end(callback?: () => void): this { + callback?.(); + return this; + } + + destroy(): this { + return this; + } + } + + const stream = new FailingWriteStream(); + vi.spyOn(fs, "createWriteStream").mockReturnValue(stream as unknown as ReturnType); + globalThis.fetch = (async (): Promise => new Response(Buffer.alloc(160 * 1024, 1), { + status: 200, + headers: { "Content-Length": String(160 * 1024) } + })) as typeof fetch; + + const pending = installLatestUpdate("owner/repo", { + updateAvailable: true, + currentVersion: APP_VERSION, + latestVersion: "9.9.9", + latestTag: "", + releaseUrl: "https://example.invalid/release", + setupAssetUrl: "https://example.invalid/", + setupAssetName: "", + setupAssetDigest: `sha256:${"a".repeat(64)}` + }); + + await stream.writeStarted; + const writeError = Object.assign(new Error("disk full"), { code: "ENOSPC" }); + stream.emit("error", writeError); + const timeout = Symbol("timeout"); + const outcome = await Promise.race([ + pending, + new Promise((resolve) => setTimeout(() => resolve(timeout), 100)) + ]); + + if (outcome === timeout) { + stream.emit("drain"); + await pending; + } + + expect(outcome).not.toBe(timeout); + expect(outcome).toEqual(expect.objectContaining({ started: false })); + expect((outcome as { message: string }).message).toMatch(/disk full|ENOSPC/i); + + const next = await installLatestUpdate("owner/repo", { + updateAvailable: false, + currentVersion: APP_VERSION, + latestVersion: APP_VERSION, + latestTag: `v${APP_VERSION}`, + releaseUrl: "https://example.invalid/release" + }); + expect(next.message).not.toBe("Update-Download läuft bereits"); + }); + + it("settles shutdown abort while an update write waits for drain", async () => { + class BackpressuredWriteStream extends EventEmitter { + private markWriteStarted!: () => void; + private wasDestroyed = false; + readonly writeStarted = new Promise((resolve) => { + this.markWriteStarted = resolve; + }); + + write(): boolean { + this.markWriteStarted(); + return false; + } + + end(callback?: () => void): this { + callback?.(); + return this; + } + + destroy(): this { + this.wasDestroyed = true; + return this; + } + + emitLateDestroyError(error: Error): boolean { + if (!this.wasDestroyed) { + throw new Error("stream was not destroyed"); + } + const handled = this.listenerCount("error") > 0; + if (handled) { + this.emit("error", error); + } + this.emit("close"); + return handled; + } + } + + const stream = new BackpressuredWriteStream(); + vi.spyOn(fs, "createWriteStream").mockReturnValue(stream as unknown as ReturnType); + globalThis.fetch = (async (): Promise => new Response(Buffer.alloc(160 * 1024, 1), { + status: 200, + headers: { "Content-Length": String(160 * 1024) } + })) as typeof fetch; + + const pending = installLatestUpdate("owner/repo", { + updateAvailable: true, + currentVersion: APP_VERSION, + latestVersion: "9.9.9", + latestTag: "", + releaseUrl: "https://example.invalid/release", + setupAssetUrl: "https://example.invalid/", + setupAssetName: "", + setupAssetDigest: `sha256:${"a".repeat(64)}` + }); + + await stream.writeStarted; + abortActiveUpdateDownload(); + const timeout = Symbol("timeout"); + const outcome = await Promise.race([ + pending, + new Promise((resolve) => setTimeout(() => resolve(timeout), 100)) + ]); + + if (outcome === timeout) { + stream.emit("drain"); + await pending; + } + + expect(outcome).not.toBe(timeout); + expect(outcome).toEqual(expect.objectContaining({ started: false })); + expect((outcome as { message: string }).message).toMatch(/aborted:update_shutdown/i); + expect(stream.emitLateDestroyError(new Error("late destroy failure"))).toBe(true); + expect(stream.listenerCount("error")).toBe(0); + + const next = await installLatestUpdate("owner/repo", { + updateAvailable: false, + currentVersion: APP_VERSION, + latestVersion: APP_VERSION, + latestTag: `v${APP_VERSION}`, + releaseUrl: "https://example.invalid/release" + }); + expect(next.message).not.toBe("Update-Download läuft bereits"); + }); + + it("does not wait for a stalled body cancellation during shutdown", async () => { + class BackpressuredWriteStream extends EventEmitter { + private markWriteStarted!: () => void; + private markDestroyed!: () => void; + readonly writeStarted = new Promise((resolve) => { + this.markWriteStarted = resolve; + }); + readonly destroyed = new Promise((resolve) => { + this.markDestroyed = resolve; + }); + + write(): boolean { + this.markWriteStarted(); + return false; + } + + end(callback?: () => void): this { + callback?.(); + return this; + } + + destroy(): this { + this.markDestroyed(); + this.emit("close"); + return this; + } + } + + const stream = new BackpressuredWriteStream(); + const cancelState: { release?: () => void } = {}; + let markCancellationRequested!: () => void; + const cancellationRequested = new Promise((resolve) => { + markCancellationRequested = resolve; + }); + const responseBody = new ReadableStream({ + start(controller) { + controller.enqueue(new Uint8Array(160 * 1024)); + }, + cancel() { + markCancellationRequested(); + return new Promise((resolve) => { + cancelState.release = resolve; + }); + } + }); + vi.spyOn(fs, "createWriteStream").mockReturnValue(stream as unknown as ReturnType); + globalThis.fetch = (async (): Promise => new Response(responseBody, { + status: 200, + headers: { "Content-Length": String(160 * 1024) } + })) as typeof fetch; + + const pending = installLatestUpdate("owner/repo", { + updateAvailable: true, + currentVersion: APP_VERSION, + latestVersion: "9.9.9", + latestTag: "", + releaseUrl: "https://example.invalid/release", + setupAssetUrl: "https://example.invalid/", + setupAssetName: "", + setupAssetDigest: `sha256:${"a".repeat(64)}` + }); + + await stream.writeStarted; + abortActiveUpdateDownload(); + const timeout = Symbol("timeout"); + const outcome = await Promise.race([ + pending, + new Promise((resolve) => setTimeout(() => resolve(timeout), 100)) + ]); + const cancellationOutcome = await Promise.race([ + cancellationRequested.then(() => true), + new Promise((resolve) => setTimeout(() => resolve(false), 100)) + ]); + + if (outcome === timeout) { + cancelState.release?.(); + await pending; + } + cancelState.release?.(); + + expect(outcome).not.toBe(timeout); + expect(cancellationOutcome).toBe(true); + expect(outcome).toEqual(expect.objectContaining({ started: false })); + expect((outcome as { message: string }).message).toMatch(/aborted:update_shutdown/i); + await expect(stream.destroyed).resolves.toBeUndefined(); + }); + + it("propagates the original write error from the end callback", async () => { + const finalWriteError = Object.assign(new Error("final disk failure"), { code: "EIO" }); + + class FinalizingWriteStream extends EventEmitter { + write(): boolean { + return true; + } + + end(callback?: (error?: Error) => void): this { + callback?.(finalWriteError); + this.emit("error", finalWriteError); + this.emit("close"); + return this; + } + + destroy(): this { + this.emit("close"); + return this; + } + } + + const stream = new FinalizingWriteStream(); + vi.spyOn(fs, "createWriteStream").mockReturnValue(stream as unknown as ReturnType); + globalThis.fetch = (async (): Promise => new Response(Buffer.alloc(160 * 1024, 1), { + status: 200, + headers: { "Content-Length": String(160 * 1024) } + })) as typeof fetch; + + const outcome = await installLatestUpdate("owner/repo", { + updateAvailable: true, + currentVersion: APP_VERSION, + latestVersion: "9.9.9", + latestTag: "", + releaseUrl: "https://example.invalid/release", + setupAssetUrl: "https://example.invalid/", + setupAssetName: "", + setupAssetDigest: `sha256:${"a".repeat(64)}` + }); + + expect(outcome).toEqual(expect.objectContaining({ started: false })); + expect(outcome.message).toMatch(/final disk failure|EIO/i); + }); + + it("waits for close and propagates a late close error after a successful end callback", async () => { + const lateCloseError = Object.assign(new Error("late close failure EIO"), { code: "EIO" }); + + class LateCloseFailureStream extends EventEmitter { + private closed = false; + + write(): boolean { + return true; + } + + end(callback?: (error?: Error) => void): this { + callback?.(); + this.emit("finish"); + queueMicrotask(() => { + this.emit("error", lateCloseError); + this.closed = true; + this.emit("close"); + }); + return this; + } + + destroy(): this { + if (!this.closed) { + this.closed = true; + this.emit("close"); + } + return this; + } + } + + const stream = new LateCloseFailureStream(); + const renameSpy = vi.spyOn(fs.promises, "rename").mockRejectedValue(new Error("rename must not run")); + vi.spyOn(fs, "createWriteStream").mockReturnValue(stream as unknown as ReturnType); + globalThis.fetch = (async (): Promise => new Response(Buffer.alloc(160 * 1024, 1), { + status: 200, + headers: { "Content-Length": String(160 * 1024) } + })) as typeof fetch; + + const outcome = await installLatestUpdate("owner/repo", { + updateAvailable: true, + currentVersion: APP_VERSION, + latestVersion: "9.9.9", + latestTag: "", + releaseUrl: "https://example.invalid/release", + setupAssetUrl: "https://example.invalid/", + setupAssetName: "", + setupAssetDigest: `sha256:${"a".repeat(64)}` + }); + + expect(outcome).toEqual(expect.objectContaining({ started: false })); + expect(outcome.message).toMatch(/late close failure|EIO/i); + expect(renameSpy).not.toHaveBeenCalled(); + expect(stream.listenerCount("error")).toBe(0); + }); + + it("settles an idle timeout while an update write waits for drain", async () => { + const previousTimeout = process.env.RD_UPDATE_BODY_IDLE_TIMEOUT_MS; + process.env.RD_UPDATE_BODY_IDLE_TIMEOUT_MS = "1000"; + + class BackpressuredWriteStream extends EventEmitter { + private markWriteStarted!: () => void; + private markDestroyed!: () => void; + readonly writeStarted = new Promise((resolve) => { + this.markWriteStarted = resolve; + }); + readonly destroyed = new Promise((resolve) => { + this.markDestroyed = resolve; + }); + + write(): boolean { + this.markWriteStarted(); + return false; + } + + end(callback?: () => void): this { + callback?.(); + return this; + } + + destroy(): this { + this.markDestroyed(); + this.emit("close"); + return this; + } + } + + const stream = new BackpressuredWriteStream(); + const externalErrorListener = (): void => undefined; + stream.on("error", externalErrorListener); + vi.spyOn(fs, "createWriteStream").mockReturnValue(stream as unknown as ReturnType); + globalThis.fetch = (async (): Promise => new Response(Buffer.alloc(160 * 1024, 1), { + status: 200, + headers: { "Content-Length": String(160 * 1024) } + })) as typeof fetch; + + try { + const pending = installLatestUpdate("owner/repo", { + updateAvailable: true, + currentVersion: APP_VERSION, + latestVersion: "9.9.9", + latestTag: "", + releaseUrl: "https://example.invalid/release", + setupAssetUrl: "https://example.invalid/", + setupAssetName: "", + setupAssetDigest: `sha256:${"a".repeat(64)}` + }); + + await stream.writeStarted; + const timeout = Symbol("timeout"); + const idleOutcome = await Promise.race([ + stream.destroyed.then(() => "destroyed" as const), + new Promise((resolve) => setTimeout(() => resolve(timeout), 1600)) + ]); + abortActiveUpdateDownload(); + const outcome = await pending; + + expect(idleOutcome).toBe("destroyed"); + expect(outcome).toEqual(expect.objectContaining({ started: false })); + expect(stream.listenerCount("drain")).toBe(0); + expect(stream.listeners("error")).toEqual([externalErrorListener]); + } finally { + if (previousTimeout === undefined) { + delete process.env.RD_UPDATE_BODY_IDLE_TIMEOUT_MS; + } else { + process.env.RD_UPDATE_BODY_IDLE_TIMEOUT_MS = previousTimeout; + } + } + }); + it("always refreshes release metadata before installing instead of using the previous check result", () => { const controller = fs.readFileSync(new URL("../src/main/app-controller.ts", import.meta.url), "utf8"); const install = controller.slice(controller.indexOf("public async installUpdate"), controller.indexOf("public addLinks")); @@ -172,8 +590,8 @@ describe("update", () => { expect(buildInstallerLaunchArgs()).toEqual(["/S", "--updated", "--force-run"]); }); - it("falls back to alternate download URL when setup asset URL returns 404", async () => { - const executablePayload = fs.readFileSync(process.execPath); + it("falls back to alternate download URL when setup asset URL returns 404", async () => { + const executablePayload = createExecutablePayload(); const executableDigest = sha256Hex(executablePayload); const requestedUrls: string[] = []; globalThis.fetch = (async (input: RequestInfo | URL): Promise => { @@ -209,8 +627,8 @@ describe("update", () => { expect(requestedUrls.filter((url) => url.includes("stale-setup.exe"))).toHaveLength(1); }); - it("skips draft tag payload and resolves setup asset from stable latest release", async () => { - const executablePayload = fs.readFileSync(process.execPath); + it("skips draft tag payload and resolves setup asset from stable latest release", async () => { + const executablePayload = createExecutablePayload(); const requestedUrls: string[] = []; globalThis.fetch = (async (input: RequestInfo | URL): Promise => { @@ -350,8 +768,8 @@ describe("update", () => { } }, 20000); - it("blocks installer start when SHA256 digest mismatches", async () => { - const executablePayload = fs.readFileSync(process.execPath); + it("blocks installer start when SHA256 digest mismatches", async () => { + const executablePayload = createExecutablePayload(); globalThis.fetch = (async (input: RequestInfo | URL): Promise => { const url = typeof input === "string" ? input : input instanceof URL ? input.toString() : input.url; if (url.includes("mismatch-setup.exe")) { @@ -379,8 +797,8 @@ describe("update", () => { expect(result.message).toMatch(/integrit|sha256|mismatch/i); }); - it("blocks installer start when no digest can be resolved", async () => { - const executablePayload = fs.readFileSync(process.execPath); + it("blocks installer start when no digest can be resolved", async () => { + const executablePayload = createExecutablePayload(); globalThis.fetch = (async (input: RequestInfo | URL): Promise => { const url = typeof input === "string" ? input : input instanceof URL ? input.toString() : input.url; if (url.includes("unsigned-setup.exe")) { @@ -408,8 +826,8 @@ describe("update", () => { expect(result.message).toMatch(/digest|integrit|sha/i); }); - it("uses latest.yml SHA512 digest when API asset digest is missing", async () => { - const executablePayload = fs.readFileSync(process.execPath); + it("uses latest.yml SHA512 digest when API asset digest is missing", async () => { + const executablePayload = createExecutablePayload(); const digestSha512Hex = sha512Hex(executablePayload); const digestSha512Base64 = Buffer.from(digestSha512Hex, "hex").toString("base64"); const requestedUrls: string[] = []; @@ -479,8 +897,8 @@ describe("update", () => { expect(requestedUrls.some((url) => url.includes("latest.yml"))).toBe(true); }); - it("rejects installer when latest.yml SHA512 digest does not match", async () => { - const executablePayload = fs.readFileSync(process.execPath); + it("rejects installer when latest.yml SHA512 digest does not match", async () => { + const executablePayload = createExecutablePayload(); const wrongDigestBase64 = Buffer.alloc(64, 0x13).toString("base64"); globalThis.fetch = (async (input: RequestInfo | URL): Promise => { @@ -547,7 +965,7 @@ describe("update", () => { }); it("emits install progress events while downloading and launching update", async () => { - const executablePayload = fs.readFileSync(process.execPath); + const executablePayload = createExecutablePayload(); const digest = sha256Hex(executablePayload); globalThis.fetch = (async (input: RequestInfo | URL): Promise => { @@ -595,7 +1013,7 @@ describe("update", () => { }); it("keeps the application running when Windows rejects the installer process", async () => { - const executablePayload = fs.readFileSync(process.execPath); + const executablePayload = createExecutablePayload(); const digest = sha256Hex(executablePayload); globalThis.fetch = (async (): Promise => new Response(executablePayload, { status: 200,