1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
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;
|