Repository navigation
Expand file tree
/
Copy pathwrite-stream.js
More file actions
478 lines (422 loc) 路 14.6 KB
/
Copy pathwrite-stream.js
File metadata and controls
478 lines (422 loc) 路 14.6 KB
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
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
import fs from 'node:fs';
import nodePath from 'node:path';
import { Meteor } from 'meteor/meteor';
import { helpers } from './lib.js';
/**
* @const {FileHandleCache} fhCache - FileHandle Cache, keyed by `${fileId}:${path}`
*/
const fhCache = new Map();
/**
* @const {function} noop - No Operation function
*/
const noop = () => {};
/**
* @function fileIdentity
* @param {{dev: *, ino: *, birth?: *}} identity - Device, inode, and birth time in nanoseconds
* @summary Normalize a file identity to strings, so it can be stored in MongoDB and compared with BigInt stats.
* A birth time of `0` means the file system does not report it, so it is left out
* @returns {{dev: string, ino: string, birth?: string}}
*/
const fileIdentity = ({ dev, ino, birth }) => {
const identity = { dev: `${dev}`, ino: `${ino}` };
if (birth !== undefined && birth !== null && `${birth}` !== '0') {
identity.birth = `${birth}`;
}
return identity;
};
/**
* @function isSameFile
* @param {{dev: string, ino: string, birth?: string}|null} identity - Stored identity, see `fileIdentity`
* @param {fs.BigIntStats|null} stats - Stats of an open or existing file
* @summary Check that `stats` belong to the file with this identity.
* Linux reuses an inode number right after unlink, so a replaced file can have the same dev and ino.
* Birth time tells them apart. Records saved without it, and file systems without birth time, fall back to dev and ino
* @returns {boolean}
*/
const isSameFile = (identity, stats) => {
if (!identity || !stats || !stats.isFile() || `${stats.dev}` !== identity.dev || `${stats.ino}` !== identity.ino) {
return false;
}
return !identity.birth || `${stats.birthtimeNs}` === identity.birth;
};
/**
* @private
* @locus Server
* @class WriteStream
* @param path {string} - Path to file on FS
* @param maxLength {number} - Max amount of chunks in stream
* @param file {Object} - Upload record, must have `chunkSize` (bytes per chunk)
* @param permissions {number} - Permissions which will be set to open descriptor (octal), like: `0o611` or `0o777`. Default: 0644
* @param parentDirPermissions {number} - Permissions which will be set to parent directory (octal), like: `0o611` or `0o777`. Default: 0755
* @param [options] {Object} - Extra options
* @param [options.exclusive=false] {boolean} - Create a new file, fail with `Meteor.Error(409)` if the file already exists. When `false` the file must already exist, it is never created
* @param [options.identity] {{dev: string, ino: string, birth?: string}} - Expected identity of an existing file. When set, opening a different file fails with `Meteor.Error(409)`
* @param [options.idleTimeout=0] {number} - Close the file handle after this many ms without writes, reopen on the next write. `0` disables
* @param [options.fileId] {string} - Upload id, used as part of the file handle cache key
* @param [options.onAbort] {function} - Called after the stream is aborted
* @summary Writes chunks at their offsets into one file and tracks which chunks are written
*/
export default class WriteStream {
constructor(path, maxLength, file, permissions, parentDirPermissions, options = {}) {
this.path = helpers.isString(path) ? path.trim() : '';
if (!this.path) {
throw new Meteor.Error(400, '[FilesCollection] [new WriteStream(path)] [constructor] {path} must be a String!');
}
this.maxLength = maxLength;
this.file = file;
this.permissions = permissions;
this.parentDirPermissions = parentDirPermissions;
this.exclusive = options.exclusive === true;
this.idleTimeout = (helpers.isNumber(options.idleTimeout) && options.idleTimeout > 0) ? options.idleTimeout : 0;
this.cacheKey = `${options.fileId || file?.fileId || file?._id || ''}:${this.path}`;
this.identity = (helpers.isObject(options.identity) && options.identity.dev && options.identity.ino) ? fileIdentity(options.identity) : null;
this.onAbort = helpers.isFunction(options.onAbort) ? options.onAbort : null;
this.opening = null;
this.fh = null;
this.ended = false;
this.aborted = false;
this.chunkIds = new Set();
this.writtenChunks = 0;
this.endRetries = 0;
this.maxEndRetries = 1000;
this.idleTimer = null;
this.pendingWrites = 0;
}
/**
* @memberOf WriteStream
* @name init
* @summary Initialize WriteStream, create fileHandle, ensure directory and file is writable
* @returns {Promise<WriteStream>}
*/
async init() {
let fh;
if (this.exclusive) {
const dir = nodePath.dirname(this.path);
try {
await fs.promises.mkdir(dir, { recursive: true, mode: this.parentDirPermissions });
} catch (mkdirError) {
throw new Meteor.Error(500, '[FilesCollection] [writeStream] [init] [mkdir] ERROR: can not make/ensure directory', mkdirError?.code);
}
try {
fh = await fs.promises.open(this.path, fs.constants.O_RDWR | fs.constants.O_CREAT | fs.constants.O_EXCL, this.permissions);
} catch (fsOpenError) {
// Do not unlink here: the file belongs to someone else
await this.stop(true);
if (fsOpenError?.code === 'EEXIST') {
throw new Meteor.Error(409, '[FilesCollection] [writeStream] [init] File already exists');
}
throw new Meteor.Error(500, '[FilesCollection] [writeStream] [init] [open] Error', fsOpenError?.code);
}
const stats = await fh.stat({ bigint: true });
this.identity = fileIdentity({ dev: stats.dev, ino: stats.ino, birth: stats.birthtimeNs });
} else {
// Resume: never create the file, it must be the same file created at the start of the upload
try {
fh = await this._openExisting();
} catch (openError) {
await this.stop(true);
throw openError;
}
const stats = await fh.stat();
if (stats.size > 0) {
// Chunks are written in order, so the existing size tells how many chunks are on disk
const written = Math.min(this.maxLength, Math.ceil(stats.size / this.file.chunkSize));
for (let i = 1; i <= written; i++) {
this.chunkIds.add(i);
}
this.writtenChunks = this.chunkIds.size;
}
}
this.fh = fh;
fhCache.set(this.cacheKey, this);
this._touch();
return this;
}
/**
* @memberOf WriteStream
* @name _isSameFile
* @param {fs.BigIntStats} stats - Stats of an open or existing file
* @summary Check that `stats` belong to the file this stream created
* @returns {boolean}
*/
_isSameFile(stats) {
return isSameFile(this.identity, stats);
}
/**
* @memberOf WriteStream
* @name _openExisting
* @summary Open the existing upload file without creating it, and check its identity
* @throws {Meteor.Error} 410 if the file is gone, 409 if it is a different file
* @returns {Promise<FileHandle>}
*/
async _openExisting() {
let fh;
try {
fh = await fs.promises.open(this.path, fs.constants.O_RDWR);
} catch (openError) {
if (openError?.code === 'ENOENT') {
throw new Meteor.Error(410, '[FilesCollection] [writeStream] Upload file is gone');
}
throw new Meteor.Error(500, '[FilesCollection] [writeStream] [open] Error', openError?.code);
}
let stats;
try {
stats = await fh.stat({ bigint: true });
} catch (_statError) {
stats = null;
}
if (!this._isSameFile(stats)) {
await fh.close().catch(noop);
throw new Meteor.Error(409, '[FilesCollection] [writeStream] Upload file was replaced');
}
return fh;
}
/**
* @memberOf WriteStream
* @name _touch
* @summary Restart the idle timer, it closes the file handle after `idleTimeout` ms without writes
* @returns {void}
*/
_touch() {
if (this.idleTimer) {
clearTimeout(this.idleTimer);
this.idleTimer = null;
}
if (!this.idleTimeout || this.ended) {
return;
}
this.idleTimer = setTimeout(() => {
this.idleTimer = null;
this._closeIdle();
}, this.idleTimeout);
if (typeof this.idleTimer?.unref === 'function') {
this.idleTimer.unref();
}
}
/**
* @memberOf WriteStream
* @name _closeIdle
* @summary Close the file handle of an idle upload. The next `write()` reopens it
* @returns {Promise<void>}
*/
async _closeIdle() {
if (this.ended || this.pendingWrites > 0 || this.opening || !this.fh) {
return;
}
const fh = this.fh;
this.fh = null;
if (fhCache.get(this.cacheKey) === this) {
fhCache.delete(this.cacheKey);
}
try {
await fh.close();
} catch (closeError) {
Meteor._debug('[FilesCollection] [writeStream] [_closeIdle] [close] Error:', closeError);
}
}
/**
* @memberOf WriteStream
* @name _ensureOpen
* @summary Reopen the file handle closed by the idle timer
* @returns {Promise<void>}
*/
async _ensureOpen() {
if (this.fh) {
return;
}
if (!this.opening) {
this.opening = (async () => {
const fh = await this._openExisting();
if (this.ended) {
await fh.close().catch(noop);
return;
}
this.fh = fh;
fhCache.set(this.cacheKey, this);
})().finally(() => {
this.opening = null;
});
}
await this.opening;
if (!this.fh && !this.ended) {
throw new Meteor.Error(500, '[FilesCollection] [writeStream] Can not reopen file');
}
}
/**
* @memberOf WriteStream
* @name write
* @param {number} num - Chunk position in a stream
* @param {Buffer} chunk - Buffer (chunk binary data)
* @summary Write chunk at its offset
* @throws {Meteor.Error} 409 or 410 if the file was replaced or removed while the handle was closed
* @returns {Promise<boolean>} - True if chunk was written to a file, false if chunk wasn't written
*/
async write(num, chunk) {
if (this.aborted || this.ended) {
return false;
}
if (this.idleTimer) {
clearTimeout(this.idleTimer);
this.idleTimer = null;
}
++this.pendingWrites;
try {
await this._ensureOpen();
if (this.aborted || this.ended) {
return false;
}
const { bytesWritten } = await this.fh.write(chunk, 0, chunk.byteLength, (num - 1) * this.file.chunkSize);
if (this.aborted || this.ended) {
return false;
}
await this.fh.sync();
if (bytesWritten !== chunk.byteLength) {
return false;
}
this.chunkIds.add(num);
this.writtenChunks = this.chunkIds.size;
return true;
} catch (error) {
const isFileLost = error?.error === 409 || error?.error === 410;
if (!isFileLost) {
Meteor._debug('[FilesCollection] [writeStream] [write] [Error:]', error);
}
if (!this.ended) {
await this.abort();
}
if (isFileLost) {
// The upload file is gone or replaced, the upload can not continue
throw error;
}
} finally {
--this.pendingWrites;
this._touch();
}
return false;
}
/**
* @memberOf WriteStream
* @name isComplete
* @summary Check if every chunk from 1 to `maxLength` is written
* @returns {boolean}
*/
isComplete() {
return this.chunkIds.size >= this.maxLength;
}
/**
* @memberOf WriteStream
* @name waitForCompletion
* @summary Waits up to 25 seconds for all chunks to complete writing
* @returns {Promise<boolean>} - `true` if file was fully written, `false` if writing was aborted and file removed
*/
async waitForCompletion() {
while (!this.aborted && !this.ended && !this.isComplete() && this.endRetries < this.maxEndRetries) {
++this.endRetries;
await new Promise(resolve => setTimeout(resolve, 25));
}
if (this.aborted) {
return false;
}
if (this.isComplete()) {
return await this.stop(false);
}
await this.abort();
return false;
}
/**
* @memberOf WriteStream
* @name end
* @summary Finishes writing, only after all chunks are written
* @returns {Promise<boolean>} - `true` if every chunk is written and the file is closed, `false` if the upload was aborted or is incomplete
*/
async end() {
if (this.aborted) {
return false;
}
if (this.ended) {
return true;
}
if (this.isComplete()) {
return await this.stop(false);
}
if (await this.waitForCompletion()) {
return true;
}
Meteor._debug('[FilesCollection] [writeStream] [end] waitForCompletion waited for 25 seconds to complete writing and failed with timeout', this.path);
return false;
}
/**
* @memberOf WriteStream
* @name abort
* @summary Aborts writing and removes created file, only if the file on disk is still the file this stream created. Does nothing when the stream already finished successfully
* @returns {Promise<boolean>} - `true` if the stream is aborted, `false` if it already finished successfully
*/
async abort() {
if (this.aborted) {
return true;
}
if (this.ended) {
// Finished file must stay on disk
return false;
}
await this.stop(true);
let stats = null;
try {
stats = await fs.promises.lstat(this.path, { bigint: true });
} catch (_statError) {
stats = null;
}
if (this._isSameFile(stats)) {
try {
await fs.promises.unlink(this.path);
} catch (unlinkError) {
Meteor._debug('[FilesCollection] [writeStream] [abort] [unlink] [ERROR:]', this.path, unlinkError);
}
}
if (this.onAbort) {
try {
await this.onAbort(this);
} catch (onAbortError) {
Meteor._debug('[FilesCollection] [writeStream] [abort] [onAbort] [ERROR:]', onAbortError);
}
}
return true;
}
/**
* @memberOf WriteStream
* @name stop
* @param {boolean} [isAborted=false] - was stop called because it was aborted?
* @summary Stop writing and close the file handle
* @returns {Promise<boolean>} - true
*/
async stop(isAborted = false) {
if (this.ended) {
return true;
}
this.ended = true;
if (this.idleTimer) {
clearTimeout(this.idleTimer);
this.idleTimer = null;
}
if (isAborted) {
this.aborted = true;
} else if (this.fh) {
try {
await this.fh.datasync();
} catch (dsError) {
Meteor._debug('[FilesCollection] [writeStream] [stop] fh.datasync resulted in Error:', this.path, dsError);
}
}
try {
await this.fh?.close();
} catch (closeError) {
Meteor._debug('[FilesCollection] [writeStream] [stop] fh.close resulted in Error:', this.path, closeError);
}
this.fh = null;
if (fhCache.get(this.cacheKey) === this) {
fhCache.delete(this.cacheKey);
}
return true;
}
}
export { fileIdentity, isSameFile };