From 5a369f75984bc0d65511ef14e366554403fb336f Mon Sep 17 00:00:00 2001 From: Sucukdeluxe <259325684+Sucukdeluxe@users.noreply.github.com> Date: Thu, 27 Aug 2026 07:15:39 +0200 Subject: [PATCH] fix: eliminate upload queue stalls under automation load Move atomic configuration writes off the Electron main event loop and preserve concurrent queue and folder monitor state. Make finish-and-pause interrupt queued admission waits without aborting active uploads, then resume only persistently marked automation jobs across active, idle, restart, and finalization races. Reuse bounded automation evidence snapshots, refresh them during long drains, and cover high-load persistence, pause, resume, and hidden Electron behavior with regression tests. --- lib/config-store.js | 75 +++++----- lib/upload-manager.js | 66 ++++++++- main.js | 45 +++--- renderer/app.js | 79 +++++++++-- renderer/i18n.js | 1 + tests/config-store.test.js | 90 ++++++++++++ tests/i18n.test.js | 1 + tests/package-build-files.test.js | 140 +++++++++++++++++-- tests/startup-renderer.test.js | 199 +++++++++++++++++++++++++- tests/upload-manager.test.js | 225 ++++++++++++++++++++++++++++++ 10 files changed, 843 insertions(+), 78 deletions(-) diff --git a/lib/config-store.js b/lib/config-store.js index cd90743..bec4f33 100644 --- a/lib/config-store.js +++ b/lib/config-store.js @@ -589,6 +589,26 @@ class ConfigStore { }, options); } + saveFolderMonitorRuntimeState(folderMonitor, options = {}) { + const snapshot = { + paused: folderMonitor?.paused === true, + pausedAt: folderMonitor?.paused === true ? (folderMonitor?.pausedAt ?? null) : null + }; + return this._enqueueWrite(() => { + const current = this.load(); + const currentGlobalSettings = current.globalSettings || {}; + current.globalSettings = { + ...currentGlobalSettings, + folderMonitor: { + ...(currentGlobalSettings.folderMonitor || {}), + ...snapshot + } + }; + this._guardHosters(current, false); + return this._commit(current); + }, options); + } + saveUploadRecovery(uploadRecovery, options = {}) { const snapshot = uploadRecovery === null || uploadRecovery === undefined ? null : this._clone(uploadRecovery); return this._enqueueWrite(() => { @@ -686,38 +706,29 @@ class ConfigStore { return config.history || []; } - _atomicWrite(data) { - return new Promise((resolve, reject) => { - const tmpPath = this.filePath + '.tmp'; - const backupPath = this.filePath + '.bak'; - let fd; - try { - fd = fs.openSync(tmpPath, 'w'); - fs.writeSync(fd, data); - fs.fsyncSync(fd); - } catch (e) { - try { if (fd !== undefined) fs.closeSync(fd); } catch {} - return reject(e); - } - try { fs.closeSync(fd); } catch {} - Promise.resolve().then(() => { - try { - try { - if (fs.existsSync(this.filePath)) { - const cur = fs.readFileSync(this.filePath, 'utf-8'); - if (cur && cur.trim().length > 2) fs.writeFileSync(backupPath, cur, 'utf-8'); - } - } catch {} - fs.renameSync(tmpPath, this.filePath); - } catch (e) { return reject(e); } - // Invalidate the read cache: the next load() re-reads + re-merges the - // freshly-written file (the on-disk format is sparse — load() fills - // defaults — so we must NOT serve a pre-merge in-memory object). - this._cache = null; - this._cacheKey = ''; - resolve(); - }); - }); + async _atomicWrite(data) { + const tmpPath = this.filePath + '.tmp'; + const backupPath = this.filePath + '.bak'; + const fileHandle = await fs.promises.open(tmpPath, 'w'); + let writeError = null; + try { + await fileHandle.writeFile(data); + await fileHandle.sync(); + } catch (error) { + writeError = error; + } + try { await fileHandle.close(); } catch {} + if (writeError) throw writeError; + try { + const current = await fs.promises.readFile(this.filePath, 'utf-8'); + if (current && current.trim().length > 2) await fs.promises.writeFile(backupPath, current, 'utf-8'); + } catch {} + await fs.promises.rename(tmpPath, this.filePath); + // Invalidate the read cache: the next load() re-reads + re-merges the + // freshly-written file (the on-disk format is sparse — load() fills + // defaults — so we must NOT serve a pre-merge in-memory object). + this._cache = null; + this._cacheKey = ''; } appendHistory(entry) { diff --git a/lib/upload-manager.js b/lib/upload-manager.js index 165befe..1e86f8e 100644 --- a/lib/upload-manager.js +++ b/lib/upload-manager.js @@ -33,6 +33,9 @@ class UploadManager extends EventEmitter { this.semaphores = {}; this.globalSemaphore = null; this.abortController = new AbortController(); + this.queueAdmissionAbortController = new AbortController(); + this.queueAdmissionPending = 0; + this.queueAdmissionWaiters = []; this.running = false; this.stopAfterActive = false; this.statsInterval = null; @@ -344,6 +347,9 @@ class UploadManager extends EventEmitter { this.running = true; this.stopAfterActive = pendingCancelAll; this.abortController = new AbortController(); + this.queueAdmissionAbortController = new AbortController(); + this.queueAdmissionPending = 0; + this.queueAdmissionWaiters = []; if (pendingCancelAll) this.abortController.abort(); this.startTime = Date.now(); this.sessionBytes = 0; @@ -462,6 +468,11 @@ class UploadManager extends EventEmitter { const jobAbortController = new AbortController(); if (this.cancelledJobIds.has(jobId)) jobAbortController.abort(); const { signal, cleanup: cleanupSignals } = this._combineSignals(batchSignal, jobAbortController.signal); + const { signal: admissionSignal, cleanup: cleanupAdmissionSignals } = this._combineSignals( + signal, + this.queueAdmissionAbortController.signal + ); + const leaveQueueAdmission = this._enterQueueAdmission(); this.jobAbortControllers.set(jobId, jobAbortController); let hosterSlotAcquired = false; @@ -546,9 +557,16 @@ class UploadManager extends EventEmitter { // queueJobs array; the first event it actually needs from main is the // 'getting-server' / 'uploading' transition for the jobs that the // semaphore lets through. - await hosterSemaphore.acquire(signal); + await hosterSemaphore.acquire(admissionSignal); hosterSlotAcquired = true; + if (this.stopAfterActive) { + const error = 'Warteschlange angehalten'; + emitFinalStatus('aborted', { error, attempt: 0 }); + recordFinalResult('aborted', { error }); + return; + } + let fileProbe = null; try { fileProbe = await probeFileHead(task.file, 512); @@ -564,14 +582,20 @@ class UploadManager extends EventEmitter { }); if (globalSemaphore) { - await globalSemaphore.acquire(signal); + await globalSemaphore.acquire(admissionSignal); globalSlotAcquired = true; } + if (this.stopAfterActive) throw new Error('Warteschlange angehalten'); + if (settings.timeIntervalSec > 0) { - await this._waitForInterval(task.hoster, settings.timeIntervalSec * 1000, signal); + await this._waitForInterval(task.hoster, settings.timeIntervalSec * 1000, admissionSignal); } + if (this.stopAfterActive) throw new Error('Warteschlange angehalten'); + + leaveQueueAdmission(); + // Pre-job-swap: if this account was marked failed WHILE this task was // waiting in the semaphore queue, jump straight to the override instead // of burning a guaranteed-to-fail upload attempt. Critical at scale: @@ -1115,6 +1139,8 @@ class UploadManager extends EventEmitter { this.activeJobs.delete(uploadId); this.jobAbortControllers.delete(jobId); cleanupSignals(); + cleanupAdmissionSignals(); + leaveQueueAdmission(); // Release in reverse order of acquire (global first, then hoster) if (globalSlotAcquired && globalSemaphore) globalSemaphore.release(); if (hosterSlotAcquired) hosterSemaphore.release(); @@ -1467,7 +1493,7 @@ class UploadManager extends EventEmitter { } addJobs(tasks) { - if (!this.running || !tasks || tasks.length === 0) { + if (!this.running || this.stopAfterActive || !tasks || tasks.length === 0) { return { added: 0, alreadyInBatchJobIds: [] }; } const { signal } = this.abortController; @@ -1505,6 +1531,38 @@ class UploadManager extends EventEmitter { finishAfterActive() { this.stopAfterActive = true; + if (!this.queueAdmissionAbortController.signal.aborted) this.queueAdmissionAbortController.abort(); + } + + async resumeAfterActive() { + const stoppedController = this.queueAdmissionAbortController; + await this._waitForQueueAdmissionIdle(); + if (this.queueAdmissionAbortController === stoppedController && stoppedController.signal.aborted) { + this.queueAdmissionAbortController = new AbortController(); + } + this.stopAfterActive = false; + } + + isStoppingAfterActive() { + return this.stopAfterActive; + } + + _enterQueueAdmission() { + this.queueAdmissionPending++; + let active = true; + return () => { + if (!active) return; + active = false; + this.queueAdmissionPending--; + if (this.queueAdmissionPending !== 0) return; + const waiters = this.queueAdmissionWaiters.splice(0); + for (const resolve of waiters) resolve(); + }; + } + + _waitForQueueAdmissionIdle() { + if (this.queueAdmissionPending === 0) return Promise.resolve(); + return new Promise(resolve => this.queueAdmissionWaiters.push(resolve)); } cancel() { diff --git a/main.js b/main.js index c0e099e..88e8c64 100644 --- a/main.js +++ b/main.js @@ -144,6 +144,15 @@ let lastSessionSummary = null; let sourceDeleteJournal = null; const pendingUploadFinalizations = new Map(); +async function waitForUploadManagerRelease(manager, timeoutMs = 300000) { + const deadline = Date.now() + timeoutMs; + for (;;) { + if (uploadManager !== manager) return true; + if (Date.now() >= deadline) return false; + await new Promise(resolve => setTimeout(resolve, 25)); + } +} + function requestUploadFinalization(summary) { const finalizationId = `upload-${Date.now()}-${Math.random().toString(36).slice(2, 10)}`; return new Promise((resolve) => { @@ -2220,7 +2229,12 @@ ipcMain.handle('start-upload', async (_event, payload) => { } if (closeFlushRequested) return { error: 'Die Anwendung wird gerade beendet' }; if (!settingsImportGate.canStartUpload()) return { error: 'Einstellungen werden gerade importiert' }; - if (uploadManager) return { error: 'Ein Upload wird bereits ausgeführt oder abgeschlossen' }; + if (uploadManager) { + const existingManager = uploadManager; + if (existingManager.running || !(await waitForUploadManagerRelease(existingManager)) || uploadManager) { + return { error: 'Ein Upload wird bereits ausgeführt oder abgeschlossen' }; + } + } const config = configStore.load(); const files = payload && Array.isArray(payload.files) ? payload.files : []; const hosters = payload && Array.isArray(payload.hosters) ? payload.hosters : []; @@ -2566,6 +2580,9 @@ ipcMain.handle('add-jobs-to-batch', async (_event, payload) => { return { error: 'Kein Upload aktiv' }; } const batchManager = uploadManager; + if (batchManager.isStoppingAfterActive()) { + return { error: 'Warteschlange angehalten' }; + } const config = configStore.load(); const jobs = payload && Array.isArray(payload.jobs) ? payload.jobs : []; const sourceCleanupGroups = payload && Array.isArray(payload.sourceCleanupGroups) ? payload.sourceCleanupGroups : []; @@ -2585,6 +2602,9 @@ ipcMain.handle('add-jobs-to-batch', async (_event, payload) => { if (uploadManager !== batchManager || !batchManager.running) { return { error: 'Kein Upload aktiv' }; } + if (batchManager.isStoppingAfterActive()) { + return { error: 'Warteschlange angehalten' }; + } if (batchManager.sourceFileCleanup) { for (const skipped of skippedJobs) batchManager.sourceFileCleanup.markSkipped(skipped.jobId); } @@ -3439,12 +3459,7 @@ ipcMain.handle('automation:pause-after-active', () => enqueueAutomationLifecycle await withAutomationStatusSuppressed(async () => { const latest = configStore.load(); const settings = latest.globalSettings?.folderMonitor || {}; - await configStore.save({ - globalSettings: { - ...latest.globalSettings, - folderMonitor: { ...settings, paused: true, pausedAt: Date.now() } - } - }); + await configStore.saveFolderMonitorRuntimeState({ ...settings, paused: true, pausedAt: Date.now() }); invalidateFolderMonitorLifecycle(); try { await folderMonitor.pause(); @@ -3475,28 +3490,20 @@ ipcMain.handle('automation:resume', () => enqueueAutomationLifecycle(async gener if (resumedSettings.enabled && resumedSettings.folderPath) { await resumeFolderMonitor(resumedSettings); } - await configStore.save({ - globalSettings: { - ...latest.globalSettings, - folderMonitor: resumedSettings - } - }); + await configStore.saveFolderMonitorRuntimeState(resumedSettings); + if (uploadManager) await uploadManager.resumeAfterActive(); if (resumedSettings.enabled && resumedSettings.folderPath) { await folderMonitor.scan({ emitFiles: true, trigger: 'resume' }); } } catch { + if (uploadManager) uploadManager.finishAfterActive(); stopFolderMonitor(); if (pausedSettings.enabled && pausedSettings.folderPath) { bindFolderMonitorEvents(pausedSettings); folderMonitor.configure(pausedSettings); } try { - await configStore.save({ - globalSettings: { - ...latest.globalSettings, - folderMonitor: pausedSettings - } - }); + await configStore.saveFolderMonitorRuntimeState(pausedSettings); } catch {} result = { error: 'Automatik konnte nicht fortgesetzt werden' }; } diff --git a/renderer/app.js b/renderer/app.js index c989cbb..1b0af66 100644 --- a/renderer/app.js +++ b/renderer/app.js @@ -73,9 +73,12 @@ let automationTestReturnFocus = null; let automationTestInertState = []; let automationTestViewState = Object.freeze({ loading: false, summary: null, error: '' }); const automationEventBatchSize = 8; +const automationEvidenceReuseMs = 5000; const automationEventQueue = new Map(); const automationEventInFlight = new Set(); let automationEventDrainPromise = null; +let automationEvidenceSnapshotCache = null; +let automationEvidenceSnapshotGeneration = 0; let managedOnlineBackups = []; let managedOnlineBackupsAuthoritative = false; let managedOnlineBackupMutationGeneration = 0; @@ -552,6 +555,31 @@ function createAutomationStatusSnapshot() { return freezeAutomationValue(snapshot); } +async function loadAutomationEvidenceSnapshot() { + const [history, uploadLog] = await Promise.all([ + window.api.getHistory(), + window.api.readOwnUploadLog() + ]); + return { history, uploadLog }; +} + +function invalidateAutomationEvidenceSnapshot() { + automationEvidenceSnapshotGeneration++; + automationEvidenceSnapshotCache = null; +} + +async function loadReusableAutomationEvidenceSnapshot() { + const now = performance.now(); + if (automationEvidenceSnapshotCache?.expiresAt > now) return automationEvidenceSnapshotCache.value; + while (true) { + const generation = automationEvidenceSnapshotGeneration; + const value = await loadAutomationEvidenceSnapshot(); + if (generation !== automationEvidenceSnapshotGeneration) continue; + automationEvidenceSnapshotCache = { value, expiresAt: performance.now() + automationEvidenceReuseMs }; + return value; + } +} + async function evaluateAutomationCandidates(files, options = {}) { const source = Array.isArray(files) ? files : []; const normalizedCandidates = source.map(normalizeAutomationCandidate); @@ -570,10 +598,7 @@ async function evaluateAutomationCandidates(files, options = {}) { const selectedHosters = Array.from(new Set((Array.isArray(options.selectedHosters) ? options.selectedHosters : folderSettings.hosters || []) .map(value => String(value || '').trim()) .filter(Boolean))); - const [history, uploadLog] = await Promise.all([ - window.api.getHistory(), - window.api.readOwnUploadLog() - ]); + const { history, uploadLog } = options.evidenceSnapshot || await loadAutomationEvidenceSnapshot(); const processed = window.AutomationControl.classifyProcessedCandidates({ candidates: matched, queuePaths: [...queueJobs.map(job => job.file), ...selectedFiles.map(file => file.path)], @@ -1171,6 +1196,11 @@ async function toggleAutomationPauseResume() { if (automationPauseResumeBusy) return; const snapshot = createAutomationStatusSnapshot(); const resume = snapshot.paused === true; + const pausingJobIds = resume + ? null + : new Set(queueJobs + .filter(job => ['queued', 'getting-server', 'uploading', 'retrying'].includes(job.status)) + .map(job => job.id)); automationPauseResumeBusy = true; updateQueueActionButtons(snapshot); try { @@ -1179,9 +1209,27 @@ async function toggleAutomationPauseResume() { : await window.api.automationPauseAfterActive(); if (result?.error) throw new Error(result.error); applyAutomationRuntimeStatus({ ...result, paused: resume ? false : true }); - if (!resume && uploading) { - lastUploadStats.state = 'stopping'; - updateStatusBar(); + if (resume) { + const resumableJobs = queueJobs.filter(job => job.automationPaused === true && ( + job.status === 'queued' + || (job.status === 'aborted' && job.error === 'Warteschlange angehalten') + )); + if (resumableJobs.length > 0) { + if (uploading) await startSelectedUpload(resumableJobs); + else await startSelectedUpload(resumableJobs); + } + } else { + for (const job of queueJobs) { + if (pausingJobIds.has(job.id) && ['queued', 'getting-server', 'uploading', 'retrying'].includes(job.status)) { + job.automationPaused = true; + } + } + queuePersistThrottle.cancel(); + await persistQueueStateNow(); + if (uploading) { + lastUploadStats.state = 'stopping'; + updateStatusBar(); + } } } catch { showCopyToast(localizeUiText(resume ? 'Automatik konnte nicht fortgesetzt werden.' : 'Automatik konnte nicht pausiert werden.')); @@ -1204,9 +1252,9 @@ function surfaceAutomationOutcome(result) { return localized; } -async function processFolderMonitorFiles(files) { +async function processFolderMonitorFiles(files, evidenceSnapshot) { window.api.debugLog('folder-monitor: received ' + files.length + ' file(s)'); - const evaluation = await evaluateAutomationCandidates(files, { dryRun: false, trigger: 'watcher' }); + const evaluation = await evaluateAutomationCandidates(files, { dryRun: false, trigger: 'watcher', evidenceSnapshot }); const result = await applyAutomationEvaluation(evaluation); surfaceAutomationOutcome(result); return result; @@ -1215,13 +1263,14 @@ async function processFolderMonitorFiles(files) { async function drainFolderMonitorFiles() { let result = freezeAutomationValue({ admittedFiles: [], deferredFiles: [], paused: false, dryRun: false }); while (automationEventQueue.size > 0) { + const evidenceSnapshot = await loadReusableAutomationEvidenceSnapshot(); const entries = [...automationEventQueue.entries()].slice(0, automationEventBatchSize); for (const [key] of entries) { automationEventQueue.delete(key); automationEventInFlight.add(key); } try { - result = await processFolderMonitorFiles(entries.map(([, file]) => file)); + result = await processFolderMonitorFiles(entries.map(([, file]) => file), evidenceSnapshot); } finally { for (const [key] of entries) automationEventInFlight.delete(key); } @@ -2124,6 +2173,7 @@ function restoreQueueStateFromConfig() { sourceCleanupCompletedHosters: Array.isArray(job.sourceCleanupCompletedHosters) ? [...job.sourceCleanupCompletedHosters] : [], sourceCleanupFingerprint: job.sourceCleanupFingerprint || null, automationAdmission: job.automationAdmission === true, + ...(job.automationPaused === true ? { automationPaused: true } : {}), attempt: 0, maxAttempts: job.maxAttempts || 0, link: '', @@ -2187,12 +2237,13 @@ function buildPersistedQueueState() { suppressedKeys, queueJobs: queueJobs.map(job => { const isTerminal = TERMINAL.has(job.status); + const automationPaused = job.automationPaused === true; return { id: job.id, file: job.file, fileName: job.fileName, hoster: job.hoster, - status: isTerminal ? job.status : 'preview', + status: automationPaused ? 'queued' : (isTerminal ? job.status : 'preview'), bytesTotal: job.bytesTotal || 0, error: isTerminal ? (job.error || null) : null, failureDetails: isTerminal ? (job.failureDetails || null) : null, @@ -2202,6 +2253,7 @@ function buildPersistedQueueState() { sourceCleanupCompletedHosters: Array.isArray(job.sourceCleanupCompletedHosters) ? [...job.sourceCleanupCompletedHosters] : [], sourceCleanupFingerprint: job.sourceCleanupFingerprint || null, automationAdmission: job.automationAdmission === true, + ...(automationPaused ? { automationPaused: true } : {}), maxAttempts: job.maxAttempts || 0 }; }) @@ -4640,6 +4692,7 @@ async function cancelUpload() { uploading = false; // Reset all non-finished jobs back to queued state for (const job of queueJobs) { + delete job.automationPaused; if (!['done', 'error', 'skipped'].includes(job.status)) { job.status = 'queued'; job.progress = 0; @@ -4701,6 +4754,8 @@ function _handleProgressImpl(data) { // Update job state job.status = data.status; + if (data.status === 'aborted' && data.error === 'Warteschlange angehalten') job.automationPaused = true; + else if (data.status !== 'queued') delete job.automationPaused; if (data.status !== 'preview') job.interrupted = false; job.bytesUploaded = data.bytesUploaded || 0; job.bytesTotal = data.bytesTotal || job.bytesTotal; @@ -4774,6 +4829,7 @@ function _handleProgressImpl(data) { } function handleBatchDone(summary) { + invalidateAutomationEvidenceSnapshot(); uploading = false; applySummaryResults(summary); _deletedJobIds.clear(); // Free memory — stale IDs no longer needed after batch completes @@ -8345,6 +8401,7 @@ async function confirmHistoryClear() { cancelButton.disabled = true; try { await runConfigWrite(() => window.api.clearHistory()); + invalidateAutomationEvidenceSnapshot(); await loadHistory(); closeHistoryClearModal(); } catch (error) { diff --git a/renderer/i18n.js b/renderer/i18n.js index d3c06b8..647c6e7 100644 --- a/renderer/i18n.js +++ b/renderer/i18n.js @@ -656,6 +656,7 @@ ['Download hängt — seit 45 s keine Daten (Netzwerk/Server überlastet). Bitte laufende Uploads stoppen und erneut versuchen.', 'The download stalled because no data was received for 45 seconds. Stop active uploads and try again.'], ['Datei nicht gefunden', 'File not found'], ['Netzwerkfehler', 'Network error'], + ['Warteschlange angehalten', 'Queue paused'], ['Bekanntes Größen-Limit auf diesem Account (frühere verdächtige Ablehnung)', 'Known size limit on this account (previous suspicious rejection)'], ['Ablehnung verdächtig - Versuch auf anderem Account', 'Suspicious rejection - trying another account'], ['Vidmoly: /api/upload/config lieferte kein JSON — evtl. nicht eingeloggt?', 'Vidmoly: /api/upload/config did not return JSON — you may not be signed in'], diff --git a/tests/config-store.test.js b/tests/config-store.test.js index 734b39c..74001e8 100644 --- a/tests/config-store.test.js +++ b/tests/config-store.test.js @@ -355,6 +355,96 @@ describe('ConfigStore', () => { assert.equal(config.globalSettings.alwaysOnTop, true); }); + it('savePendingQueue does not block the event loop on slow synchronous filesystem methods', async () => { + await store.save({ + hosters: { 'byse.sx': [{ id: 'non-blocking-account', enabled: true, authType: 'api', apiKey: 'test-key' }] }, + globalSettings: { alwaysOnTop: true } + }); + store.load(); + + const syncMethods = ['openSync', 'writeSync', 'fsyncSync', 'readFileSync', 'writeFileSync', 'renameSync']; + const originals = new Map(syncMethods.map(name => [name, fs[name]])); + const waitBuffer = new Int32Array(new SharedArrayBuffer(4)); + const stallMs = 35; + let eventLoopDelayMs; + + for (const name of syncMethods) { + fs[name] = (...args) => { + Atomics.wait(waitBuffer, 0, 0, stallMs); + return originals.get(name)(...args); + }; + } + + try { + const startedAt = performance.now(); + const eventLoopTick = new Promise(resolve => { + setTimeout(() => { + eventLoopDelayMs = performance.now() - startedAt; + resolve(); + }, 0); + }); + const save = store.savePendingQueue({ savedAt: 5, queueJobs: [{ id: 'non-blocking' }] }); + + await eventLoopTick; + await save; + } finally { + for (const [name, original] of originals) fs[name] = original; + } + + assert.ok( + eventLoopDelayMs < stallMs * 3, + `queue save blocked the event loop for ${eventLoopDelayMs.toFixed(1)} ms` + ); + }); + + it('folder monitor runtime saves preserve a concurrently queued pending queue snapshot', async () => { + const pendingQueue = { + savedAt: 1787712000000, + queueJobs: [{ id: 'paused-job', automationPaused: true }] + }; + + const queueSave = store.savePendingQueue(pendingQueue); + const runtimeSave = store.saveFolderMonitorRuntimeState({ paused: false, pausedAt: null }); + await Promise.all([queueSave, runtimeSave]); + + const current = store.load().globalSettings; + assert.deepEqual(current.pendingQueue, pendingQueue); + assert.equal(current.folderMonitor.paused, false); + assert.equal(current.folderMonitor.pausedAt, null); + }); + + it('folder monitor runtime saves cannot revert concurrently queued monitor settings', async () => { + const stale = store.load().globalSettings.folderMonitor; + const currentSettings = { + ...stale, + folderPath: 'D:\\new-watch', + hosters: ['byse.sx'], + filterMode: 'exclude', + autoStart: false + }; + + const settingsSave = store.save({ + globalSettings: { + ...store.load().globalSettings, + folderMonitor: currentSettings + } + }); + const runtimeSave = store.saveFolderMonitorRuntimeState({ + ...stale, + paused: true, + pausedAt: 1787712000000 + }); + await Promise.all([settingsSave, runtimeSave]); + + const folderMonitor = store.load().globalSettings.folderMonitor; + assert.equal(folderMonitor.folderPath, 'D:\\new-watch'); + assert.deepEqual(folderMonitor.hosters, ['byse.sx']); + assert.equal(folderMonitor.filterMode, 'exclude'); + assert.equal(folderMonitor.autoStart, false); + assert.equal(folderMonitor.paused, true); + assert.equal(folderMonitor.pausedAt, 1787712000000); + }); + it('drainWrites waits for config and history writes appended while draining', async () => { assert.equal(typeof store.drainWrites, 'function'); await store.save({ globalSettings: { alwaysOnTop: false } }); diff --git a/tests/i18n.test.js b/tests/i18n.test.js index 8f6f96b..7bfae8d 100644 --- a/tests/i18n.test.js +++ b/tests/i18n.test.js @@ -175,6 +175,7 @@ test('runtime queue, account, toast, and shutdown copy translates completely', ( const cases = [ ['Wartet', 'Waiting'], ['Abgebrochen', 'Canceled'], + ['Warteschlange angehalten', 'Queue paused'], ['Fehlgeschlagen: Verbindung verloren', 'Failed: Connection lost'], ['Retry 2/3 · Primär nicht verfügbar', 'Retry 2/3 · Primary unavailable'], ['Link kopiert', 'Link copied'], diff --git a/tests/package-build-files.test.js b/tests/package-build-files.test.js index 4847c32..1fc7343 100644 --- a/tests/package-build-files.test.js +++ b/tests/package-build-files.test.js @@ -48,6 +48,7 @@ function createAutomationLifecycleHarness(mainSource) { const resumeDeferred = createDeferred(); const configuredSettings = []; const startedSettings = []; + let scanError = null; let publishStatus = () => {}; let state = { globalSettings: { @@ -98,13 +99,23 @@ function createAutomationLifecycleHarness(mainSource) { }; folderMonitor.scan = async options => { order.push(`scan:${options.trigger}:${options.emitFiles}`); + if (scanError) throw scanError; return { reachable: true, trigger: options.trigger }; }; const configStore = { load: () => structuredClone(state), - save: config => { + saveFolderMonitorRuntimeState: folderMonitor => { const deferred = createDeferred(); - const snapshot = structuredClone(config); + const snapshot = structuredClone({ + ...state, + globalSettings: { + ...state.globalSettings, + folderMonitor: { + ...state.globalSettings.folderMonitor, + ...folderMonitor + } + } + }); saves.push({ paused: snapshot.globalSettings.folderMonitor.paused, deferred }); order.push(`save:${snapshot.globalSettings.folderMonitor.paused}`); return deferred.promise.then(() => { state = snapshot; }); @@ -112,6 +123,7 @@ function createAutomationLifecycleHarness(mainSource) { }; const uploadManager = { finishAfterActive: () => order.push('finish'), + resumeAfterActive: () => order.push('resume-manager'), startBatch: () => order.push('startBatch') }; const webContents = {}; @@ -155,6 +167,9 @@ function createAutomationLifecycleHarness(mainSource) { resumeDeferred, saves, sent, + setScanError(error) { + scanError = error; + }, setFolderMonitorState(value) { state.globalSettings.folderMonitor = { ...state.globalSettings.folderMonitor, ...value }; }, @@ -165,6 +180,41 @@ function createAutomationLifecycleHarness(mainSource) { }; } +function createAddJobsToBatchHarness(mainSource) { + const blockStart = mainSource.indexOf("ipcMain.handle('add-jobs-to-batch'"); + const blockEnd = mainSource.indexOf("\nipcMain.handle('finish-after-active'", blockStart); + assert.notEqual(blockStart, -1, 'add-jobs-to-batch handler missing'); + assert.notEqual(blockEnd, -1, 'add-jobs-to-batch handler boundary missing'); + const handlers = new Map(); + const effects = []; + const uploadManager = { + running: true, + isStoppingAfterActive: () => true, + sourceFileCleanup: { + registerGroups: async () => { effects.push('sourceCleanup'); return {}; }, + markSkipped: () => effects.push('markSkipped') + }, + addJobs: tasks => { effects.push('addJobs'); return { added: tasks.length, alreadyInBatchJobIds: [] }; } + }; + const config = { globalSettings: { folderMonitor: { paused: false } } }; + vm.runInNewContext(mainSource.slice(blockStart, blockEnd), { + appendUploadPlanAudit: async () => effects.push('audit'), + buildUploadTasksFromJobs: (_config, jobs) => { + effects.push('buildUploadTasks'); + return jobs.map(job => ({ ...job, jobId: job.id })); + }, + closeFlushRequested: false, + configStore: { load: () => config }, + debugLog: () => effects.push('debugLog'), + ipcMain: { handle: (channel, handler) => handlers.set(channel, handler) }, + makeAccountPicker: () => { effects.push('makeAccountPicker'); return {}; }, + persistRotation: () => effects.push('persistRotation'), + summarizeBatchPlan: value => value, + uploadManager + }); + return { effects, handlers }; +} + test('packages every Electron preload referenced by the main process', () => { assert.ok(packageJson.build.files.includes('preload.js')); assert.ok(packageJson.build.files.includes('preload-drop-target.js')); @@ -571,6 +621,47 @@ test('every batch start and extension IPC fails closed before account and cleanu } }); +test('add-jobs-to-batch rejects a stopping manager before account and cleanup side effects', async () => { + const mainSource = fs.readFileSync(path.join(projectRoot, 'main.js'), 'utf8'); + const harness = createAddJobsToBatchHarness(mainSource); + + const result = await harness.handlers.get('add-jobs-to-batch')(null, { + jobs: [{ id: 'job-1', file: 'C:\\watch\\video.mp4', hoster: 'doodstream.com' }], + sourceCleanupGroups: [{ id: 'group-1' }] + }); + + assert.deepEqual({ ...result }, { error: 'Warteschlange angehalten' }); + assert.deepEqual(harness.effects, []); +}); + +test('start-upload waits for a finalizing manager to release before deciding availability', async () => { + const mainSource = fs.readFileSync(path.join(projectRoot, 'main.js'), 'utf8'); + const helperStart = mainSource.indexOf('async function waitForUploadManagerRelease'); + const helperEnd = mainSource.indexOf('\nfunction requestUploadFinalization', helperStart); + const handlerStart = mainSource.indexOf("ipcMain.handle('start-upload'"); + const configLoad = mainSource.indexOf(' const config = configStore.load();', handlerStart); + const waitCall = mainSource.indexOf('await waitForUploadManagerRelease(existingManager)', handlerStart); + assert.notEqual(helperStart, -1); + assert.notEqual(helperEnd, -1); + assert.notEqual(handlerStart, -1); + assert.ok(waitCall > handlerStart && waitCall < configLoad); + + const context = { manager: { running: false }, setTimeout, Date }; + vm.runInNewContext(` +let uploadManager = manager; +${mainSource.slice(helperStart, helperEnd)} +globalThis.waitForRelease = timeout => waitForUploadManagerRelease(manager, timeout); +globalThis.release = () => { uploadManager = null; }; +`, context); + + const released = context.waitForRelease(500); + setTimeout(context.release, 20); + assert.equal(await released, true); + + vm.runInNewContext('uploadManager = manager;', context); + assert.equal(await context.waitForRelease(10), false); +}); + test('automation pause save commits before lifecycle effects and save failure is inert', async () => { const mainSource = fs.readFileSync(path.join(projectRoot, 'main.js'), 'utf8'); const blockStart = Math.max( @@ -607,15 +698,19 @@ test('automation pause save commits before lifecycle effects and save failure is let rejectSave = false; const configStore = { load: () => structuredClone(state), - save: async config => { - const paused = config.globalSettings.folderMonitor.paused; + saveFolderMonitorRuntimeState: async folderMonitor => { + const paused = folderMonitor.paused; order.push(`save:${paused}`); if (rejectSave) throw new Error('save failed'); - state = structuredClone(config); + state.globalSettings.folderMonitor = { + ...state.globalSettings.folderMonitor, + ...structuredClone(folderMonitor) + }; } }; const uploadManager = { finishAfterActive: () => order.push('finish'), + resumeAfterActive: () => order.push('resume-manager'), startBatch: () => order.push('startBatch') }; vm.runInNewContext(mainSource.slice(blockStart, blockEnd), { @@ -658,7 +753,7 @@ test('automation pause save commits before lifecycle effects and save failure is state.globalSettings.folderMonitor.pausedAt = 1; folderMonitor.running = false; await handlers.get('automation:resume')(); - assert.deepEqual(order, ['resume', 'save:false', 'scan:resume:true']); + assert.deepEqual(order, ['resume', 'save:false', 'resume-manager', 'scan:resume:true']); assert.equal(order.includes('startBatch'), false); assert.equal(sent.length, 1); assert.equal(sent[0][1].paused, false); @@ -684,7 +779,7 @@ test('automation lifecycle serializes pause then resume so the newer intent wins harness.saves[1].deferred.resolve(); await Promise.all([pause, resume]); - assert.deepEqual(harness.order, ['save:true', 'pause', 'finish', 'resume', 'save:false', 'scan:resume:true']); + assert.deepEqual(harness.order, ['save:true', 'pause', 'finish', 'resume', 'save:false', 'resume-manager', 'scan:resume:true']); assert.equal(harness.state().globalSettings.folderMonitor.paused, false); assert.equal(harness.sent.length, 1); assert.equal(harness.sent[0][1].paused, false); @@ -711,7 +806,7 @@ test('automation lifecycle serializes resume then pause so the newer intent wins harness.pauseDeferred.resolve(); await Promise.all([resume, pause]); - assert.deepEqual(harness.order, ['resume', 'save:false', 'scan:resume:true', 'save:true', 'pause', 'finish']); + assert.deepEqual(harness.order, ['resume', 'save:false', 'resume-manager', 'scan:resume:true', 'save:true', 'pause', 'finish']); assert.equal(harness.state().globalSettings.folderMonitor.paused, true); assert.equal(harness.sent.length, 1); assert.equal(harness.sent[0][1].paused, true); @@ -808,12 +903,39 @@ test('resume keeps pause authoritative until monitor success and restores the pr assert.equal(result.value.paused, true); assert.equal(result.value.pausedAt, 1); assert.deepEqual(harness.saves.map(save => save.paused), [true]); - assert.deepEqual(harness.order, ['resume', 'stop', 'configure', 'save:true']); + assert.deepEqual(harness.order, ['resume', 'finish', 'stop', 'configure', 'save:true']); assert.equal(harness.state().globalSettings.folderMonitor.paused, true); assert.equal(harness.state().globalSettings.folderMonitor.pausedAt, 1); assert.equal(JSON.stringify(result.value).includes('resume-secret'), false); }); +test('resume reopens the manager before scanning and relatches it when the scan fails', async () => { + const mainSource = fs.readFileSync(path.join(projectRoot, 'main.js'), 'utf8'); + const harness = createAutomationLifecycleHarness(mainSource); + harness.setScanError(new Error('scan failed')); + + const resume = harness.handlers.get('automation:resume')(); + harness.resumeDeferred.resolve(); + await waitForCondition(() => harness.saves.length === 1); + harness.saves[0].deferred.resolve(); + await waitForCondition(() => harness.saves.length === 2); + harness.saves[1].deferred.resolve(); + const result = await resume; + + assert.deepEqual(harness.order, [ + 'resume', + 'save:false', + 'resume-manager', + 'scan:resume:true', + 'finish', + 'stop', + 'configure', + 'save:true' + ]); + assert.equal(result.error, 'Automatik konnte nicht fortgesetzt werden'); + assert.equal(result.paused, true); +}); + test('prepared upload start waits for the final tick and clears recovery when pause wins', async () => { const mainSource = fs.readFileSync(path.join(projectRoot, 'main.js'), 'utf8'); const blockStart = mainSource.indexOf('async function rejectPreparedUploadStart'); diff --git a/tests/startup-renderer.test.js b/tests/startup-renderer.test.js index 18b5b5b..8d06c08 100644 --- a/tests/startup-renderer.test.js +++ b/tests/startup-renderer.test.js @@ -197,6 +197,10 @@ contextBridge.exposeInMainWorld('api', { savedSettings: [] }; }, + setAutomationEvidence(value = {}) { + if (Array.isArray(value.history)) automationProbe.history = value.history; + if (Array.isArray(value.uploadLog)) automationProbe.uploadLog = value.uploadLog; + }, getAutomationProbeState() { return { readCalls: { ...automationProbe.readCalls }, @@ -327,7 +331,7 @@ contextBridge.exposeInMainWorld('api', { return Promise.resolve(automationProbe.addResult || { added: payload?.jobs?.length || 0 }); }, startUpload(payload) { - automationProbe.mutationCalls.push(['start', payload?.jobs?.length || 0]); + automationProbe.mutationCalls.push(['start', payload?.jobs?.length || 0, (payload?.jobs || []).map(job => job.id)]); if (automationProbe.startError) return Promise.reject(new Error(automationProbe.startError)); return Promise.resolve(automationProbe.startResult || { started: true }); }, @@ -824,6 +828,73 @@ contextBridge.exposeInMainWorld('api', { .map(job => normalizeAutomationPath(job.file)))], queuedTelemetry: config.globalSettings.folderMonitor.telemetry.queued }; + configureAtomicState(0); + config.globalSettings.folderMonitor.hosters = ['doodstream.com']; + hosterSettings = {}; + handleBatchDone({ files: [] }); + const evidenceSnapshotFiles = Array.from({ length: 66 }, (_, index) => ({ + path: 'C:\\\\evidence-snapshot\\\\file-' + String(index).padStart(2, '0') + '.mkv', + name: 'file-' + String(index).padStart(2, '0') + '.mkv', + size: 1, + mtimeMs: index + })); + await handleFolderMonitorFiles(evidenceSnapshotFiles); + const evidenceSnapshotProbe = await window.api.getAutomationProbeState(); + const evidenceSnapshotDrain = { + historyCalls: evidenceSnapshotProbe.readCalls.history, + uploadLogCalls: evidenceSnapshotProbe.readCalls.uploadLog, + inspectCalls: evidenceSnapshotProbe.readCalls.inspect, + batchSizes: evidenceSnapshotProbe.logs + .filter(message => message.startsWith('folder-monitor: received ')) + .map(message => Number(message.split(' ')[2] || 0)), + queuedFiles: new Set(queueJobs.filter(job => job.file.startsWith('C:\\\\evidence-snapshot\\\\')).map(job => job.file)).size + }; + configureAtomicState(0); + config.globalSettings.folderMonitor.hosters = ['doodstream.com']; + hosterSettings = {}; + handleBatchDone({ files: [] }); + const separatedEventFiles = Array.from({ length: 66 }, (_, index) => ({ + path: 'C:\\\\separated-events\\\\file-' + String(index).padStart(2, '0') + '.mkv', + name: 'file-' + String(index).padStart(2, '0') + '.mkv', + size: 1, + mtimeMs: index + })); + for (const file of separatedEventFiles) { + await handleFolderMonitorFiles([file]); + await new Promise(resolve => setTimeout(resolve, 0)); + } + const separatedBurstProbe = await window.api.getAutomationProbeState(); + const invalidatedEvidenceFile = { path: 'C:\\\\separated-events\\\\invalidated.mkv', name: 'invalidated.mkv', size: 1, mtimeMs: 100 }; + window.api.setAutomationEvidence({ + history: [{ files: [{ ...invalidatedEvidenceFile, results: [{ hoster: 'doodstream.com', status: 'done' }] }] }] + }); + handleBatchDone({ files: [] }); + await handleFolderMonitorFiles([invalidatedEvidenceFile]); + const invalidatedProbe = await window.api.getAutomationProbeState(); + const expiredEvidenceFile = { path: 'C:\\\\separated-events\\\\expired.mkv', name: 'expired.mkv', size: 1, mtimeMs: 101 }; + window.api.setAutomationEvidence({ + history: [{ files: [{ ...expiredEvidenceFile, results: [{ hoster: 'doodstream.com', status: 'done' }] }] }] + }); + automationEvidenceSnapshotCache.expiresAt = 0; + await handleFolderMonitorFiles([expiredEvidenceFile]); + const expiredProbe = await window.api.getAutomationProbeState(); + const separatedEventEvidence = { + afterBurst: { + historyCalls: separatedBurstProbe.readCalls.history, + uploadLogCalls: separatedBurstProbe.readCalls.uploadLog, + queuedFiles: new Set(queueJobs.filter(job => job.file.startsWith('C:\\\\separated-events\\\\file-')).map(job => job.file)).size + }, + afterInvalidation: { + historyCalls: invalidatedProbe.readCalls.history, + uploadLogCalls: invalidatedProbe.readCalls.uploadLog, + queued: queueJobs.some(job => normalizeAutomationPath(job.file) === normalizeAutomationPath(invalidatedEvidenceFile.path)) + }, + afterExpiry: { + historyCalls: expiredProbe.readCalls.history, + uploadLogCalls: expiredProbe.readCalls.uploadLog, + queued: queueJobs.some(job => normalizeAutomationPath(job.file) === normalizeAutomationPath(expiredEvidenceFile.path)) + } + }; configureAtomicState(18); config.globalSettings.folderMonitor.queueLimitJobs = 20; config.globalSettings.folderMonitor.hosters = ['doodstream.com']; @@ -1838,6 +1909,76 @@ contextBridge.exposeInMainWorld('api', { secretExposed: JSON.stringify({ pausedInjection, unconfirmedInjection, exceptionInjection, telemetryFailure }).includes('secret') }; + configureAtomicState(0); + const pauseMarkerJob = makePauseRaceJob('pause-marker.mkv'); + pauseMarkerJob.status = 'queued'; + queueJobs = [pauseMarkerJob]; + uploading = true; + rebuildJobIndex(); + window.api.configureAutomationProbe({ paused: false }); + applyAutomationRuntimeStatus({ paused: false }); + await toggleAutomationPauseResume(); + const pauseMarkerPersisted = buildPersistedQueueState()?.queueJobs.find(entry => entry.id === pauseMarkerJob.id); + const pauseMarker = { + marked: pauseMarkerJob.automationPaused === true, + persisted: pauseMarkerPersisted?.automationPaused === true, + persistedStatus: pauseMarkerPersisted?.status || null + }; + + const runResumeQueueCase = async ({ active, resumeError = '' }) => { + configureAtomicState(0); + const job = makePauseRaceJob(active ? 'resume-active.mkv' : 'resume-idle.mkv'); + job.id = active ? 'resume-active' : 'resume-idle'; + job.status = 'aborted'; + job.error = 'Warteschlange angehalten'; + job.automationPaused = true; + queueJobs = [ + job, + { ...makePauseRaceJob('manual-preview.mkv'), id: 'manual-preview', status: 'preview' }, + { ...makePauseRaceJob('manual-queued.mkv'), id: 'manual-queued', status: 'queued' }, + { ...makePauseRaceJob('manual-error.mkv'), id: 'manual-error', status: 'error' }, + { ...makePauseRaceJob('manual-skipped.mkv'), id: 'manual-skipped', status: 'skipped' } + ]; + selectedFiles = []; + selectedUploadHosters = ['doodstream.com']; + config.globalSettings.folderMonitor.paused = true; + uploading = active; + rebuildJobIndex(); + window.api.configureAutomationProbe({ + paused: true, + runtimeStatus: resumeError ? { error: resumeError } : {}, + addResult: { added: 1 }, + startResult: { started: true } + }); + applyAutomationRuntimeStatus({ paused: true }); + await toggleAutomationPauseResume(); + const probe = await window.api.getAutomationProbeState(); + const acceptedStatus = job.status; + const persistedJob = buildPersistedQueueState()?.queueJobs.find(entry => entry.id === job.id); + if (!resumeError) { + handleProgress({ + jobId: job.id, + fileName: job.fileName, + hoster: job.hoster, + status: 'getting-server', + bytesUploaded: 0, + bytesTotal: job.bytesTotal + }); + } + return { + status: acceptedStatus, + uploading, + markerPersisted: persistedJob?.automationPaused === true, + markerAfterProgress: job.automationPaused === true, + mutations: probe.mutationCalls.map(call => ({ kind: call[0], count: call[1] || 0, ids: call[2] || [] })) + }; + }; + const resumeQueue = { + active: await runResumeQueueCase({ active: true }), + idle: await runResumeQueueCase({ active: false }), + rollback: await runResumeQueueCase({ active: false, resumeError: 'Automatik konnte nicht fortgesetzt werden.' }) + }; + configureAtomicState(0); config.globalSettings.folderMonitor.paused = true; window.api.configureAutomationProbe({ paused: true }); @@ -1879,7 +2020,7 @@ contextBridge.exposeInMainWorld('api', { startCalls: pausedProbe.mutationCalls.filter(call => call[0] === 'start').length, injectCalls: pausedProbe.mutationCalls.filter(call => call[0] === 'inject').length }; - return { dry, manualTest, historyEvidence, pendingDedup, parallelAdmission, distinctParallel, disjointClassification, manualHostTransactional, atomic, status, zeroAdmission, stress, persistedQueueExactness, stale, replannedEligibility, mainPauseResponses, cleanupRollback, crossPathCleanupRollback, partialAddOutcomes, collisionResolver, collisionAdmission, pauseBetweenApplyAndStart, startAcceptance, fulfilledFeedback, injectionOutcomes, paused }; + return { dry, manualTest, historyEvidence, pendingDedup, parallelAdmission, evidenceSnapshotDrain, separatedEventEvidence, distinctParallel, disjointClassification, manualHostTransactional, atomic, status, zeroAdmission, stress, persistedQueueExactness, stale, replannedEligibility, mainPauseResponses, cleanupRollback, crossPathCleanupRollback, partialAddOutcomes, collisionResolver, collisionAdmission, pauseBetweenApplyAndStart, startAcceptance, fulfilledFeedback, injectionOutcomes, pauseMarker, resumeQueue, paused }; })()`; const automationControlCenterScript = `(async () => { const waitFor = async predicate => { @@ -2670,6 +2811,52 @@ app.whenReady().then(async () => { matchingPaths: ['c:/watch/parallel.mkv'], queuedTelemetry: 1 }); + assert.deepEqual(result.automationPipeline.evidenceSnapshotDrain, { + historyCalls: 1, + uploadLogCalls: 1, + inspectCalls: 9, + batchSizes: [8, 8, 8, 8, 8, 8, 8, 8, 2], + queuedFiles: 66 + }); + assert.deepEqual(result.automationPipeline.separatedEventEvidence, { + afterBurst: { historyCalls: 1, uploadLogCalls: 1, queuedFiles: 66 }, + afterInvalidation: { historyCalls: 2, uploadLogCalls: 2, queued: false }, + afterExpiry: { historyCalls: 3, uploadLogCalls: 3, queued: false } + }); + assert.deepEqual(result.automationPipeline.resumeQueue, { + active: { + status: 'queued', + uploading: true, + markerPersisted: true, + markerAfterProgress: false, + mutations: [ + { kind: 'resume', count: 0, ids: [] }, + { kind: 'inject', count: 1, ids: ['resume-active'] } + ] + }, + idle: { + status: 'queued', + uploading: true, + markerPersisted: true, + markerAfterProgress: false, + mutations: [ + { kind: 'resume', count: 0, ids: [] }, + { kind: 'start', count: 1, ids: ['resume-idle'] } + ] + }, + rollback: { + status: 'aborted', + uploading: false, + markerPersisted: true, + markerAfterProgress: true, + mutations: [{ kind: 'resume', count: 0, ids: [] }] + } + }); + assert.deepEqual(result.automationPipeline.pauseMarker, { + marked: true, + persisted: true, + persistedStatus: 'queued' + }); assert.deepEqual(result.automationPipeline.distinctParallel, { inspectCalls: 3, maxConcurrentInspections: 1, @@ -3670,6 +3857,7 @@ let walkCalls = 0; let addJobsCalls = 0; let startBatchCalls = 0; let finishCalls = 0; +let stoppingAfterActive = true; let cleanupRelease; let cleanupStartedResolve; const cleanupStarted = new Promise(resolve => { cleanupStartedResolve = resolve; }); @@ -3727,11 +3915,16 @@ let uploadManager = { addJobsCalls++; return { added: tasks.length, alreadyInBatchJobIds: [] }; }, + isStoppingAfterActive: () => stoppingAfterActive, + resumeAfterActive: () => { stoppingAfterActive = false; }, startBatch: () => { startBatchCalls++; return Promise.resolve(); }, - finishAfterActive: () => { finishCalls++; } + finishAfterActive: () => { + finishCalls++; + stoppingAfterActive = true; + } }; ${productionHandlers} ${automationHandlers} diff --git a/tests/upload-manager.test.js b/tests/upload-manager.test.js index c1c8a00..58edfbe 100644 --- a/tests/upload-manager.test.js +++ b/tests/upload-manager.test.js @@ -648,6 +648,231 @@ describe('UploadManager', () => { assert.equal(batchDoneEvents.length, 1); }); + it('addJobs rejects new work while stopping and accepts it after resume', async () => { + let releaseActive; + const started = []; + mockUploadFile.mock.mockImplementation(async (hoster, filePath, apiKey, onProgress) => { + started.push(filePath); + if (filePath.endsWith('/active.mp4')) { + await new Promise(resolve => { releaseActive = resolve; }); + } + if (onProgress) onProgress(fakeFileSize, fakeFileSize); + return { download_url: `https://${hoster}/d/ok123`, embed_url: null, file_code: 'ok123' }; + }); + const mgr = new UploadManager({ + 'doodstream.com': { retries: 0, parallelCount: 1, maxSpeedKbs: 0, restartBelowKbs: 0, timeIntervalSec: 0, maxSizeMb: 0 } + }); + const batch = mgr.startBatch([ + { jobId: 'active', file: '/test/active.mp4', hoster: 'doodstream.com', apiKey: 'key1' } + ]); + + for (let attempt = 0; attempt < 50 && !releaseActive; attempt++) { + await new Promise(resolve => setTimeout(resolve, 5)); + } + assert.equal(typeof releaseActive, 'function'); + + mgr.finishAfterActive(); + const stopping = typeof mgr.isStoppingAfterActive === 'function' + ? mgr.isStoppingAfterActive() + : undefined; + const result = mgr.addJobs([ + { jobId: 'rejected', file: '/test/rejected.mp4', hoster: 'doodstream.com', apiKey: 'key1' } + ]); + await mgr.resumeAfterActive(); + const resumedResult = mgr.addJobs([ + { jobId: 'resumed', file: '/test/resumed.mp4', hoster: 'doodstream.com', apiKey: 'key1' } + ]); + + releaseActive(); + await batch; + + assert.equal(stopping, true); + assert.deepEqual(result, { added: 0, alreadyInBatchJobIds: [] }); + assert.deepEqual(resumedResult, { added: 1, alreadyInBatchJobIds: [] }); + assert.equal(mgr.isStoppingAfterActive(), false); + assert.deepEqual(started, ['/test/active.mp4', '/test/resumed.mp4']); + }); + + it('finishAfterActive bypasses queued interval waits', async () => { + let releaseActive; + const started = []; + mockUploadFile.mock.mockImplementation(async (hoster, filePath, apiKey, onProgress) => { + started.push(filePath); + if (filePath.endsWith('/active.mp4')) { + await new Promise(resolve => { releaseActive = resolve; }); + } + if (onProgress) onProgress(fakeFileSize, fakeFileSize); + return { download_url: `https://${hoster}/d/ok123`, embed_url: null, file_code: 'ok123' }; + }); + const mgr = new UploadManager({ + 'doodstream.com': { retries: 0, parallelCount: 1, maxSpeedKbs: 0, restartBelowKbs: 0, timeIntervalSec: 1, maxSizeMb: 0 } + }); + const batch = mgr.startBatch([ + { jobId: 'active', file: '/test/active.mp4', hoster: 'doodstream.com', apiKey: 'key1' }, + { jobId: 'queued-1', file: '/test/queued-1.mp4', hoster: 'doodstream.com', apiKey: 'key1' }, + { jobId: 'queued-2', file: '/test/queued-2.mp4', hoster: 'doodstream.com', apiKey: 'key1' }, + { jobId: 'queued-3', file: '/test/queued-3.mp4', hoster: 'doodstream.com', apiKey: 'key1' } + ]); + + for (let attempt = 0; attempt < 50 && !releaseActive; attempt++) { + await new Promise(resolve => setTimeout(resolve, 5)); + } + assert.equal(typeof releaseActive, 'function'); + + mgr.finishAfterActive(); + const stoppedAt = Date.now(); + releaseActive(); + await batch; + + assert.ok(Date.now() - stoppedAt < 500, `queued jobs took ${Date.now() - stoppedAt} ms to stop`); + assert.deepEqual(started, ['/test/active.mp4']); + }); + + it('finishAfterActive interrupts a job already waiting inside the upload interval', async () => { + let releaseActive; + let intervalEnteredResolve; + const intervalEntered = new Promise(resolve => { intervalEnteredResolve = resolve; }); + let intervalCalls = 0; + const started = []; + mockUploadFile.mock.mockImplementation(async (hoster, filePath, apiKey, onProgress) => { + started.push(filePath); + if (filePath.endsWith('/active.mp4')) { + await new Promise(resolve => { releaseActive = resolve; }); + } + if (onProgress) onProgress(fakeFileSize, fakeFileSize); + return { download_url: `https://${hoster}/d/ok123`, embed_url: null, file_code: 'ok123' }; + }); + const mgr = new UploadManager({ + 'doodstream.com': { retries: 0, parallelCount: 2, maxSpeedKbs: 0, restartBelowKbs: 0, timeIntervalSec: 2, maxSizeMb: 0 } + }); + mgr._waitForInterval = (hoster, intervalMs, signal) => new Promise((resolve, reject) => { + intervalCalls++; + if (intervalCalls === 1) { + resolve(); + return; + } + intervalEnteredResolve(); + if (signal.aborted) reject(new Error('Aborted')); + else signal.addEventListener('abort', () => reject(new Error('Aborted')), { once: true }); + }); + const batch = mgr.startBatch([ + { jobId: 'active', file: '/test/active.mp4', hoster: 'doodstream.com', apiKey: 'key1' } + ]); + + for (let attempt = 0; attempt < 50 && !releaseActive; attempt++) { + await new Promise(resolve => setTimeout(resolve, 5)); + } + assert.equal(typeof releaseActive, 'function'); + assert.deepEqual(mgr.addJobs([ + { jobId: 'interval-waiter', file: '/test/interval-waiter.mp4', hoster: 'doodstream.com', apiKey: 'key1' } + ]), { added: 1, alreadyInBatchJobIds: [] }); + await intervalEntered; + + const stoppedAt = Date.now(); + mgr.finishAfterActive(); + releaseActive(); + await batch; + + assert.ok(Date.now() - stoppedAt < 500, `interval waiter took ${Date.now() - stoppedAt} ms to stop`); + assert.deepEqual(started, ['/test/active.mp4']); + }); + + it('finishAfterActive interrupts a job already waiting for the global upload slot', async () => { + let releaseActive; + const started = []; + mockUploadFile.mock.mockImplementation(async (hoster, filePath, apiKey, onProgress) => { + started.push(filePath); + if (filePath.endsWith('/active.mp4')) { + await new Promise(resolve => { releaseActive = resolve; }); + } + if (onProgress) onProgress(fakeFileSize, fakeFileSize); + return { download_url: `https://${hoster}/d/ok123`, embed_url: null, file_code: 'ok123' }; + }); + const mgr = new UploadManager({ + 'doodstream.com': { retries: 0, parallelCount: 1, maxSpeedKbs: 0, restartBelowKbs: 0, timeIntervalSec: 0, maxSizeMb: 0 }, + 'byse.sx': { retries: 0, parallelCount: 1, maxSpeedKbs: 0, restartBelowKbs: 0, timeIntervalSec: 0, maxSizeMb: 0 } + }, { parallelUploadCount: 1 }); + const batch = mgr.startBatch([ + { jobId: 'active', file: '/test/active.mp4', hoster: 'doodstream.com', apiKey: 'key1' } + ]); + + for (let attempt = 0; attempt < 50 && !releaseActive; attempt++) { + await new Promise(resolve => setTimeout(resolve, 5)); + } + assert.equal(typeof releaseActive, 'function'); + assert.deepEqual(mgr.addJobs([ + { jobId: 'global-waiter', file: '/test/global-waiter.mp4', hoster: 'byse.sx', apiKey: 'key2' } + ]), { added: 1, alreadyInBatchJobIds: [] }); + for (let attempt = 0; attempt < 50 && mgr.globalSemaphore.pending === 0; attempt++) { + await new Promise(resolve => setTimeout(resolve, 5)); + } + assert.equal(mgr.globalSemaphore.pending, 1); + + const stoppedAt = Date.now(); + mgr.finishAfterActive(); + releaseActive(); + await batch; + + assert.ok(Date.now() - stoppedAt < 500, `global waiter took ${Date.now() - stoppedAt} ms to stop`); + assert.deepEqual(started, ['/test/active.mp4']); + }); + + it('resumeAfterActive waits for stopped admission jobs before reopening the queue', async () => { + let releaseActive; + let intervalEnteredResolve; + const intervalEntered = new Promise(resolve => { intervalEnteredResolve = resolve; }); + let intervalCalls = 0; + const started = []; + const terminal = []; + mockUploadFile.mock.mockImplementation(async (hoster, filePath, apiKey, onProgress) => { + started.push(filePath); + if (filePath.endsWith('/active.mp4')) { + await new Promise(resolve => { releaseActive = resolve; }); + } + if (onProgress) onProgress(fakeFileSize, fakeFileSize); + return { download_url: `https://${hoster}/d/ok123`, embed_url: null, file_code: 'ok123' }; + }); + const mgr = new UploadManager({ + 'doodstream.com': { retries: 0, parallelCount: 2, maxSpeedKbs: 0, restartBelowKbs: 0, timeIntervalSec: 2, maxSizeMb: 0 } + }); + mgr._waitForInterval = (hoster, intervalMs, signal) => new Promise((resolve, reject) => { + intervalCalls++; + if (intervalCalls === 1) { + resolve(); + return; + } + intervalEnteredResolve(); + if (signal.aborted) reject(new Error('Aborted')); + else signal.addEventListener('abort', () => reject(new Error('Aborted')), { once: true }); + }); + mgr.on('progress', value => { + if (value.jobId === 'interval-waiter' && ['aborted', 'error'].includes(value.status)) terminal.push(value.status); + }); + const batch = mgr.startBatch([ + { jobId: 'active', file: '/test/active.mp4', hoster: 'doodstream.com', apiKey: 'key1' } + ]); + + for (let attempt = 0; attempt < 50 && !releaseActive; attempt++) { + await new Promise(resolve => setTimeout(resolve, 5)); + } + assert.equal(typeof releaseActive, 'function'); + assert.deepEqual(mgr.addJobs([ + { jobId: 'interval-waiter', file: '/test/interval-waiter.mp4', hoster: 'doodstream.com', apiKey: 'key1' } + ]), { added: 1, alreadyInBatchJobIds: [] }); + await intervalEntered; + + const resumedAt = Date.now(); + mgr.finishAfterActive(); + await mgr.resumeAfterActive(); + + assert.ok(Date.now() - resumedAt < 500, `resume waited ${Date.now() - resumedAt} ms for admission shutdown`); + assert.equal(mgr.isStoppingAfterActive(), false); + assert.deepEqual(terminal, ['aborted']); + releaseActive(); + await batch; + assert.deepEqual(started, ['/test/active.mp4']); + }); + it('_combineSignals propagates abort from either source', () => { const mgr = new UploadManager({}); const ac1 = new AbortController();