Close post-upload and admission race gaps
Treat fallback result fetch and body failures as uncertain Doodstream commits, and coordinate interval timing with slot acquisition so upload starts remain spaced without occupying another host's global slot during the wait.
This commit is contained in:
@@ -640,15 +640,26 @@ class DoodstreamUploader {
|
|||||||
if (formAction) {
|
if (formAction) {
|
||||||
_debugLog(`Fallback: following form action ${safeEndpoint(formAction[1]) || 'unknown endpoint'}`);
|
_debugLog(`Fallback: following form action ${safeEndpoint(formAction[1]) || 'unknown endpoint'}`);
|
||||||
const formData = new URLSearchParams(hiddenFields);
|
const formData = new URLSearchParams(hiddenFields);
|
||||||
const followRes = await this._fetch(formAction[1], {
|
let followText;
|
||||||
method: 'POST',
|
try {
|
||||||
headers: {
|
const followRes = await this._fetch(formAction[1], {
|
||||||
'Content-Type': 'application/x-www-form-urlencoded',
|
method: 'POST',
|
||||||
'Referer': BASE_URL + '/'
|
headers: {
|
||||||
},
|
'Content-Type': 'application/x-www-form-urlencoded',
|
||||||
body: formData.toString()
|
'Referer': BASE_URL + '/'
|
||||||
});
|
},
|
||||||
const followText = await followRes.text();
|
body: formData.toString()
|
||||||
|
});
|
||||||
|
followText = await followRes.text();
|
||||||
|
} catch {
|
||||||
|
throw createTransportError('Doodstream Upload: Redirect-Antwort konnte nicht gelesen werden', {
|
||||||
|
phase: 'upload-result-submit',
|
||||||
|
endpoint: formAction[1],
|
||||||
|
retryable: true,
|
||||||
|
transientNetwork: true,
|
||||||
|
remoteCommitUncertain: true
|
||||||
|
});
|
||||||
|
}
|
||||||
_debugLog(`Fallback response: ${summarizeResponse(followText, '')}`);
|
_debugLog(`Fallback response: ${summarizeResponse(followText, '')}`);
|
||||||
|
|
||||||
const fallbackCode = this._findFilecodeInHtml(followText);
|
const fallbackCode = this._findFilecodeInHtml(followText);
|
||||||
|
|||||||
@@ -1264,20 +1264,24 @@ class UploadManager extends EventEmitter {
|
|||||||
const globalSemaphore = this._getGlobalSemaphore();
|
const globalSemaphore = this._getGlobalSemaphore();
|
||||||
let hosterSlotAcquired = false;
|
let hosterSlotAcquired = false;
|
||||||
let globalSlotAcquired = false;
|
let globalSlotAcquired = false;
|
||||||
try {
|
const acquireSlots = async () => {
|
||||||
await hosterSemaphore.acquire(signal);
|
await hosterSemaphore.acquire(signal);
|
||||||
hosterSlotAcquired = true;
|
hosterSlotAcquired = true;
|
||||||
if (globalSemaphore) {
|
if (globalSemaphore) {
|
||||||
await globalSemaphore.acquire(signal);
|
await globalSemaphore.acquire(signal);
|
||||||
globalSlotAcquired = true;
|
globalSlotAcquired = true;
|
||||||
}
|
}
|
||||||
|
};
|
||||||
|
try {
|
||||||
if (this._swapFailedAccount(task, jobId, path.basename(task.file))) {
|
if (this._swapFailedAccount(task, jobId, path.basename(task.file))) {
|
||||||
retryAdmission = true;
|
retryAdmission = true;
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
const settings = this._getSettings(task.hoster);
|
const settings = this._getSettings(task.hoster);
|
||||||
if (settings.timeIntervalSec > 0) {
|
if (settings.timeIntervalSec > 0) {
|
||||||
await this._waitForInterval(task.hoster, settings.timeIntervalSec * 1000, signal);
|
await this._waitForInterval(task.hoster, settings.timeIntervalSec * 1000, signal, acquireSlots);
|
||||||
|
} else {
|
||||||
|
await acquireSlots();
|
||||||
}
|
}
|
||||||
try {
|
try {
|
||||||
return await this._executeUpload(task, progressCb, signal, throttle, fileProbe, context);
|
return await this._executeUpload(task, progressCb, signal, throttle, fileProbe, context);
|
||||||
@@ -1618,7 +1622,7 @@ class UploadManager extends EventEmitter {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
_waitForInterval(hoster, intervalMs, signal) {
|
_waitForInterval(hoster, intervalMs, signal, acquireSlots) {
|
||||||
// Serialize interval waits per hoster so concurrent jobs queue up properly
|
// Serialize interval waits per hoster so concurrent jobs queue up properly
|
||||||
const prev = this.intervalLocks[hoster] || Promise.resolve();
|
const prev = this.intervalLocks[hoster] || Promise.resolve();
|
||||||
const next = prev.then(async () => {
|
const next = prev.then(async () => {
|
||||||
@@ -1628,6 +1632,7 @@ class UploadManager extends EventEmitter {
|
|||||||
if (elapsed < intervalMs) {
|
if (elapsed < intervalMs) {
|
||||||
await this._sleep(intervalMs - elapsed, signal);
|
await this._sleep(intervalMs - elapsed, signal);
|
||||||
}
|
}
|
||||||
|
await acquireSlots();
|
||||||
this.lastStartTime[hoster] = Date.now();
|
this.lastStartTime[hoster] = Date.now();
|
||||||
});
|
});
|
||||||
this.intervalLocks[hoster] = next.catch(() => {});
|
this.intervalLocks[hoster] = next.catch(() => {});
|
||||||
|
|||||||
@@ -329,3 +329,35 @@ test('redirect fetch failure after upload is marked as an uncertain remote commi
|
|||||||
fs.rmSync(root, { recursive: true, force: true });
|
fs.rmSync(root, { recursive: true, force: true });
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test('fallback form fetch failure after upload is marked as an uncertain remote commit', async () => {
|
||||||
|
const up = new DoodstreamUploader();
|
||||||
|
up._fetch = async () => {
|
||||||
|
throw new Error('fallback request failed');
|
||||||
|
};
|
||||||
|
await assert.rejects(
|
||||||
|
() => up._parseUploadResponse('<form action="https://doodstream.com/result"></form>'),
|
||||||
|
(err) => {
|
||||||
|
assert.equal(err.remoteCommitUncertain, true);
|
||||||
|
assert.equal(err.diagnostic.phase, 'upload-result-submit');
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
test('fallback form body failure after upload is marked as an uncertain remote commit', async () => {
|
||||||
|
const up = new DoodstreamUploader();
|
||||||
|
up._fetch = async () => ({
|
||||||
|
text: async () => {
|
||||||
|
throw new Error('fallback body failed');
|
||||||
|
}
|
||||||
|
});
|
||||||
|
await assert.rejects(
|
||||||
|
() => up._parseUploadResponse('<form action="https://doodstream.com/result"></form>'),
|
||||||
|
(err) => {
|
||||||
|
assert.equal(err.remoteCommitUncertain, true);
|
||||||
|
assert.equal(err.diagnostic.phase, 'upload-result-submit');
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|||||||
@@ -749,8 +749,9 @@ test('upload interval is enforced at the admitted upload start', async () => {
|
|||||||
const hosterSettings = settings('byse.sx', 1);
|
const hosterSettings = settings('byse.sx', 1);
|
||||||
hosterSettings['byse.sx'].timeIntervalSec = 1;
|
hosterSettings['byse.sx'].timeIntervalSec = 1;
|
||||||
const manager = new UploadManager(hosterSettings);
|
const manager = new UploadManager(hosterSettings);
|
||||||
manager._waitForInterval = async () => {
|
manager._waitForInterval = async (hoster, intervalMs, signal, acquireSlots) => {
|
||||||
events.push('interval');
|
events.push('interval');
|
||||||
|
await acquireSlots();
|
||||||
};
|
};
|
||||||
const batch = runBatch(manager, [
|
const batch = runBatch(manager, [
|
||||||
{ jobId: 'interval-first', file: firstPath, hoster: 'byse.sx', accountId: 'BYSE_ACCOUNT', apiKey: 'BYSE_KEY' },
|
{ jobId: 'interval-first', file: firstPath, hoster: 'byse.sx', accountId: 'BYSE_ACCOUNT', apiKey: 'BYSE_KEY' },
|
||||||
@@ -769,6 +770,61 @@ test('upload interval is enforced at the admitted upload start', async () => {
|
|||||||
assert.deepEqual(events.filter(event => event === 'interval'), ['interval', 'interval', 'interval']);
|
assert.deepEqual(events.filter(event => event === 'interval'), ['interval', 'interval', 'interval']);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test('an interval wait never occupies the global slot of another hoster', async () => {
|
||||||
|
let markIntervalWaiting;
|
||||||
|
let releaseInterval;
|
||||||
|
let markOtherStarted;
|
||||||
|
const intervalWaiting = new Promise(resolve => {
|
||||||
|
markIntervalWaiting = resolve;
|
||||||
|
});
|
||||||
|
const intervalGate = new Promise(resolve => {
|
||||||
|
releaseInterval = resolve;
|
||||||
|
});
|
||||||
|
const otherStarted = new Promise(resolve => {
|
||||||
|
markOtherStarted = resolve;
|
||||||
|
});
|
||||||
|
loadManager(async (hoster) => {
|
||||||
|
if (hoster === 'voe.sx') markOtherStarted();
|
||||||
|
const code = hoster === 'voe.sx' ? 'VOEINTERVAL1' : 'BYSEINTERVAL1';
|
||||||
|
return {
|
||||||
|
file_code: code,
|
||||||
|
download_url: hoster === 'voe.sx' ? `https://voe.sx/${code}` : `https://byse.sx/d/${code}`,
|
||||||
|
embed_url: hoster === 'voe.sx' ? `https://voe.sx/e/${code}` : `https://byse.sx/e/${code}`
|
||||||
|
};
|
||||||
|
});
|
||||||
|
const hosterSettings = {
|
||||||
|
...settings('byse.sx', 1),
|
||||||
|
...settings('voe.sx', 1)
|
||||||
|
};
|
||||||
|
hosterSettings['byse.sx'].timeIntervalSec = 1;
|
||||||
|
const manager = new UploadManager(hosterSettings, { parallelUploadCount: 1 });
|
||||||
|
const originalWait = manager._waitForInterval.bind(manager);
|
||||||
|
manager._waitForInterval = async (hoster, intervalMs, signal, acquireSlots) => {
|
||||||
|
if (hoster === 'byse.sx') {
|
||||||
|
markIntervalWaiting();
|
||||||
|
await intervalGate;
|
||||||
|
}
|
||||||
|
return originalWait(hoster, 0, signal, acquireSlots);
|
||||||
|
};
|
||||||
|
const batch = runBatch(manager, [
|
||||||
|
{ jobId: 'interval-waiting-hoster', file: firstPath, hoster: 'byse.sx', accountId: 'BYSE_ACCOUNT', apiKey: 'BYSE_KEY' },
|
||||||
|
{ jobId: 'interval-independent-hoster', file: distinctPath, hoster: 'voe.sx', accountId: 'VOE_ACCOUNT', apiKey: 'VOE_KEY' }
|
||||||
|
]);
|
||||||
|
|
||||||
|
await waitFor(intervalWaiting, 500, 'Configured interval did not start waiting');
|
||||||
|
let blockedError = null;
|
||||||
|
try {
|
||||||
|
await waitFor(otherStarted, 500, 'Interval wait occupied the global upload slot');
|
||||||
|
} catch (error) {
|
||||||
|
blockedError = error;
|
||||||
|
} finally {
|
||||||
|
releaseInterval();
|
||||||
|
}
|
||||||
|
const summary = await batch;
|
||||||
|
if (blockedError) throw blockedError;
|
||||||
|
assert.equal(summary.succeeded, 2);
|
||||||
|
});
|
||||||
|
|
||||||
for (const scenario of [
|
for (const scenario of [
|
||||||
{
|
{
|
||||||
label: 'VOE',
|
label: 'VOE',
|
||||||
|
|||||||
Reference in New Issue
Block a user