Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes15kdownloads
pack.js469 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        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
codekingpro/portable-devtools · Team Ai