summaryrefslogtreecommitdiffstats
path: root/packages/meshbay-hub/src/meshbay_hub/static/transfers.js
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-hub/src/meshbay_hub/static/transfers.js')
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/transfers.js203
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`;
+}