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.js306
1 files changed, 300 insertions, 6 deletions
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/transport.js b/packages/meshbay-hub/src/meshbay_hub/static/transport.js
index c4a24c5..0b2fed5 100644
--- a/packages/meshbay-hub/src/meshbay_hub/static/transport.js
+++ b/packages/meshbay-hub/src/meshbay_hub/static/transport.js
@@ -40,6 +40,16 @@ async function _pkEdFromSk(skPkcs8B64) {
// 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;
// Segments of 256 KB: 24 in flight is 6 MB, enough to keep playback fed over a
@@ -251,8 +261,12 @@ window.addEventListener('hashchange', () => {
// The `v: '0.1'` on every other message in this file is the historical value
// and is read by nothing; it is left alone deliberately. The range is
// negotiated once, at the start, not restated per message.
-const MNP_V = '2.0';
-const MNP_V_MIN = '1.0';
+const MNP_V = '3.0';
+// Raised with it on the 3.0 flag day. A node older than 3.0 cannot grant the
+// lease this client opens for every download and upload, so talking to one
+// would mean every transfer failing for a reason the person cannot act on.
+// Refusing it at the handshake says so once, in a sentence.
+const MNP_V_MIN = '3.0';
// Codes a NODE sends us, in its own vocabulary (meshbay_common/handshake.py's
// check_version): `version_too_old` means *we* are too old for it,
@@ -283,6 +297,129 @@ const JOIN_REFUSALS = {
group_mismatch: 'The node refused a request naming a different group.',
};
+/**
+ * One transfer's slot on the node, from this side.
+ *
+ * The contract the transfer store depends on: `acquire()` resolves when the
+ * node has granted the slot (immediately on a node that hands out none), and
+ * `release(reason)` gives it back exactly once. Nothing else in the client
+ * speaks to the node about slots.
+ *
+ * Two things here exist only because a queue can lie, and both are the
+ * difference between "waiting" and "waiting for ever":
+ *
+ * - **the watchdog.** A grant is pushed, not polled, so a lost push leaves
+ * this side waiting on a node that believes it has started. Re-asking is
+ * free — the node is idempotent on `tr` — and it is the only thing that
+ * recovers a message that did not arrive.
+ * - **release is idempotent and unconditional.** A slot given back twice
+ * costs nothing; one never given back is a member who cannot transfer
+ * again until a timeout the node runs on its own.
+ */
+const LEASE_WATCHDOG_MS = 60000;
+
+class Lease {
+ constructor(transport, tr, kind, bytes, chunks, onState) {
+ this.transport = transport;
+ this.tr = tr;
+ this.kind = kind;
+ this.bytes = bytes;
+ this.chunks = chunks;
+ this.state = 'opening';
+ this.ahead = 0;
+ this.closed = false;
+ this._onState = onState;
+ this._granted = null;
+ this._watchdog = 0;
+ this._wait = new Promise((resolve) => { this._granted = resolve; });
+ }
+
+ /** No slots on this node: behave as though one was granted at once. */
+ _skip() {
+ this.state = 'granted';
+ this._granted();
+ }
+
+ _request() {
+ // A closed channel is not a failure here, and must not throw: the transport
+ // reconnects on its own, `_reopenTransfers` re-asks for every live lease
+ // when it does, and the watchdog below asks again meanwhile.
+ //
+ // This is the same tolerance `_fetchChunkResilient` already gives a chunk
+ // request — and before leases existed, a chunk request was the first thing
+ // to touch the channel, so a download started on a briefly dead connection
+ // simply retried. Asking for a slot first made `_send` the first contact
+ // and threw "DataChannel not open (state: closed)" out of `downloadEntry`,
+ // where nothing catches it: a download that used to recover became an
+ // error with no row in the widget to show it. Found by downloading a file
+ // right after a connection dropped.
+ try {
+ this.transport._send({
+ type: 'transfer_open', v: '0.1', tr: this.tr, kind: this.kind,
+ bytes: this.bytes, chunks: this.chunks,
+ });
+ } catch (err) {
+ console.warn('[MeshBay] could not ask for a slot yet:', err.message);
+ }
+ this._arm();
+ }
+
+ _arm() {
+ clearTimeout(this._watchdog);
+ if (this.closed || this.state === 'granted') return;
+ this._watchdog = setTimeout(() => {
+ if (this.closed || this.state === 'granted') return;
+ console.warn('[MeshBay] no answer for transfer', this.tr.slice(0, 8),
+ '- asking again');
+ this._request();
+ }, LEASE_WATCHDOG_MS);
+ }
+
+ _apply(msg) {
+ if (this.closed) return;
+ this.state = msg.state;
+ this.ahead = msg.ahead || 0;
+ this.used = msg.used;
+ this.cap = msg.cap;
+ if (msg.state === 'granted') {
+ clearTimeout(this._watchdog);
+ this._granted();
+ } else if (msg.state === 'closed') {
+ // The node ended it: reclaimed as idle, or revoked. Not an error here —
+ // whoever is running the transfer finds out through its own failure — but
+ // the slot is gone and asking again is the only way back.
+ clearTimeout(this._watchdog);
+ } else {
+ this._arm();
+ }
+ if (this._onState) {
+ try { this._onState(this); } catch (e) {
+ console.error('[MeshBay] lease state handler threw:', e);
+ }
+ }
+ }
+
+ /** Resolves once the node has granted the slot. */
+ acquire() { return this._wait; }
+
+ /**
+ * Give the slot back. Safe to call twice, and safe on a dead transport: a
+ * lease that is not released is a member who cannot start another transfer
+ * until the node times it out, so this must never be conditional on anything.
+ */
+ release(reason = 'done') {
+ if (this.closed) return;
+ this.closed = true;
+ clearTimeout(this._watchdog);
+ this.transport._leases.delete(this.tr);
+ if (!this.transport.supportsTransferSlots) return;
+ try {
+ this.transport._send({ type: 'transfer_close', v: '0.1', tr: this.tr,
+ reason });
+ } catch { /* the connection is gone, and so is the lease with it */ }
+ }
+}
+
class MeshBayTransport {
constructor(hubUrl, accessToken) {
this._hubUrl = hubUrl;
@@ -315,6 +452,13 @@ class MeshBayTransport {
// Names, not ids: the "already being uploaded" guard is about the file the
// caller passed, and two `uploadFile` calls for one file draw two ids.
this._inFlightUploads = new Set();
+ // tr → Lease. A transfer's slot on the node, from the client's side.
+ this._leases = new Map();
+ // Set from the handshake ack: a node that answers with `transfer_limits`
+ // speaks transfer slots. Used instead of a timeout, because "no answer
+ // yet" and "this node will never answer" are indistinguishable in time and
+ // guessing wrong either stalls every download or defeats the cap.
+ this._transferLimits = null;
// Set once close() runs — stops the automatic reconnect from firing on a
// connection the caller tore down on purpose (leaving the group, page
// unload), which would otherwise race back in right as everything else
@@ -377,6 +521,20 @@ class MeshBayTransport {
get nodeVersion() { return this._nodeVersion || ''; }
/**
+ * Whether this node hands out transfer slots.
+ *
+ * Read from the handshake ack rather than from the MNP version: the caps
+ * shipped before the version bump that will make leases compulsory, so for
+ * now a node either answers with `transfer_limits` or it predates all of
+ * this. A node that does not is asked for nothing and enforces nothing —
+ * every download behaves exactly as it did.
+ */
+ get supportsTransferSlots() { return this._transferLimits !== null; }
+
+ /** This member's own caps in this group, or null when the node said nothing. */
+ get transferLimits() { return this._transferLimits; }
+
+ /**
* Whether the node speaks the per-root and per-app operations MNP 1.1 added:
* `root_update`/`root_eject`/`root_plug`, `app_directories`,
* `chat_directory`, `chat_link_preview`.
@@ -880,6 +1038,7 @@ class MeshBayTransport {
delete ack.nonce;
delete ack.ct;
Object.assign(ack, config);
+ this._transferLimits = ack.transfer_limits || null;
// Tell the node which of this account's devices is on this connection.
// Deliberately after the ack, and gated on the node's own version rather
@@ -1016,6 +1175,10 @@ class MeshBayTransport {
}
trace('reconnect_ok', { attempt: this._reconnectAttempts });
console.log('[MeshBay] Reconnected after', this._reconnectAttempts, 'attempt(s)');
+ // Before the caller's own hook: a transfer that resumes mid-chunk must
+ // have asked for its slot back first, or its next `file_req` carries a
+ // `tr` the node has never heard of.
+ this._reopenTransfers();
if (this._onReconnected) {
try { this._onReconnected(); } catch (e) {
console.error('[MeshBay] onReconnected handler threw:', e);
@@ -1100,12 +1263,52 @@ class MeshBayTransport {
return msg;
}
- async fetchChunk(fileId, chunkIndex) {
+ // ── Transfer slots ─────────────────────────────────────────────────────────
+
+ /**
+ * Ask the node for a slot, and wait until it says yes.
+ *
+ * `tr` is drawn here, not by the node — 16 random bytes, exactly like
+ * `upload_id` — which is what makes re-opening after a reconnect idempotent
+ * rather than a second charge against the member's cap.
+ *
+ * On a node that predates transfer slots this resolves at once and costs
+ * nothing: there is no cap to respect and no message that would be
+ * understood.
+ */
+ openTransfer({ kind = 'download', bytes = 0, chunks = 0, onState = null } = {}) {
+ const tr = _hex(crypto.getRandomValues(new Uint8Array(16)));
+ const lease = new Lease(this, tr, kind, bytes, chunks, onState);
+ if (!this.supportsTransferSlots) {
+ lease._skip();
+ return lease;
+ }
+ this._leases.set(tr, lease);
+ lease._request();
+ return lease;
+ }
+
+ /** Re-ask for every live lease. Called after a reconnect. */
+ _reopenTransfers() {
+ if (!this.supportsTransferSlots) return;
+ for (const lease of this._leases.values()) {
+ // The node lost the lease with the session, so this is a fresh request
+ // for the same `tr` — which the node treats as the same transfer rather
+ // than a second one.
+ if (!lease.closed) lease._request();
+ }
+ }
+
+ async fetchChunk(fileId, chunkIndex, tr = '') {
const msg = await this._sendAndWait({
type: 'file_req',
v: '0.1',
file_id: fileId,
chunk_index: chunkIndex,
+ // Present only when this download holds a slot. The node does not require
+ // it yet; carrying it is what lets the node see the transfer is alive and
+ // not reclaim its slot as idle.
+ ...(tr ? { tr } : {}),
});
if (msg.type === 'error') throw new Error(msg.detail);
return msg;
@@ -2189,7 +2392,8 @@ class MeshBayTransport {
* 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 } = {}) {
+ 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 name — and would race for the same destination. The guard
// is by name for that reason, even though the map below is keyed by id.
@@ -2219,8 +2423,23 @@ class MeshBayTransport {
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;
@@ -2233,19 +2452,75 @@ class MeshBayTransport {
.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;
})
- .finally(wake);
+ // 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 {
- for (let i = 0; i < total; i++) {
+ // 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) {
@@ -2273,6 +2548,7 @@ class MeshBayTransport {
upload_id: uploadId,
chunk_index: i,
total_chunks: total,
+ ...(tr ? { tr } : {}),
...sealed,
});
}
@@ -2992,6 +3268,24 @@ class MeshBayTransport {
// "arrived" — and every message after that is one slot off too. Found
// live: a group mid-scan corrupted its own handshake and chat history
// this way, arriving roughly every 2s for as long as scanning ran.
+ // Routed by `tr`, and only by `tr`. A grant arrives unsolicited, minutes
+ // after the request that produced it, so falling through to "the oldest
+ // pending request" would hand a chat send or a handshake somebody else's
+ // slot — the class of defect `req_id` was introduced for.
+ if (msg.type === 'transfer_state') {
+ const lease = this._leases.get(msg.tr);
+ if (lease) lease._apply(msg);
+ else if (msg.state === 'granted') {
+ // A grant for a transfer this page has forgotten (a reload, a cancel
+ // that raced the grant). Handing it back at once matters: otherwise the
+ // node holds it until the 30 s acceptance deadline, and everyone behind
+ // it waits for nothing.
+ this._send({ type: 'transfer_close', v: '0.1', tr: msg.tr,
+ reason: 'cancelled' });
+ }
+ return;
+ }
+
if (msg.type === 'index_progress') {
if (this._onIndexProgress) {
this._onIndexProgress({