diff options
Diffstat (limited to 'packages/meshbay-hub/src/meshbay_hub/static/transfers.js')
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/static/transfers.js | 203 |
1 files changed, 203 insertions, 0 deletions
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/transfers.js b/packages/meshbay-hub/src/meshbay_hub/static/transfers.js new file mode 100644 index 0000000..1ad1ec5 --- /dev/null +++ b/packages/meshbay-hub/src/meshbay_hub/static/transfers.js @@ -0,0 +1,203 @@ +/** + * Transfers that outlive the page that started them. + * + * Downloads and uploads used to be state inside GroupPage, which meant leaving + * a group killed them — the component unmounted, its effect closed the + * DataChannel, and a half-written file was all you had. They live here instead: + * a module-level store that nothing unmounts, with the group page as one of + * several possible views onto it. + * + * Two consequences worth stating, because they are the reason this exists: + * + * - The transport cannot be closed just because a page went away. A group + * page hands its transport over with `releaseWhenIdle()`, and the last + * transfer using it closes it. + * - Signing out is different from navigating. It cancels everything and + * closes what it was using, because the tokens those transfers are running + * on are about to stop being ours. + * + * No browser globals: exercised under Node by + * packages/meshbay-hub/tests/test_transfers.py. + */ + +const SPEED_WINDOW_MS = 5000; + +let _nextId = 1; + +export class TransferStore { + constructor(now = () => Date.now()) { + this._now = now; + this._items = []; + this._subs = new Set(); + this._releasing = new Set(); + } + + subscribe(fn) { + this._subs.add(fn); + return () => this._subs.delete(fn); + } + + _emit() { + for (const fn of this._subs) fn(this.list()); + } + + /** + * What a view needs to render, as plain data — never the internals, so a + * render cannot accidentally hold a transport alive. + */ + list() { + return this._items.map(it => ({ + id: it.id, + kind: it.kind, + name: it.name, + total: it.total, + done: it.done, + status: it.status, + error: it.error || '', + speed: this._speed(it), + percent: it.total ? Math.min(100, Math.round(it.done / it.total * 100)) : 0, + })); + } + + get active() { + return this._items.filter(it => it.status === 'running').length; + } + + _speed(it) { + // Over a window rather than since the start: a transfer that stalls should + // read as slow immediately, not as its own historical average. + const s = it.samples; + if (s.length < 2) return 0; + const dt = (s[s.length - 1].t - s[0].t) / 1000; + if (dt <= 0) return 0; + return (s[s.length - 1].done - s[0].done) / dt; + } + + /** + * Start a transfer. + * + * `run` receives `{ signal, onProgress }`. It must poll `signal.aborted` — a + * cancel that only sets a flag nobody reads is a button that lies. + */ + start({ kind, name, total = 0, transport = null, run }) { + const item = { + id: _nextId++, + kind, name, total, transport, + done: 0, + status: 'running', + error: '', + samples: [{ t: this._now(), done: 0 }], + signal: { aborted: false }, + }; + this._items.push(item); + this._emit(); + + const onProgress = (done, total) => { + item.done = done; + if (total) item.total = total; + const t = this._now(); + item.samples.push({ t, done }); + while (item.samples.length > 2 && t - item.samples[0].t > SPEED_WINDOW_MS) { + item.samples.shift(); + } + this._emit(); + }; + + const finish = (status, error = '') => { + item.status = status; + item.error = error; + this._emit(); + this._maybeRelease(item.transport); + }; + + const promise = Promise.resolve() + .then(() => run({ signal: item.signal, onProgress })) + .then(() => { + if (item.signal.aborted) finish('cancelled'); + else { + if (item.total) item.done = item.total; + finish('done'); + } + }) + .catch(err => { + if (item.signal.aborted || err.name === 'AbortError') finish('cancelled'); + else finish('failed', err.message || String(err)); + }); + + item.promise = promise; + return item.id; + } + + cancel(id) { + const item = this._items.find(it => it.id === id); + if (!item || item.status !== 'running') return; + item.signal.aborted = true; + // Marked at once. The work stops when it next looks, but a cancelled + // transfer should not keep reporting progress in the meantime. + item.status = 'cancelled'; + this._emit(); + this._maybeRelease(item.transport); + } + + cancelAll() { + for (const it of this._items) { + if (it.status === 'running') this.cancel(it.id); + } + } + + /** Drop everything finished, keeping what is still running. */ + clearFinished() { + this._items = this._items.filter(it => it.status === 'running'); + this._emit(); + } + + _busy(transport) { + return this._items.some( + it => it.transport === transport && it.status === 'running'); + } + + /** + * The group page is going away. Close its transport once nothing is using it, + * which may be now or may be in twenty minutes. + */ + releaseWhenIdle(transport) { + if (!transport) return; + this._releasing.add(transport); + this._maybeRelease(transport); + } + + _maybeRelease(transport) { + if (!transport || !this._releasing.has(transport)) return; + if (this._busy(transport)) return; + this._releasing.delete(transport); + try { + transport.onIndexSync = null; + transport.close(); + } catch { /* already gone */ } + } + + /** Signing out: stop everything and let go of every transport. */ + reset() { + this.cancelAll(); + for (const transport of [...this._releasing]) { + this._releasing.delete(transport); + try { transport.close(); } catch { /* already gone */ } + } + for (const it of this._items) { + if (it.transport) { + try { it.transport.close(); } catch { /* already gone */ } + } + } + this._items = []; + this._emit(); + } +} + +export const transfers = new TransferStore(); + +/** Human-readable rate, for a widget that updates several times a second. */ +export function formatSpeed(bytesPerSecond) { + if (!bytesPerSecond || bytesPerSecond < 1) return ''; + if (bytesPerSecond < 1024 * 1024) return `${Math.round(bytesPerSecond / 1024)} KB/s`; + return `${(bytesPerSecond / (1024 * 1024)).toFixed(1)} MB/s`; +} |