// 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, }; }, }; }