codekingpro/portable-devtools
115k
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 super();
90 this.opt = opt;
91 this.file = opt.file || '';
92 this.cwd = opt.cwd || process.cwd();
93 this.maxReadSize = opt.maxReadSize;
94 this.preservePaths = !!opt.preservePaths;
95 this.strict = !!opt.strict;
96 this.noPax = !!opt.noPax;
97 this.prefix = normalizeWindowsPath(opt.prefix || '');
98 this.linkCache = opt.linkCache || new Map();
99 this.statCache = opt.statCache || new Map();
100 this.readdirCache = opt.readdirCache || new Map();
101 this.onWriteEntry = opt.onWriteEntry;
102 this[WRITEENTRYCLASS] = WriteEntry;
103 if (typeof opt.onwarn === 'function') {
104 this.on('warn', opt.onwarn);
105 }
106 this.portable = !!opt.portable;
107 if (opt.gzip || opt.brotli || opt.zstd) {
108 if ((opt.gzip ? 1 : 0) + (opt.brotli ? 1 : 0) + (opt.zstd ? 1 : 0) >
109 1) {
110 throw new TypeError('gzip, brotli, zstd are mutually exclusive');
111 }
112 if (opt.gzip) {
113 if (typeof opt.gzip !== 'object') {
114 opt.gzip = {};
115 }
116 if (this.portable) {
117 opt.gzip.portable = true;
118 }
119 this.zip = new zlib.Gzip(opt.gzip);
120 }
121 if (opt.brotli) {
122 if (typeof opt.brotli !== 'object') {
123 opt.brotli = {};
124 }
125 this.zip = new zlib.BrotliCompress(opt.brotli);
126 }
127 if (opt.zstd) {
128 if (typeof opt.zstd !== 'object') {
129 opt.zstd = {};
130 }
131 this.zip = new zlib.ZstdCompress(opt.zstd);
132 }
133 /* c8 ignore next */
134 if (!this.zip)
135 throw new Error('impossible');
136 const zip = this.zip;
137 zip.on('data', chunk => super.write(chunk));
138 zip.on('end', () => super.end());
139 zip.on('drain', () => this[ONDRAIN]());
140 this.on('resume', () => zip.resume());
141 }
142 else {
143 this.on('drain', this[ONDRAIN]);
144 }
145 this.noDirRecurse = !!opt.noDirRecurse;
146 this.follow = !!opt.follow;
147 this.noMtime = !!opt.noMtime;
148 if (opt.mtime)
149 this.mtime = opt.mtime;
150 this.filter =
151 typeof opt.filter === 'function' ? opt.filter : () => true;
152 this[QUEUE] = new Yallist();
153 this[JOBS] = 0;
154 this.jobs = Number(opt.jobs) || 4;
155 this[PROCESSING] = false;
156 this[ENDED] = false;
157 }
158 [WRITE](chunk) {
159 return super.write(chunk);
160 }
161 add(path) {
162 this.write(path);
163 return this;
164 }
165 end(path, encoding, cb) {
166 /* c8 ignore start */
167 if (typeof path === 'function') {
168 cb = path;
169 path = undefined;
170 }
171 if (typeof encoding === 'function') {
172 cb = encoding;
173 encoding = undefined;
174 }
175 /* c8 ignore stop */
176 if (path) {
177 this.add(path);
178 }
179 this[ENDED] = true;
180 this[PROCESS]();
181 /* c8 ignore next */
182 if (cb)
183 cb();
184 return this;
185 }
186 write(path) {
187 if (this[ENDED]) {
188 throw new Error('write after end');
189 }
190 if (path instanceof ReadEntry) {
191 this[ADDTARENTRY](path);
192 }
193 else {
194 this[ADDFSENTRY](path);
195 }
196 return this.flowing;
197 }
198 [ADDTARENTRY](p) {
199 const absolute = normalizeWindowsPath(path.resolve(this.cwd, p.path));
200 // in this case, we don't have to wait for the stat
201 if (!this.filter(p.path, p)) {
202 p.resume();
203 }
204 else {
205 const job = new PackJob(p.path, absolute);
206 job.entry = new WriteEntryTar(p, this[ENTRYOPT](job));
207 job.entry.on('end', () => this[JOBDONE](job));
208 this[JOBS] += 1;
209 this[QUEUE].push(job);
210 }
211 this[PROCESS]();
212 }
213 [ADDFSENTRY](p) {
214 const absolute = normalizeWindowsPath(path.resolve(this.cwd, p));
215 this[QUEUE].push(new PackJob(p, absolute));
216 this[PROCESS]();
217 }
218 [STAT](job) {
219 job.pending = true;
220 this[JOBS] += 1;
221 const stat = this.follow ? 'stat' : 'lstat';
222 fs[stat](job.absolute, (er, stat) => {
223 job.pending = false;
224 this[JOBS] -= 1;
225 if (er) {
226 this.emit('error', er);
227 }
228 else {
229 this[ONSTAT](job, stat);
230 }
231 });
232 }
233 [ONSTAT](job, stat) {
234 this.statCache.set(job.absolute, stat);
235 job.stat = stat;
236 // now we have the stat, we can filter it.
237 if (!this.filter(job.path, stat)) {
238 job.ignore = true;
239 }
240 else if (stat.isFile() &&
241 stat.nlink > 1 &&
242 job === this[CURRENT] &&
243 !this.linkCache.get(`${stat.dev}:${stat.ino}`) &&
244 !this.sync) {
245 // if it's not filtered, and it's a new File entry,
246 // jump the queue in case any pending Link entries are about
247 // to try to link to it. This prevents a hardlink from coming ahead
248 // of its target in the archive.
249 this[PROCESSJOB](job);
250 }
251 this[PROCESS]();
252 }
253 [READDIR](job) {
254 job.pending = true;
255 this[JOBS] += 1;
256 fs.readdir(job.absolute, (er, entries) => {
257 job.pending = false;
258 this[JOBS] -= 1;
259 if (er) {
260 return this.emit('error', er);
261 }
262 this[ONREADDIR](job, entries);
263 });
264 }
265 [ONREADDIR](job, entries) {
266 this.readdirCache.set(job.absolute, entries);
267 job.readdir = entries;
268 this[PROCESS]();
269 }
270 [PROCESS]() {
271 if (this[PROCESSING]) {
272 return;
273 }
274 this[PROCESSING] = true;
275 for (let w = this[QUEUE].head; !!w && this[JOBS] < this.jobs; w = w.next) {
276 this[PROCESSJOB](w.value);
277 if (w.value.ignore) {
278 const p = w.next;
279 this[QUEUE].removeNode(w);
280 w.next = p;
281 }
282 }
283 this[PROCESSING] = false;
284 if (this[ENDED] && this[QUEUE].length === 0 && this[JOBS] === 0) {
285 if (this.zip) {
286 this.zip.end(EOF);
287 }
288 else {
289 super.write(EOF);
290 super.end();
291 }
292 }
293 }
294 get [CURRENT]() {
295 return this[QUEUE] && this[QUEUE].head && this[QUEUE].head.value;
296 }
297 [JOBDONE](_job) {
298 this[QUEUE].shift();
299 this[JOBS] -= 1;
300 this[PROCESS]();
301 }
302 [PROCESSJOB](job) {
303 if (job.pending) {
304 return;
305 }
306 if (job.entry) {
307 if (job === this[CURRENT] && !job.piped) {
308 this[PIPE](job);
309 }
310 return;
311 }
312 if (!job.stat) {
313 const sc = this.statCache.get(job.absolute);
314 if (sc) {
315 this[ONSTAT](job, sc);
316 }
317 else {
318 this[STAT](job);
319 }
320 }
321 if (!job.stat) {
322 return;
323 }
324 // filtered out!
325 if (job.ignore) {
326 return;
327 }
328 if (!this.noDirRecurse && job.stat.isDirectory() && !job.readdir) {
329 const rc = this.readdirCache.get(job.absolute);
330 if (rc) {
331 this[ONREADDIR](job, rc);
332 }
333 else {
334 this[READDIR](job);
335 }
336 if (!job.readdir) {
337 return;
338 }
339 }
340 // we know it doesn't have an entry, because that got checked above
341 job.entry = this[ENTRY](job);
342 if (!job.entry) {
343 job.ignore = true;
344 return;
345 }
346 if (job === this[CURRENT] && !job.piped) {
347 this[PIPE](job);
348 }
349 }
350 [ENTRYOPT](job) {
351 return {
352 onwarn: (code, msg, data) => this.warn(code, msg, data),
353 noPax: this.noPax,
354 cwd: this.cwd,
355 absolute: job.absolute,
356 preservePaths: this.preservePaths,
357 maxReadSize: this.maxReadSize,
358 strict: this.strict,
359 portable: this.portable,
360 linkCache: this.linkCache,
361 statCache: this.statCache,
362 noMtime: this.noMtime,
363 mtime: this.mtime,
364 prefix: this.prefix,
365 onWriteEntry: this.onWriteEntry,
366 };
367 }
368 [ENTRY](job) {
369 this[JOBS] += 1;
370 try {
371 const e = new this[WRITEENTRYCLASS](job.path, this[ENTRYOPT](job));
372 return e
373 .on('end', () => this[JOBDONE](job))
374 .on('error', er => this.emit('error', er));
375 }
376 catch (er) {
377 this.emit('error', er);
378 }
379 }
380 [ONDRAIN]() {
381 if (this[CURRENT] && this[CURRENT].entry) {
382 this[CURRENT].entry.resume();
383 }
384 }
385 // like .pipe() but using super, because our write() is special
386 [PIPE](job) {
387 job.piped = true;
388 if (job.readdir) {
389 job.readdir.forEach(entry => {
390 const p = job.path;
391 const base = p === './' ? '' : p.replace(/\/*$/, '/');
392 this[ADDFSENTRY](base + entry);
393 });
394 }
395 const source = job.entry;
396 const zip = this.zip;
397 /* c8 ignore start */
398 if (!source)
399 throw new Error('cannot pipe without source');
400 /* c8 ignore stop */
401 if (zip) {
402 source.on('data', chunk => {
403 if (!zip.write(chunk)) {
404 source.pause();
405 }
406 });
407 }
408 else {
409 source.on('data', chunk => {
410 if (!super.write(chunk)) {
411 source.pause();
412 }
413 });
414 }
415 }
416 pause() {
417 if (this.zip) {
418 this.zip.pause();
419 }
420 return super.pause();
421 }
422 warn(code, message, data = {}) {
423 warnMethod(this, code, message, data);
424 }
425}
426export class PackSync extends Pack {
427 sync = true;
428 constructor(opt) {
429 super(opt);
430 this[WRITEENTRYCLASS] = WriteEntrySync;
431 }
432 // pause/resume are no-ops in sync streams.
433 pause() { }
434 resume() { }
435 [STAT](job) {
436 const stat = this.follow ? 'statSync' : 'lstatSync';
437 this[ONSTAT](job, fs[stat](job.absolute));
438 }
439 [READDIR](job) {
440 this[ONREADDIR](job, fs.readdirSync(job.absolute));
441 }
442 // gotta get it all in this tick
443 [PIPE](job) {
444 const source = job.entry;
445 const zip = this.zip;
446 if (job.readdir) {
447 job.readdir.forEach(entry => {
448 const p = job.path;
449 const base = p === './' ? '' : p.replace(/\/*$/, '/');
450 this[ADDFSENTRY](base + entry);
451 });
452 }
453 /* c8 ignore start */
454 if (!source)
455 throw new Error('Cannot pipe without source');
456 /* c8 ignore stop */
457 if (zip) {
458 source.on('data', chunk => {
459 zip.write(chunk);
460 });
461 }
462 else {
463 source.on('data', chunk => {
464 super[WRITE](chunk);
465 });
466 }
467 }
468}
469//# sourceMappingURL=pack.js.map