Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
pack.js474 linesDownload Raw Back to esm
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
codekingpro/portable-devtools · Team Ai