codekingpro/portable-devtools
114k
1// A readable tar stream creator
2// Technically, this is a transform stream that you write paths into,
3// and tar format comes out of.
4// The `add()` method is like `write()` but returns this,
5// and end() return `this` as well, so you can
6// do `new Pack(opt).add('files').add('dir').end().pipe(output)
7// You could also do something like:
8// streamOfPaths().pipe(new Pack()).pipe(new fs.WriteStream('out.tar'))
9import fs from 'fs';
10import { WriteEntry, WriteEntrySync, WriteEntryTar, } from './write-entry.js';
11export class PackJob {
12 path;
13 absolute;
14 entry;
15 stat;
16 readdir;
17 pending = false;
18 ignore = false;
19 piped = false;
20 constructor(path, absolute) {
21 this.path = path || './';
22 this.absolute = absolute;
23 }
24}
25import { Minipass } from 'minipass';
26import * as zlib from 'minizlib';
27import { Yallist } from 'yallist';
28import { ReadEntry } from './read-entry.js';
29import { warnMethod, } from './warn-method.js';
30const EOF = Buffer.alloc(1024);
31const ONSTAT = Symbol('onStat');
32const ENDED = Symbol('ended');
33const QUEUE = Symbol('queue');
34const CURRENT = Symbol('current');
35const PROCESS = Symbol('process');
36const PROCESSING = Symbol('processing');
37const PROCESSJOB = Symbol('processJob');
38const JOBS = Symbol('jobs');
39const JOBDONE = Symbol('jobDone');
40const ADDFSENTRY = Symbol('addFSEntry');
41const ADDTARENTRY = Symbol('addTarEntry');
42const STAT = Symbol('stat');
43const READDIR = Symbol('readdir');
44const ONREADDIR = Symbol('onreaddir');
45const PIPE = Symbol('pipe');
46const ENTRY = Symbol('entry');
47const ENTRYOPT = Symbol('entryOpt');
48const WRITEENTRYCLASS = Symbol('writeEntryClass');
49const WRITE = Symbol('write');
50const ONDRAIN = Symbol('ondrain');
51import path from 'path';
52import { normalizeWindowsPath } from './normalize-windows-path.js';
53export class Pack extends Minipass {
54 sync = false;
55 opt;
56 cwd;
57 maxReadSize;
58 preservePaths;
59 strict;
60 noPax;
61 prefix;
62 linkCache;
63 statCache;
64 file;
65 portable;
66 zip;
67 readdirCache;
68 noDirRecurse;
69 follow;
70 noMtime;
71 mtime;
72 filter;
73 jobs;
74 [WRITEENTRYCLASS];
75 onWriteEntry;
76 // Note: we actually DO need a linked list here, because we
77 // shift() to update the head of the list where we start, but still
78 // while that happens, need to know what the next item in the queue
79 // will be. Since we do multiple jobs in parallel, it's not as simple
80 // as just an Array.shift(), since that would lose the information about
81 // the next job in the list. We could add a .next field on the PackJob
82 // class, but then we'd have to be tracking the tail of the queue the
83 // whole time, and Yallist just does that for us anyway.
84 [QUEUE];
85 [JOBS] = 0;
86 [PROCESSING] = false;
87 [ENDED] = false;
88 constructor(opt = {}) {
89 //@ts-ignore
90 super();
91 this.opt = opt;
92 this.file = opt.file || '';
93 this.cwd = opt.cwd || process.cwd();
94 this.maxReadSize = opt.maxReadSize;
95 this.preservePaths = !!opt.preservePaths;
96 this.strict = !!opt.strict;
97 this.noPax = !!opt.noPax;
98 this.prefix = normalizeWindowsPath(opt.prefix || '');
99 this.linkCache = opt.linkCache || new Map();
100 this.statCache = opt.statCache || new Map();
101 this.readdirCache = opt.readdirCache || new Map();
102 this.onWriteEntry = opt.onWriteEntry;
103 this[WRITEENTRYCLASS] = WriteEntry;
104 if (typeof opt.onwarn === 'function') {
105 this.on('warn', opt.onwarn);
106 }
107 this.portable = !!opt.portable;
108 if (opt.gzip || opt.brotli || opt.zstd) {
109 if ((opt.gzip ? 1 : 0) +
110 (opt.brotli ? 1 : 0) +
111 (opt.zstd ? 1 : 0) >
112 1) {
113 throw new TypeError('gzip, brotli, zstd are mutually exclusive');
114 }
115 if (opt.gzip) {
116 if (typeof opt.gzip !== 'object') {
117 opt.gzip = {};
118 }
119 if (this.portable) {
120 opt.gzip.portable = true;
121 }
122 this.zip = new zlib.Gzip(opt.gzip);
123 }
124 if (opt.brotli) {
125 if (typeof opt.brotli !== 'object') {
126 opt.brotli = {};
127 }
128 this.zip = new zlib.BrotliCompress(opt.brotli);
129 }
130 if (opt.zstd) {
131 if (typeof opt.zstd !== 'object') {
132 opt.zstd = {};
133 }
134 this.zip = new zlib.ZstdCompress(opt.zstd);
135 }
136 /* c8 ignore next */
137 if (!this.zip)
138 throw new Error('impossible');
139 const zip = this.zip;
140 zip.on('data', chunk => super.write(chunk));
141 zip.on('end', () => super.end());
142 zip.on('drain', () => this[ONDRAIN]());
143 this.on('resume', () => zip.resume());
144 }
145 else {
146 this.on('drain', this[ONDRAIN]);
147 }
148 this.noDirRecurse = !!opt.noDirRecurse;
149 this.follow = !!opt.follow;
150 this.noMtime = !!opt.noMtime;
151 if (opt.mtime)
152 this.mtime = opt.mtime;
153 this.filter =
154 typeof opt.filter === 'function' ? opt.filter : () => true;
155 this[QUEUE] = new Yallist();
156 this[JOBS] = 0;
157 this.jobs = Number(opt.jobs) || 4;
158 this[PROCESSING] = false;
159 this[ENDED] = false;
160 }
161 [WRITE](chunk) {
162 return super.write(chunk);
163 }
164 add(path) {
165 this.write(path);
166 return this;
167 }
168 end(path, encoding, cb) {
169 /* c8 ignore start */
170 if (typeof path === 'function') {
171 cb = path;
172 path = undefined;
173 }
174 if (typeof encoding === 'function') {
175 cb = encoding;
176 encoding = undefined;
177 }
178 /* c8 ignore stop */
179 if (path) {
180 this.add(path);
181 }
182 this[ENDED] = true;
183 this[PROCESS]();
184 /* c8 ignore next */
185 if (cb)
186 cb();
187 return this;
188 }
189 write(path) {
190 if (this[ENDED]) {
191 throw new Error('write after end');
192 }
193 if (path instanceof ReadEntry) {
194 this[ADDTARENTRY](path);
195 }
196 else {
197 this[ADDFSENTRY](path);
198 }
199 return this.flowing;
200 }
201 [ADDTARENTRY](p) {
202 const absolute = normalizeWindowsPath(path.resolve(this.cwd, p.path));
203 // in this case, we don't have to wait for the stat
204 if (!this.filter(p.path, p)) {
205 p.resume();
206 }
207 else {
208 const job = new PackJob(p.path, absolute);
209 job.entry = new WriteEntryTar(p, this[ENTRYOPT](job));
210 job.entry.on('end', () => this[JOBDONE](job));
211 this[JOBS] += 1;
212 this[QUEUE].push(job);
213 }
214 this[PROCESS]();
215 }
216 [ADDFSENTRY](p) {
217 const absolute = normalizeWindowsPath(path.resolve(this.cwd, p));
218 this[QUEUE].push(new PackJob(p, absolute));
219 this[PROCESS]();
220 }
221 [STAT](job) {
222 job.pending = true;
223 this[JOBS] += 1;
224 const stat = this.follow ? 'stat' : 'lstat';
225 fs[stat](job.absolute, (er, stat) => {
226 job.pending = false;
227 this[JOBS] -= 1;
228 if (er) {
229 this.emit('error', er);
230 }
231 else {
232 this[ONSTAT](job, stat);
233 }
234 });
235 }
236 [ONSTAT](job, stat) {
237 this.statCache.set(job.absolute, stat);
238 job.stat = stat;
239 // now we have the stat, we can filter it.
240 if (!this.filter(job.path, stat)) {
241 job.ignore = true;
242 }
243 else if (stat.isFile() &&
244 stat.nlink > 1 &&
245 job === this[CURRENT] &&
246 !this.linkCache.get(`${stat.dev}:${stat.ino}`) &&
247 !this.sync) {
248 // if it's not filtered, and it's a new File entry,
249 // jump the queue in case any pending Link entries are about
250 // to try to link to it. This prevents a hardlink from coming ahead
251 // of its target in the archive.
252 this[PROCESSJOB](job);
253 }
254 this[PROCESS]();
255 }
256 [READDIR](job) {
257 job.pending = true;
258 this[JOBS] += 1;
259 fs.readdir(job.absolute, (er, entries) => {
260 job.pending = false;
261 this[JOBS] -= 1;
262 if (er) {
263 return this.emit('error', er);
264 }
265 this[ONREADDIR](job, entries);
266 });
267 }
268 [ONREADDIR](job, entries) {
269 this.readdirCache.set(job.absolute, entries);
270 job.readdir = entries;
271 this[PROCESS]();
272 }
273 [PROCESS]() {
274 if (this[PROCESSING]) {
275 return;
276 }
277 this[PROCESSING] = true;
278 for (let w = this[QUEUE].head; !!w && this[JOBS] < this.jobs; w = w.next) {
279 this[PROCESSJOB](w.value);
280 if (w.value.ignore) {
281 const p = w.next;
282 this[QUEUE].removeNode(w);
283 w.next = p;
284 }
285 }
286 this[PROCESSING] = false;
287 if (this[ENDED] && !this[QUEUE].length && this[JOBS] === 0) {
288 if (this.zip) {
289 this.zip.end(EOF);
290 }
291 else {
292 super.write(EOF);
293 super.end();
294 }
295 }
296 }
297 get [CURRENT]() {
298 return this[QUEUE] && this[QUEUE].head && this[QUEUE].head.value;
299 }
300 [JOBDONE](_job) {
301 this[QUEUE].shift();
302 this[JOBS] -= 1;
303 this[PROCESS]();
304 }
305 [PROCESSJOB](job) {
306 if (job.pending) {
307 return;
308 }
309 if (job.entry) {
310 if (job === this[CURRENT] && !job.piped) {
311 this[PIPE](job);
312 }
313 return;
314 }
315 if (!job.stat) {
316 const sc = this.statCache.get(job.absolute);
317 if (sc) {
318 this[ONSTAT](job, sc);
319 }
320 else {
321 this[STAT](job);
322 }
323 }
324 if (!job.stat) {
325 return;
326 }
327 // filtered out!
328 if (job.ignore) {
329 return;
330 }
331 if (!this.noDirRecurse &&
332 job.stat.isDirectory() &&
333 !job.readdir) {
334 const rc = this.readdirCache.get(job.absolute);
335 if (rc) {
336 this[ONREADDIR](job, rc);
337 }
338 else {
339 this[READDIR](job);
340 }
341 if (!job.readdir) {
342 return;
343 }
344 }
345 // we know it doesn't have an entry, because that got checked above
346 job.entry = this[ENTRY](job);
347 if (!job.entry) {
348 job.ignore = true;
349 return;
350 }
351 if (job === this[CURRENT] && !job.piped) {
352 this[PIPE](job);
353 }
354 }
355 [ENTRYOPT](job) {
356 return {
357 onwarn: (code, msg, data) => this.warn(code, msg, data),
358 noPax: this.noPax,
359 cwd: this.cwd,
360 absolute: job.absolute,
361 preservePaths: this.preservePaths,
362 maxReadSize: this.maxReadSize,
363 strict: this.strict,
364 portable: this.portable,
365 linkCache: this.linkCache,
366 statCache: this.statCache,
367 noMtime: this.noMtime,
368 mtime: this.mtime,
369 prefix: this.prefix,
370 onWriteEntry: this.onWriteEntry,
371 };
372 }
373 [ENTRY](job) {
374 this[JOBS] += 1;
375 try {
376 const e = new this[WRITEENTRYCLASS](job.path, this[ENTRYOPT](job));
377 return e
378 .on('end', () => this[JOBDONE](job))
379 .on('error', er => this.emit('error', er));
380 }
381 catch (er) {
382 this.emit('error', er);
383 }
384 }
385 [ONDRAIN]() {
386 if (this[CURRENT] && this[CURRENT].entry) {
387 this[CURRENT].entry.resume();
388 }
389 }
390 // like .pipe() but using super, because our write() is special
391 [PIPE](job) {
392 job.piped = true;
393 if (job.readdir) {
394 job.readdir.forEach(entry => {
395 const p = job.path;
396 const base = p === './' ? '' : p.replace(/\/*$/, '/');
397 this[ADDFSENTRY](base + entry);
398 });
399 }
400 const source = job.entry;
401 const zip = this.zip;
402 /* c8 ignore start */
403 if (!source)
404 throw new Error('cannot pipe without source');
405 /* c8 ignore stop */
406 if (zip) {
407 source.on('data', chunk => {
408 if (!zip.write(chunk)) {
409 source.pause();
410 }
411 });
412 }
413 else {
414 source.on('data', chunk => {
415 if (!super.write(chunk)) {
416 source.pause();
417 }
418 });
419 }
420 }
421 pause() {
422 if (this.zip) {
423 this.zip.pause();
424 }
425 return super.pause();
426 }
427 warn(code, message, data = {}) {
428 warnMethod(this, code, message, data);
429 }
430}
431export class PackSync extends Pack {
432 sync = true;
433 constructor(opt) {
434 super(opt);
435 this[WRITEENTRYCLASS] = WriteEntrySync;
436 }
437 // pause/resume are no-ops in sync streams.
438 pause() { }
439 resume() { }
440 [STAT](job) {
441 const stat = this.follow ? 'statSync' : 'lstatSync';
442 this[ONSTAT](job, fs[stat](job.absolute));
443 }
444 [READDIR](job) {
445 this[ONREADDIR](job, fs.readdirSync(job.absolute));
446 }
447 // gotta get it all in this tick
448 [PIPE](job) {
449 const source = job.entry;
450 const zip = this.zip;
451 if (job.readdir) {
452 job.readdir.forEach(entry => {
453 const p = job.path;
454 const base = p === './' ? '' : p.replace(/\/*$/, '/');
455 this[ADDFSENTRY](base + entry);
456 });
457 }
458 /* c8 ignore start */
459 if (!source)
460 throw new Error('Cannot pipe without source');
461 /* c8 ignore stop */
462 if (zip) {
463 source.on('data', chunk => {
464 zip.write(chunk);
465 });
466 }
467 else {
468 source.on('data', chunk => {
469 super[WRITE](chunk);
470 });
471 }
472 }
473}
474//# sourceMappingURL=pack.js.map