82 lines
2.8 KiB
JavaScript
82 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,
|
||
|
|
};
|
||
|
|
},
|
||
|
|
};
|
||
|
|
}
|