diff options
Diffstat (limited to 'packages/meshbay-hub/src')
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/static/transport.js | 99 |
1 files changed, 88 insertions, 11 deletions
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/transport.js b/packages/meshbay-hub/src/meshbay_hub/static/transport.js index 179292e..e57ad8b 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/transport.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/transport.js @@ -287,6 +287,11 @@ class MeshBayTransport { this._pc = null; this._channel = null; this._pending = new Map(); + // Set the first time this connection sees a reply that names the request + // it answers (see _dispatch). A node either stamps every reply or none, + // so one is proof for the connection — and once there is proof, the + // arrival-order fallback at the bottom of _dispatch is never right again. + this._correlates = false; this._seqId = 0; this._recvBuf = new Uint8Array(0); this._connected = false; @@ -1429,6 +1434,12 @@ class MeshBayTransport { thread_id: threadId || null, sender_name: senderName || null, }); + // Every other request in this file refuses an `error` reply; this one + // returned it as though the node had accepted the message. It never + // mattered while a refusal reached the wrong caller anyway — now that a + // reply finds the request that made it, a message the node rejected + // would otherwise appear in the conversation as sent. + if (msg.type === 'error') throw new Error(msg.detail || 'chat send refused'); return msg; } @@ -2191,12 +2202,15 @@ class MeshBayTransport { if (msg.type === 'index_sync') { if (this._onIndexSync) this._onIndexSync(opened); - for (const [, handler] of this._pending) { - if (handler._reqType === 'index_sync') { - handler.resolve(opened); - break; - } - } + // The node's first push to a newly connected peer is an index_sync + // nobody asked for, so there is not always a request to resolve. When + // there is, `req_id` says which one — the type match below is what a + // node too old to stamp one leaves us, and it is why two fetches in + // flight at once used to resolve the wrong one. + const handler = opened.req_id !== undefined && opened.req_id !== null + ? this._pending.get(opened.req_id) + : [...this._pending.values()].find(h => h._reqType === 'index_sync'); + if (handler) handler.resolve(opened); return; } if (this._onIndexDelta) this._onIndexDelta(opened); @@ -2283,7 +2297,13 @@ class MeshBayTransport { resolve: (msg) => { clearTimeout(timeout); this._pending.delete(id); resolve(msg); }, reject: (err) => { clearTimeout(timeout); this._pending.delete(id); reject(err); }, }); - this._send(obj); + // The id goes on the wire (MNP 1.1+): a node that understands it stamps + // the reply with it, and _dispatch matches on that alone. It used to be + // local to this map, which is why every reply had to be recognised by + // some field of its own — and why the ones that carry no such field + // reached their caller by luck. An older node ignores the extra key and + // is routed by the per-type fallbacks below, exactly as before. + this._send({ ...obj, req_id: id }); }); } @@ -2323,6 +2343,45 @@ class MeshBayTransport { } _dispatch(msg) { + // A reply that names the request it answers. Nothing below this needs to + // recognise it, and nothing below this may see it: every remaining branch + // exists to identify a reply by some field of its own, which is the job + // this makes unnecessary. + // + // What is left underneath is genuinely unsolicited — a broadcast to every + // connected client, a push, a challenge — or a reply from a node too old + // to stamp one, which is what the per-type keys are for now. + if (msg.req_id !== undefined && msg.req_id !== null) { + this._correlates = true; + // The one exception, and the only one: an index message is sealed under + // the GEK and cannot be handed to its caller until it is opened, which + // is not something this synchronous function can do. Resolving it here + // would give `fetchIndex` the envelope — nonce and ciphertext, no + // entries — and skip `_onIndexSync` entirely. `_queueIndexMessage` + // opens it and then resolves, by this same id. + const sealed = msg.type === 'index_sync' || msg.type === 'index_delta'; + if (!sealed) { + const handler = this._pending.get(msg.req_id); + if (handler) { + handler.resolve(msg); + // The acks whose *broadcast* half their own requester also needs: + // every other client learns the change from the broadcast, and the + // one that asked for it is the only one that would not, because its + // own request swallowed its copy. Same call the keyed `_ack` branch + // below makes, for the same reason. + if (BROADCAST_ACK_TYPES.has(msg.type)) _replayBroadcast(this, msg); + return; + } + // Answers a request that is no longer waiting: it gave up at its own + // timeout, or a reconnect rejected everything in flight. It belongs to + // nobody, and the whole point of this change is that it is not offered + // to somebody else instead. + console.warn('[MeshBay] late reply to req', msg.req_id, '(', msg.type, + ') — nothing waiting'); + return; + } + } + // Two-step admin-op flow (_authorizeAdminOp, ADMIN_OP_TYPES) — resolve // by (op) key before anything below gets a chance to steal it via the // generic "oldest pending" fallback further down. Returns as soon as a @@ -2715,10 +2774,28 @@ class MeshBayTransport { return; } - // Everything above is routed by something in the message. What is left is - // matched by arrival order, which is only ever a guess — and a wrong guess - // here hands one request's answer to another, which then waits for a reply - // that already came. Logged so that guess is visible. + // Everything above is routed by something in the message. What is left + // used to be matched by arrival order — a guess, and a wrong guess hands + // one request's answer to another, which then waits out its own 30s + // timeout for a reply that already came and went. That is how the Chat + // composer, disabled while a send is in flight, could stay disabled for + // thirty seconds on a message the node had already stored. + // + // A node that stamps its replies (`req_id`, handled at the top) has taken + // every one of its answers out of this path, so anything arriving here is + // unsolicited and the guess can only ever be wrong. Dropping it loses + // nothing and stops the theft. + if (this._correlates) { + console.warn('[MeshBay] unsolicited', msg.type, '— dropped (pending:', + this._pending.size, ')'); + return; + } + + // Only a node too old to stamp anything reaches here, where arrival order + // is still the only thing there is. Kept deliberately, and no wider than + // it was: the alternative for such a node is that half the protocol + // (device_list_result, join_result, the handshake's own replies) reaches + // nobody at all. const oldest = this._pending.entries().next(); if (!oldest.done) { const [, handler] = oldest.value; |