theseus/bundled-addons/pithos/core/autoflush.js
Local Dev d355bbbf87 Theseus: bundle Pithos 0.3.11
New installs carry the same Pithos the extension channel serves, including
drives on Silent Mode and the Sia account card.
2026-10-04 03:36:50 +02:00

81 lines
2.8 KiB
JavaScript

// Optional automatic flush of leftovers. s3d packs uploads into full slabs
// (about 40 MB) before sending them to Sia, so a small batch can sit in the
// local buffer indefinitely. When the user opts in, this watches s3d's upload
// stats and flushes once the pending set has stopped changing for N minutes,
// i.e. nobody is still uploading into it.
//
// The watcher only sees stats, so it also catches uploads made by other S3
// clients, not just the ones that went through Pithos.
export const DEFAULT_MINUTES = 10;
export const MIN_MINUTES = 1;
export const MAX_MINUTES = 24 * 60;
export function clampMinutes(m) {
const n = Math.round(Number(m));
if (!Number.isFinite(n)) return DEFAULT_MINUTES;
return Math.min(MAX_MINUTES, Math.max(MIN_MINUTES, n));
}
// settings: () => ({ enabled, minutes }); stats: async () => /stats/uploads;
// flush: async () => void; log: (line) => void.
export function createAutoFlush({ settings, stats, flush, log = () => {}, now = Date.now }) {
let signature = null; // what was pending last time we looked
let quietSince = 0;
let flushing = false;
let lastFlushAt = null;
let lastError = null;
let failures = 0; // in a row; each one doubles the wait, up to an hour
function reset() { signature = null; quietSince = 0; }
function waitMs(minutes) {
const base = clampMinutes(minutes) * 60_000;
return failures ? Math.max(base, Math.min(60 * 60_000, base * 2 ** failures)) : base;
}
async function tick() {
const { enabled, minutes } = settings();
if (!enabled || flushing) { if (!enabled) reset(); return; }
let x;
try { x = await stats(); } catch { reset(); return; } // s3d not reachable
if (!x.pendingObjects) { reset(); return; }
const sig = `${x.pendingObjects}:${x.pendingSize}`;
if (sig !== signature) { signature = sig; quietSince = now(); return; }
if (now() - quietSince < waitMs(minutes)) return;
flushing = true;
log(`uploading ${x.pendingObjects} leftover object(s) to Sia after ${clampMinutes(minutes)} min with no new uploads`);
try {
await flush();
lastFlushAt = now();
lastError = null;
failures = 0;
log('leftovers are on Sia');
reset();
} catch (e) {
lastError = e.message;
failures++;
log(`automatic flush failed: ${e.message}`);
quietSince = now(); // a broken indexer is retried less and less often
} finally {
flushing = false;
}
}
return {
tick,
status() {
const { enabled, minutes } = settings();
const m = clampMinutes(minutes);
return {
enabled: !!enabled,
minutes: m,
flushing,
lastFlushAt,
lastError,
// When the next automatic flush happens if nothing new arrives.
dueAt: enabled && signature && !flushing ? quietSince + waitMs(m) : null,
};
},
};
}