aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-hub/src/meshbay_hub/static/transport.js
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-hub/src/meshbay_hub/static/transport.js')
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/transport.js102
1 files changed, 89 insertions, 13 deletions
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/transport.js b/packages/meshbay-hub/src/meshbay_hub/static/transport.js
index d904286..e93c601 100644
--- a/packages/meshbay-hub/src/meshbay_hub/static/transport.js
+++ b/packages/meshbay-hub/src/meshbay_hub/static/transport.js
@@ -35,6 +35,13 @@ async function _pkEdFromSk(skPkcs8B64) {
return pad ? b64 + '='.repeat(4 - pad) : b64;
}
+// 48 KB is what fits comfortably in one SCTP message across stacks; the window is
+// what makes the rate independent of the round trip. 32 × 48 KB = 1.5 MB in
+// flight, which saturates any path up to roughly 100 Mb/s at 100 ms.
+const UPLOAD_CHUNK_SIZE = 48 * 1024;
+const UPLOAD_WINDOW = 32;
+const UPLOAD_BUFFER_HIGH = 1024 * 1024;
+
const JOIN_REFUSALS = {
code_required: 'This node does not know this browser yet. Ask the node operator '
+ 'for a pairing code (meshbay-node operator pair).',
@@ -66,6 +73,7 @@ class MeshBayTransport {
this._onStreamData = null;
this._onStreamEnd = null;
this._onIndexSync = null;
+ this._onUploadAck = null;
}
get connected() { return this._connected; }
@@ -489,19 +497,79 @@ class MeshBayTransport {
this._send({ type: 'stream_req', v: '0.1', file_id: fileId });
}
- async uploadChunk(filename, chunkIndex, totalChunks, data) {
- // The node decides where this lands (uploads/) and under what name — it
- // finds a free one rather than replacing anything. The ack says which.
- const msg = await this._sendAndWait({
- type: 'file_upload',
- v: '0.1',
- filename,
- chunk_index: chunkIndex,
- total_chunks: totalChunks,
- data: data,
- });
- if (msg.type === 'error') throw new Error(msg.detail);
- return msg;
+ /**
+ * Push a whole file, several chunks in flight at once.
+ *
+ * One chunk per round trip is 48 KB of throughput per RTT no matter how much
+ * bandwidth there is: 4.8 MB/s on a 10 ms path, 480 KB/s on a 100 ms one, and
+ * the sender is idle for almost all of it — which also keeps SCTP's congestion
+ * window shut, so the transport never gets a chance to speed up either. A
+ * window of chunks makes the rate depend on bandwidth rather than distance.
+ *
+ * Order is not at risk: a DataChannel is ordered and reliable by default, and
+ * the node refuses any chunk that is not the one it expects next.
+ *
+ * The node decides where this lands (uploads/) and under what name — it finds a
+ * free one rather than replacing anything. The ack says which, and that is what
+ * this returns.
+ */
+ async uploadFile(file, { chunkSize, onProgress } = {}) {
+ // One at a time: the acks are matched by arrival, so two uploads sharing the
+ // channel would credit each other's progress and finish at the wrong moment.
+ if (this._onUploadAck) throw new Error('Another upload is already running');
+ const size = chunkSize || UPLOAD_CHUNK_SIZE;
+ const total = Math.max(1, Math.ceil(file.size / size));
+ let acked = 0;
+ let stored = null;
+ let failure = null;
+
+ const acks = [];
+ this._onUploadAck = (msg) => {
+ if (msg.type === 'error') {
+ failure = new Error(msg.detail || 'Upload refused');
+ } else if (msg.stored_as) {
+ stored = msg;
+ }
+ acked += 1;
+ if (onProgress) onProgress(Math.min(file.size, acked * size), file.size);
+ const waiter = acks.shift();
+ if (waiter) waiter();
+ };
+
+ const nextAck = () => new Promise(r => acks.push(r));
+
+ try {
+ for (let i = 0; i < total; i++) {
+ // Backpressure: without it the whole file lands in the browser's send
+ // buffer in seconds and the progress bar becomes a work of fiction.
+ while (this._channel && this._channel.bufferedAmount > UPLOAD_BUFFER_HIGH) {
+ await new Promise(r => setTimeout(r, 20));
+ }
+ while (i - acked >= UPLOAD_WINDOW) {
+ await nextAck();
+ if (failure) throw failure;
+ }
+ if (failure) throw failure;
+
+ const buf = new Uint8Array(
+ await file.slice(i * size, (i + 1) * size).arrayBuffer());
+ this._send({
+ type: 'file_upload',
+ v: '0.1',
+ filename: file.name,
+ chunk_index: i,
+ total_chunks: total,
+ data: buf,
+ });
+ }
+ while (acked < total) {
+ await nextAck();
+ if (failure) throw failure;
+ }
+ } finally {
+ this._onUploadAck = null;
+ }
+ return stored || {};
}
/** Create a directory under the current one. Any member may. */
@@ -681,6 +749,14 @@ class MeshBayTransport {
}
_dispatch(msg) {
+ // While an upload is in flight the acks are its own, and there are many of
+ // them: they must not be handed to whatever request happens to be oldest in
+ // the pending map.
+ if (this._onUploadAck
+ && (msg.type === 'file_upload_ack' || msg.type === 'error')) {
+ this._onUploadAck(msg);
+ return;
+ }
if (msg.type === 'chat_msg' && this._onChat) {
this._onChat(msg);
return;