diff options
| author | Christophe Besson <cbesson@gmail.com> | 2026-09-25 11:28:31 +0200 |
|---|---|---|
| committer | Christophe Besson <cbesson@gmail.com> | 2026-09-25 12:07:49 +0200 |
| commit | b7bf11c077bdd165400cf8d80c5cd1ad4248d854 (patch) | |
| tree | e80590cd7b798e21a4ac663fd1e653a66e5f3b34 /packages/meshbay-hub/src/meshbay_hub/static/transport-upload.js | |
| parent | 6348698322e5bfcb80d0ef0dcd833ba4e6179640 (diff) | |
| download | meshbay-b7bf11c077bdd165400cf8d80c5cd1ad4248d854.tar.gz | |
refactor(client): split transport.js into classic scripts
transport.js keeps the core (connection, reconnect, leases, dispatch).
Chat, media, admin, upload and device methods move, cut as text, into
transport-*.js scripts that hand a class of their own to extendTransport,
which copies each method onto MeshBayTransport.prototype; the codec,
roster checks, node pins and the rewrap fan-out move as they were. Both
shells load them after transport.js. Every prototype member, class
property and top-level function has the same source text as before.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Diffstat (limited to 'packages/meshbay-hub/src/meshbay_hub/static/transport-upload.js')
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/static/transport-upload.js | 228 |
1 files changed, 228 insertions, 0 deletions
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/transport-upload.js b/packages/meshbay-hub/src/meshbay_hub/static/transport-upload.js new file mode 100644 index 0000000..9f4b531 --- /dev/null +++ b/packages/meshbay-hub/src/meshbay_hub/static/transport-upload.js @@ -0,0 +1,228 @@ +// Uploads: a whole file, sealed chunk by chunk with several in flight, and the +// folders it lands in. +// +// Methods of MeshBayTransport, copied onto its prototype by extendTransport +// (transport.js, which the shell loads first). + +extendTransport(class { + /** + * 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. + * + * `dir` names the folder to upload into, as a virtual path + * (`Media/Films/1999`) — where the sender is actually looking. The node + * resolves it against the group's own roots, which refuses `..`, absolute + * segments and anything escaping its root; it is a place among the group's + * folders, never a path on the operator's filesystem. + * + * `root` is the older, coarser form: the root's name and nothing below it. + * Kept because a node that predates `dir` reads it, and because Chat has no + * folder on screen to name. Omitting both leaves the node to pick, which it + * only does for a client old enough to have had one destination. + */ + async uploadFile(file, { chunkSize, onProgress, signal, root, dir, + tr = '' } = {}) { + // The same file twice at once would confuse the node, which keys its own + // upload state by folder and name — and would race for the same + // destination. The guard uses the same key: by name alone, a dropped folder + // holding a `cover.jpg` in two albums failed the second one for nothing. + const inFlightKey = `${dir || ''}/${file.name}`; + if (this._inFlightUploads.has(inFlightKey)) { + throw new Error(`${file.name} is already being uploaded`); + } + if (!this._gekRaw) throw new Error('This group has no key on this device'); + const C = window.MeshBayCrypto; + const groupId = (this._connectArgs && this._connectArgs.groupId) || ''; + this._inFlightUploads.add(inFlightKey); + const uploadId = _hex(crypto.getRandomValues(new Uint8Array(16))); + 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 = []; + const wake = () => { + acked += 1; + if (onProgress) onProgress(Math.min(file.size, acked * size), file.size); + const waiter = acks.shift(); + if (waiter) waiter(); + }; + // "Where am I?" — resolved by the node's answer to the probe chunk below, + // or by anything that says this node cannot answer it. + let settleProbe = null; + const probed = new Promise((r) => { settleProbe = r; }); + const answerProbe = (from) => { + if (!settleProbe) return false; + const done = settleProbe; + settleProbe = null; + done(from); + return true; + }; + this._uploaders.set(uploadId, (msg) => { + if (msg.type === 'error') { + // A node that predates the probe refuses its index. That is not a + // failure — it is the answer "start from the beginning", which is what + // this client did before there was anything to ask. + if (answerProbe(0)) return; + failure = new Error(msg.detail || 'Upload refused'); + wake(); + return; + } + // The ack is sealed too — `stored_as` and the folder it landed in name + // the operator's content. Opening it is what makes the result usable, so + // a failure here fails the upload rather than being swallowed: a chat + // attachment that cannot learn its stored name would point at nothing. + C.openGroup(this._gekRaw, 'upload', 'file_upload_ack', groupId, msg) + .then((plain) => { + const payload = msgpack_decode(plain); + if (payload.stored_as) stored = payload; + // Only the probe's answer carries this, so the two are told apart + // without trusting the index the node echoed back in clear. + if (typeof payload.resume_from === 'number') return answerProbe(payload.resume_from); + return false; + }) + .catch((e) => { + failure = new Error( + `The node's upload reply did not open under the group key (${e.message})`); + return false; + }) + // A probe's answer is not a chunk: waking here would credit the + // progress bar with a chunk that was never sent. + .then((wasProbe) => { if (!wasProbe) wake(); }); + }); + + const nextAck = () => new Promise(r => acks.push(r)); + + try { + // Ask before sending anything. An upload interrupted at 99% used to start + // again from zero, because the node kept its position on the connection + // that was lost — see `uploads.py`. The question goes inside the seal, as + // a chunk with no bytes, because naming the file on a clear message is + // exactly what sealing this path was for. + // Sealed first, spread second — the same shape as the chunk loop below, + // and not only for symmetry: `test_the_upload_itself_is_sealed` reads + // this call and fails if a filename appears in it, which is how it can + // tell a field outside the seal from one inside it. + const probeSealed = await C.sealGroup( + this._gekRaw, 'upload', 'file_upload', groupId, + msgpack_encode({ filename: file.name, data: new Uint8Array(0), + dir: dir || '', root: root || '' })); + this._send({ + type: 'file_upload', + v: '0.1', + upload_id: uploadId, + chunk_index: UPLOAD_PROBE_INDEX, + total_chunks: total, + ...(tr ? { tr } : {}), + ...probeSealed, + }); + // Bounded: a node that answers neither the probe nor its refusal must not + // leave an upload waiting for ever. Starting over is always safe. + let from = await Promise.race([ + probed, + new Promise((r) => setTimeout(() => { answerProbe(0); r(0); }, + UPLOAD_PROBE_TIMEOUT_MS)), + ]); + // Defensive: a node reporting a position at or past the end would have + // renamed the file and dropped its state, so this cannot happen — and if + // it does, sending everything again is the answer that cannot corrupt. + if (!(from > 0) || from >= total) from = 0; + if (from > 0) { + acked = from; + if (onProgress) onProgress(Math.min(file.size, from * size), file.size); + } + + for (let i = from; i < total; i++) { + if (signal && signal.aborted) throw _aborted(); + // Between two chunks, never inside one — the node refuses a chunk that + // is not the one it expects, so a position is the only thing worth + // remembering. Nothing is recorded here beyond that: the node holds the + // real position, and the probe above is what asks for it on the way + // back in, which makes resuming correct even across a reconnect. + if (signal && signal.paused) { + signal.resumeFrom = i; + const paused = new Error('Paused'); + paused.name = 'PausedError'; + throw paused; + } + // 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) { + if (signal && signal.aborted) throw _aborted(); + 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()); + // The name, the destination and the bytes go inside the seal together. + // Mirrors `file_upload_wire` in meshbay_common/protocol.py; only the + // fields the node routes on stay outside it. + const sealed = await C.sealGroup( + this._gekRaw, 'upload', 'file_upload', groupId, + msgpack_encode({ filename: file.name, data: buf, + dir: dir || '', root: root || '' })); + this._send({ + type: 'file_upload', + v: '0.1', + upload_id: uploadId, + chunk_index: i, + total_chunks: total, + ...(tr ? { tr } : {}), + ...sealed, + }); + } + while (acked < total) { + await nextAck(); + if (failure) throw failure; + } + } finally { + this._uploaders.delete(uploadId); + this._inFlightUploads.delete(inFlightKey); + } + return stored || {}; + } + + /** Create a directory under the current one. Any member may. */ + async createDirectory(dir, name) { + const msg = await this._sendAndWait({ + type: 'dir_create', v: '0.1', dir: dir || '', name, + }); + if (msg.type === 'error') throw new Error(msg.detail); + return msg; + } +}); + +// 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; +// "Where am I?", asked as an ordinary sealed upload chunk with no bytes rather +// than on a clear message. Mirrors UPLOAD_PROBE_INDEX in +// meshbay_common/protocol.py; the node writes nothing and answers with +// `resume_from`, and one that predates it refuses the index, which reads as +// "start from the beginning". +const UPLOAD_PROBE_INDEX = -1; +// How long to wait for that answer before assuming there is none. A node that +// answers neither the probe nor its refusal must not leave an upload waiting +// for ever, and starting over is always safe. +const UPLOAD_PROBE_TIMEOUT_MS = 5000; +const UPLOAD_BUFFER_HIGH = 1024 * 1024; |