Harden source cleanup durability barriers

Separate provisional upload success from persistently confirmed cleanup state, invalidate prior confirmations before retries, and promote completion only inside the final queue persistence handshake. Preserve source files when history or queue persistence fails and add regression coverage for restart, retry, rollback, and real Electron finalization paths.
This commit is contained in:
Sucukdeluxe
2026-08-13 19:45:31 +02:00
parent bfb3a39fed
commit 67978cb81f
7 changed files with 540 additions and 91 deletions
+107 -21
View File
@@ -3,6 +3,7 @@
const pathApi = typeof require === 'function' ? require('path') : null;
const protectedStatuses = new Set(['done', 'error', 'aborted', 'skipped']);
const metadataVersion = 2;
function normalizeFile(file, platform) {
const value = typeof file === 'string' ? file.trim() : '';
@@ -59,18 +60,26 @@
return uniqueHosters(values);
}
function completedHosters(jobs, requiredHosters) {
function confirmedHosters(jobs, requiredHosters) {
const values = [];
for (const job of jobs) {
if (Array.isArray(job.sourceCleanupCompletedHosters)) {
values.push(...job.sourceCleanupCompletedHosters);
if (job.sourceCleanupMetadataVersion === metadataVersion && Array.isArray(job.sourceCleanupConfirmedHosters)) {
values.push(...job.sourceCleanupConfirmedHosters);
}
}
const confirmed = new Set(uniqueHosters(values));
return requiredHosters.filter((hoster) => confirmed.has(hoster));
}
function provisionalHosters(jobs, requiredHosters) {
const values = [];
for (const job of jobs) {
if (job.status === 'done') values.push(job.hoster);
if (job.sourceCleanupMetadataVersion === metadataVersion && Array.isArray(job.sourceCleanupProvisionalHosters)) {
values.push(...job.sourceCleanupProvisionalHosters);
}
}
const completed = new Set(uniqueHosters(values));
return requiredHosters.filter((hoster) => completed.has(hoster));
const provisional = new Set(uniqueHosters(values));
return requiredHosters.filter((hoster) => provisional.has(hoster));
}
function storedToken(jobs) {
@@ -90,12 +99,15 @@
return null;
}
function assignMetadata(jobs, token, requiredHosters, completed, fingerprint, touchedJobs, touchedSet) {
function assignMetadata(jobs, token, requiredHosters, confirmed, provisional, fingerprint, touchedJobs, touchedSet) {
for (const job of jobs) {
job.sourceCleanupMetadataVersion = metadataVersion;
job.sourceCleanupToken = token;
job.sourceCleanupRequiredHosters = [...requiredHosters];
job.sourceCleanupCompletedHosters = [...completed];
job.sourceCleanupConfirmedHosters = [...confirmed];
job.sourceCleanupProvisionalHosters = [...provisional];
job.sourceCleanupFingerprint = cloneFingerprint(fingerprint);
delete job.sourceCleanupCompletedHosters;
if (!touchedSet.has(job)) {
touchedSet.add(job);
touchedJobs.push(job);
@@ -108,7 +120,10 @@
const touchedJobs = [];
const touchedSet = new Set();
const preparedFiles = new Set();
if (!Array.isArray(queueJobs) || !Array.isArray(jobsToStart)) return { groups, touchedJobs };
const revokedHosters = [];
const revokedSet = new Set();
if (!Array.isArray(queueJobs) || !Array.isArray(jobsToStart)) return { groups, touchedJobs, revokedHosters };
const currentRoundJobs = new Set(jobsToStart);
for (const selectedJob of jobsToStart) {
const file = normalizeFile(selectedJob && selectedJob.file, platform);
@@ -124,14 +139,26 @@
...persistedRequired,
...siblings.map((job) => job.hoster)
]);
const completed = completedHosters(siblings, requiredHosters);
const startedHosters = new Set(uniqueHosters(
siblings.filter((job) => currentRoundJobs.has(job)).map((job) => job.hoster)
));
const storedConfirmed = confirmedHosters(siblings, requiredHosters);
const confirmed = storedConfirmed.filter((hoster) => !startedHosters.has(hoster));
const provisional = provisionalHosters(siblings, requiredHosters)
.filter((hoster) => !startedHosters.has(hoster));
for (const hoster of storedConfirmed) {
if (!startedHosters.has(hoster) || revokedSet.has(hoster)) continue;
revokedSet.add(hoster);
revokedHosters.push(hoster);
}
const fingerprint = storedFingerprint(siblings);
assignMetadata(
siblings,
token,
requiredHosters,
completed,
confirmed,
provisional,
fingerprint,
touchedJobs,
touchedSet
@@ -141,33 +168,84 @@
token,
file: selectedJob.file,
requiredHosters: [...requiredHosters],
completedHosters: [...completed],
confirmedHosters: [...confirmed],
fingerprint: cloneFingerprint(fingerprint),
jobs: siblings.map((job) => ({
jobId: job.id,
hoster: normalizeHoster(job.hoster),
status: job.status
status: job.status,
currentRound: currentRoundJobs.has(job)
}))
});
}
return { groups, touchedJobs };
return { groups, touchedJobs, revokedHosters };
}
function markCompleted(queueJobs, job, platform) {
const siblings = relatedJobs(queueJobs, job, platform);
if (siblings.length === 0) return [];
const requiredHosters = storedRequiredHosters(siblings);
const completed = new Set(completedHosters(siblings, requiredHosters));
const provisional = new Set(provisionalHosters(siblings, requiredHosters));
const hoster = normalizeHoster(job.hoster);
if (requiredHosters.includes(hoster)) completed.add(hoster);
const orderedCompleted = requiredHosters.filter((required) => completed.has(required));
if (requiredHosters.includes(hoster)) provisional.add(hoster);
const orderedConfirmed = confirmedHosters(siblings, requiredHosters);
const orderedProvisional = requiredHosters.filter((required) => provisional.has(required));
for (const sibling of siblings) {
sibling.sourceCleanupCompletedHosters = [...orderedCompleted];
sibling.sourceCleanupMetadataVersion = metadataVersion;
sibling.sourceCleanupConfirmedHosters = [...orderedConfirmed];
sibling.sourceCleanupProvisionalHosters = [...orderedProvisional];
delete sibling.sourceCleanupCompletedHosters;
}
return siblings;
}
async function persistRoundCompletions(queueJobs, options = {}) {
if (!Array.isArray(queueJobs) || typeof options.persist !== 'function') return false;
const historyPersisted = options.historyPersisted === true;
const groupsByToken = new Map();
for (const job of queueJobs) {
if (!job || typeof job.sourceCleanupToken !== 'string' || !job.sourceCleanupToken) continue;
if (!groupsByToken.has(job.sourceCleanupToken)) groupsByToken.set(job.sourceCleanupToken, []);
groupsByToken.get(job.sourceCleanupToken).push(job);
}
const snapshots = [];
for (const jobs of groupsByToken.values()) {
const requiredHosters = storedRequiredHosters(jobs);
const confirmed = confirmedHosters(jobs, requiredHosters);
const provisional = provisionalHosters(jobs, requiredHosters);
const promoted = historyPersisted
? uniqueHosters([...confirmed, ...provisional])
: confirmed;
const orderedPromoted = requiredHosters.filter((hoster) => promoted.includes(hoster));
for (const job of jobs) {
snapshots.push({
job,
confirmedHosters: job.sourceCleanupMetadataVersion === metadataVersion
? uniqueHosters(job.sourceCleanupConfirmedHosters)
: []
});
job.sourceCleanupMetadataVersion = metadataVersion;
job.sourceCleanupConfirmedHosters = [...orderedPromoted];
job.sourceCleanupProvisionalHosters = [];
delete job.sourceCleanupCompletedHosters;
}
}
let persisted = false;
try {
persisted = (await options.persist()) === true;
} catch {}
if (!persisted) {
for (const snapshot of snapshots) {
snapshot.job.sourceCleanupMetadataVersion = metadataVersion;
snapshot.job.sourceCleanupConfirmedHosters = [...snapshot.confirmedHosters];
snapshot.job.sourceCleanupProvisionalHosters = [];
delete snapshot.job.sourceCleanupCompletedHosters;
}
}
return persisted;
}
function removeRequirement(queueJobs, job, platform) {
if (!job || protectedStatuses.has(job.status)) return [];
const siblings = relatedJobs(queueJobs, job, platform);
@@ -176,13 +254,20 @@
const required = Array.isArray(sibling.sourceCleanupRequiredHosters)
? sibling.sourceCleanupRequiredHosters
: [];
const completed = Array.isArray(sibling.sourceCleanupCompletedHosters)
? sibling.sourceCleanupCompletedHosters
const confirmed = sibling.sourceCleanupMetadataVersion === metadataVersion && Array.isArray(sibling.sourceCleanupConfirmedHosters)
? sibling.sourceCleanupConfirmedHosters
: [];
const provisional = sibling.sourceCleanupMetadataVersion === metadataVersion && Array.isArray(sibling.sourceCleanupProvisionalHosters)
? sibling.sourceCleanupProvisionalHosters
: [];
sibling.sourceCleanupMetadataVersion = metadataVersion;
sibling.sourceCleanupRequiredHosters = uniqueHosters(required)
.filter((hoster) => hoster !== removedHoster);
sibling.sourceCleanupCompletedHosters = uniqueHosters(completed)
sibling.sourceCleanupConfirmedHosters = uniqueHosters(confirmed)
.filter((hoster) => hoster !== removedHoster);
sibling.sourceCleanupProvisionalHosters = uniqueHosters(provisional)
.filter((hoster) => hoster !== removedHoster);
delete sibling.sourceCleanupCompletedHosters;
}
return siblings;
}
@@ -211,6 +296,7 @@
const api = {
prepareGroups,
markCompleted,
persistRoundCompletions,
removeRequirement,
applyFingerprints
};
+43 -23
View File
@@ -83,30 +83,29 @@ function createSourceFileCleanup(options) {
function createManifest(input, canonicalFile) {
const requiredHosters = uniqueStrings(input.requiredHosters);
const completedHosters = new Set(uniqueStrings(input.completedHosters));
const confirmedHosters = new Set(
uniqueStrings(input.confirmedHosters).filter((hoster) => requiredHosters.includes(hoster))
);
const jobs = new Map();
for (const job of Array.isArray(input.jobs) ? input.jobs : []) {
if (!job || typeof job.jobId !== 'string' || typeof job.hoster !== 'string') continue;
const status = completedHosters.has(job.hoster) ? 'done' : normalizeStatus(job.status);
const currentRound = job.currentRound !== false;
const status = currentRound ? 'pending' : normalizeStatus(job.status);
jobs.set(job.jobId, Object.freeze({
jobId: job.jobId,
hoster: job.hoster,
status
status,
currentRound
}));
}
for (const hoster of completedHosters) {
if (requiredHosters.includes(hoster)) continue;
completedHosters.delete(hoster);
}
return {
token: input.token || input.sourceCleanupToken,
file: path.resolve(input.file),
canonicalFile,
requiredHosters: Object.freeze(requiredHosters),
completedHosters,
confirmedHosters,
jobs,
suppliedFingerprint: isFingerprint(input.fingerprint) ? cloneFingerprint(input.fingerprint) : null,
fingerprint: null,
@@ -152,16 +151,32 @@ function createSourceFileCleanup(options) {
...existing.requiredHosters,
...input.requiredHosters
]));
for (const hoster of uniqueStrings(input.completedHosters)) {
if (existing.requiredHosters.includes(hoster)) existing.completedHosters.add(hoster);
const incomingJobs = Array.isArray(input.jobs) ? input.jobs : [];
const currentRoundHosters = new Set(
incomingJobs
.filter((job) => job && typeof job.hoster === 'string' && job.currentRound !== false)
.map((job) => job.hoster)
);
for (const hoster of currentRoundHosters) existing.confirmedHosters.delete(hoster);
for (const hoster of uniqueStrings(input.confirmedHosters)) {
if (existing.requiredHosters.includes(hoster) && !currentRoundHosters.has(hoster)) {
existing.confirmedHosters.add(hoster);
}
}
for (const job of Array.isArray(input.jobs) ? input.jobs : []) {
for (const job of incomingJobs) {
if (!job || typeof job.jobId !== 'string' || typeof job.hoster !== 'string') continue;
const previous = existing.jobs.get(job.jobId);
const status = existing.completedHosters.has(job.hoster)
? 'done'
const incomingCurrentRound = job.currentRound !== false;
const currentRound = Boolean((previous && previous.currentRound) || incomingCurrentRound);
const status = incomingCurrentRound
? 'pending'
: normalizeStatus(previous ? previous.status : job.status);
existing.jobs.set(job.jobId, Object.freeze({ jobId: job.jobId, hoster: job.hoster, status }));
existing.jobs.set(job.jobId, Object.freeze({
jobId: job.jobId,
hoster: job.hoster,
status,
currentRound
}));
}
fingerprints[token] = cloneFingerprint(existing.fingerprint);
continue;
@@ -194,10 +209,9 @@ function createSourceFileCleanup(options) {
const manifest = groups.get(token);
if (!manifest || manifest.finalizationPromise) return false;
const job = manifest.jobs.get(event.jobId);
if (!job || job.hoster !== event.hoster) return false;
if (!job || !job.currentRound || job.hoster !== event.hoster) return false;
if (typeof event.file === 'string' && canonicalize(event.file) !== manifest.canonicalFile) return false;
manifest.jobs.set(event.jobId, Object.freeze({ ...job, status: event.status }));
if (event.status === 'done') manifest.completedHosters.add(event.hoster);
return true;
}
@@ -206,7 +220,7 @@ function createSourceFileCleanup(options) {
for (const manifest of groups.values()) {
if (manifest.finalizationPromise) continue;
const job = manifest.jobs.get(jobId);
if (!job) continue;
if (!job || !job.currentRound) continue;
manifest.jobs.set(jobId, Object.freeze({ ...job, status: 'skipped' }));
changed = true;
}
@@ -216,11 +230,17 @@ function createSourceFileCleanup(options) {
function blockingStatuses(manifest) {
const blocking = [];
for (const hoster of manifest.requiredHosters) {
if (manifest.completedHosters.has(hoster)) continue;
const statuses = [...manifest.jobs.values()]
.filter((job) => job.hoster === hoster)
.map((job) => job.status);
if (statuses.includes('done')) continue;
const jobs = [...manifest.jobs.values()].filter((job) => job.hoster === hoster);
const currentJobs = jobs.filter((job) => job.currentRound);
if (currentJobs.length > 0) {
const currentBlocker = currentJobs.find((job) => job.status !== 'done');
if (!currentBlocker) continue;
const status = currentJobs.map((job) => job.status).find((value) => value !== 'pending') || 'pending';
blocking.push({ hoster, status });
continue;
}
if (manifest.confirmedHosters.has(hoster)) continue;
const statuses = jobs.map((job) => job.status).filter((status) => status !== 'done');
const status = statuses.find((value) => value !== 'pending') || 'pending';
blocking.push({ hoster, status });
}