Fix concurrent upload batch accounting
Reserve accepted job IDs before asynchronous upload preparation so parallel additions cannot start duplicate jobs or replace their abort controllers. Track late-added jobs in the batch total and keep failure counts non-negative, with focused race and summary regressions.
This commit is contained in:
@@ -53,6 +53,8 @@ class UploadManager extends EventEmitter {
|
|||||||
this._doodApiKeyCache = new Map(); // accountId/username -> derived doodstream API key ('' = tried, none)
|
this._doodApiKeyCache = new Map(); // accountId/username -> derived doodstream API key ('' = tried, none)
|
||||||
this._baselineCache = new Map(); // hoster:apiKey -> Promise<Set<file_code>> (one fetch shared across all jobs in batch)
|
this._baselineCache = new Map(); // hoster:apiKey -> Promise<Set<file_code>> (one fetch shared across all jobs in batch)
|
||||||
this._recoveryClaims = createRecoveryClaimRegistry();
|
this._recoveryClaims = createRecoveryClaimRegistry();
|
||||||
|
this._batchJobIds = new Set();
|
||||||
|
this._batchTotal = 0;
|
||||||
}
|
}
|
||||||
|
|
||||||
updateAccountPools(accountPools) {
|
updateAccountPools(accountPools) {
|
||||||
@@ -336,6 +338,8 @@ class UploadManager extends EventEmitter {
|
|||||||
async startBatch(tasks, opts = {}) {
|
async startBatch(tasks, opts = {}) {
|
||||||
const pendingCancelledJobIds = new Set(this.pendingCancelledJobIds);
|
const pendingCancelledJobIds = new Set(this.pendingCancelledJobIds);
|
||||||
const pendingCancelAll = this.pendingCancelAll;
|
const pendingCancelAll = this.pendingCancelAll;
|
||||||
|
this._batchJobIds = new Set(tasks.map((task) => task.jobId).filter(Boolean));
|
||||||
|
this._batchTotal = tasks.length;
|
||||||
this.pendingCancelledJobIds.clear();
|
this.pendingCancelledJobIds.clear();
|
||||||
this.pendingCancelAll = false;
|
this.pendingCancelAll = false;
|
||||||
this.running = true;
|
this.running = true;
|
||||||
@@ -422,7 +426,7 @@ class UploadManager extends EventEmitter {
|
|||||||
this._emitStats();
|
this._emitStats();
|
||||||
|
|
||||||
const files = Array.from(results.values());
|
const files = Array.from(results.values());
|
||||||
const total = tasks.length;
|
const total = this._batchTotal;
|
||||||
const succeeded = files.reduce((count, file) => count + file.results.filter((result) => result.status === 'done').length, 0);
|
const succeeded = files.reduce((count, file) => count + file.results.filter((result) => result.status === 'done').length, 0);
|
||||||
const skipped = files.reduce((count, file) => count + file.results.filter((result) => result.status === 'skipped').length, 0);
|
const skipped = files.reduce((count, file) => count + file.results.filter((result) => result.status === 'skipped').length, 0);
|
||||||
|
|
||||||
@@ -431,7 +435,7 @@ class UploadManager extends EventEmitter {
|
|||||||
timestamp: new Date().toISOString(),
|
timestamp: new Date().toISOString(),
|
||||||
total,
|
total,
|
||||||
succeeded,
|
succeeded,
|
||||||
failed: total - succeeded - skipped,
|
failed: Math.max(0, total - succeeded - skipped),
|
||||||
skipped,
|
skipped,
|
||||||
files
|
files
|
||||||
};
|
};
|
||||||
@@ -1461,16 +1465,18 @@ class UploadManager extends EventEmitter {
|
|||||||
const addResult = { added: 0, alreadyInBatchJobIds: [] };
|
const addResult = { added: 0, alreadyInBatchJobIds: [] };
|
||||||
for (const task of tasks) {
|
for (const task of tasks) {
|
||||||
// Skip if this job is already being processed (prevent duplicates)
|
// Skip if this job is already being processed (prevent duplicates)
|
||||||
if (task.jobId && this.jobAbortControllers.has(task.jobId)) {
|
if (task.jobId && this._batchJobIds.has(task.jobId)) {
|
||||||
addResult.alreadyInBatchJobIds.push(task.jobId);
|
addResult.alreadyInBatchJobIds.push(task.jobId);
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
if (task.jobId) this._batchJobIds.add(task.jobId);
|
||||||
const fileName = path.basename(task.file);
|
const fileName = path.basename(task.file);
|
||||||
if (!results.has(task.file)) {
|
if (!results.has(task.file)) {
|
||||||
let size = 0;
|
let size = 0;
|
||||||
try { size = fs.statSync(task.file).size; } catch {}
|
try { size = fs.statSync(task.file).size; } catch {}
|
||||||
results.set(task.file, { name: fileName, size, results: [] });
|
results.set(task.file, { name: fileName, size, results: [] });
|
||||||
}
|
}
|
||||||
|
this._batchTotal++;
|
||||||
this._additionalPromises.push(this._runJob(task, results, signal));
|
this._additionalPromises.push(this._runJob(task, results, signal));
|
||||||
addResult.added++;
|
addResult.added++;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -648,6 +648,70 @@ describe('UploadManager', () => {
|
|||||||
assert.ok(statuses.some((entry) => entry.jobId === 'job-third' && entry.status === 'done'));
|
assert.ok(statuses.some((entry) => entry.jobId === 'job-third' && entry.status === 'done'));
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it('accepts each job ID only once across concurrent addJobs calls', async () => {
|
||||||
|
const anchorPath = '/race/anchor.mp4';
|
||||||
|
const racePath = '/race/concurrent.mp4';
|
||||||
|
const originalStatSync = fs.statSync;
|
||||||
|
const originalStat = fs.promises.stat;
|
||||||
|
const pendingStats = [];
|
||||||
|
let concurrentStarts = 0;
|
||||||
|
let batchPromise;
|
||||||
|
|
||||||
|
fs.statSync = function(p) {
|
||||||
|
if (p === anchorPath || p === racePath) return { size: 0 };
|
||||||
|
return originalStatSync.call(this, p);
|
||||||
|
};
|
||||||
|
fs.promises.stat = function(p) {
|
||||||
|
if (p === anchorPath || p === racePath) {
|
||||||
|
return new Promise((resolve) => pendingStats.push({ path: p, resolve }));
|
||||||
|
}
|
||||||
|
return originalStat.call(this, p);
|
||||||
|
};
|
||||||
|
|
||||||
|
try {
|
||||||
|
mockUploadFile.mock.mockImplementation(async (hoster, filePath, apiKey, onProgress) => {
|
||||||
|
if (filePath === racePath) concurrentStarts++;
|
||||||
|
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: 3, maxSpeedKbs: 0, restartBelowKbs: 0, timeIntervalSec: 0, maxSizeMb: 0 }
|
||||||
|
});
|
||||||
|
const controllerRegistrations = [];
|
||||||
|
const registerController = mgr.jobAbortControllers.set;
|
||||||
|
mgr.jobAbortControllers.set = function(jobId, controller) {
|
||||||
|
if (jobId === 'job-concurrent') controllerRegistrations.push(controller);
|
||||||
|
return registerController.call(this, jobId, controller);
|
||||||
|
};
|
||||||
|
|
||||||
|
batchPromise = mgr.startBatch([
|
||||||
|
{ jobId: 'job-anchor', file: anchorPath, hoster: 'doodstream.com', apiKey: 'k' }
|
||||||
|
]);
|
||||||
|
assert.equal(mgr.running, true);
|
||||||
|
|
||||||
|
const task = { jobId: 'job-concurrent', file: racePath, hoster: 'doodstream.com', apiKey: 'k' };
|
||||||
|
const addResults = await Promise.all([
|
||||||
|
Promise.resolve().then(() => mgr.addJobs([task])),
|
||||||
|
Promise.resolve().then(() => mgr.addJobs([{ ...task }]))
|
||||||
|
]);
|
||||||
|
const concurrentStats = pendingStats.filter((entry) => entry.path === racePath);
|
||||||
|
for (const entry of pendingStats) entry.resolve({ size: fakeFileSize });
|
||||||
|
await batchPromise;
|
||||||
|
|
||||||
|
assert.equal(addResults.reduce((total, result) => total + result.added, 0), 1);
|
||||||
|
assert.deepEqual(addResults.flatMap((result) => result.alreadyInBatchJobIds), ['job-concurrent']);
|
||||||
|
assert.equal(concurrentStats.length, 1);
|
||||||
|
assert.equal(controllerRegistrations.length, 1);
|
||||||
|
assert.equal(concurrentStarts, 1);
|
||||||
|
} finally {
|
||||||
|
for (const entry of pendingStats) entry.resolve({ size: fakeFileSize });
|
||||||
|
if (batchPromise) await batchPromise.catch(() => {});
|
||||||
|
fs.statSync = originalStatSync;
|
||||||
|
fs.promises.stat = originalStat;
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
it('_combineSignals propagates abort from either source', () => {
|
it('_combineSignals propagates abort from either source', () => {
|
||||||
const mgr = new UploadManager({});
|
const mgr = new UploadManager({});
|
||||||
const ac1 = new AbortController();
|
const ac1 = new AbortController();
|
||||||
@@ -771,7 +835,7 @@ describe('UploadManager', () => {
|
|||||||
assert.ok(maxConcurrent <= 2, `scaleParallelUploads should cap at 2, was ${maxConcurrent}`);
|
assert.ok(maxConcurrent <= 2, `scaleParallelUploads should cap at 2, was ${maxConcurrent}`);
|
||||||
});
|
});
|
||||||
|
|
||||||
it('addJobs injects new tasks into running batch', async () => {
|
it('addJobs includes newly injected tasks in the batch summary', async () => {
|
||||||
let started = 0;
|
let started = 0;
|
||||||
mockUploadFile.mock.mockImplementation(async (hoster, filePath, apiKey, onProgress) => {
|
mockUploadFile.mock.mockImplementation(async (hoster, filePath, apiKey, onProgress) => {
|
||||||
started++;
|
started++;
|
||||||
@@ -802,6 +866,9 @@ describe('UploadManager', () => {
|
|||||||
await batchPromise;
|
await batchPromise;
|
||||||
assert.ok(summary);
|
assert.ok(summary);
|
||||||
assert.equal(started, 4, 'all 4 jobs should have run');
|
assert.equal(started, 4, 'all 4 jobs should have run');
|
||||||
|
assert.equal(summary.total, 4);
|
||||||
|
assert.equal(summary.succeeded, 4);
|
||||||
|
assert.equal(summary.failed, 0);
|
||||||
});
|
});
|
||||||
|
|
||||||
it('addJobs rejects duplicates already in running batch', async () => {
|
it('addJobs rejects duplicates already in running batch', async () => {
|
||||||
|
|||||||
Reference in New Issue
Block a user