diff --git a/README.md b/README.md index be5927d..9437b7d 100644 --- a/README.md +++ b/README.md @@ -49,6 +49,7 @@ Positionals: Options: -u, --up slice a path off the bottom of the paths [number] -a, --all include files & directories begining with a dot (.) [boolean] + -c, --concurrency maximum simultaneous file copies [number] -d, --dry-run show what would be copied, without actually copying anything [boolean] -f, --flat flatten the output [boolean] -e, --exclude pattern or glob to exclude (may be passed multiple times) [string|string[]] @@ -242,9 +243,12 @@ copyfiles(source[s], destination, options, callback); 3. third argument is the optional "options" argument 4. and finally the last argument is an optional callback function that will be executed after the copy process ended +By default, each call runs up to `os.availableParallelism()` file copies at once, capped at 32. Set `concurrency` in the API or `--concurrency` / `-c` in the CLI to override this limit with a positive integer. If a copy fails, queued copies stop and active streams close before the callback receives the error. Files already written may remain in the destination. + ```js { verbose: boolean; // print more information to console + concurrency: number; // maximum simultaneous copies (default: available parallelism, capped at 32) up: number; // slice a path off the bottom of the paths exclude: string; // exclude pattern all: boolean; // include dot files diff --git a/src/__tests__/performance.spec.ts b/src/__tests__/performance.spec.ts new file mode 100644 index 0000000..ceb13e0 --- /dev/null +++ b/src/__tests__/performance.spec.ts @@ -0,0 +1,315 @@ +import { + createReadStream, + createWriteStream, + existsSync, + globSync, + mkdirSync, + mkdtempSync, + readdirSync, + readFileSync, + rmSync, + statSync, + symlinkSync, + writeFileSync, +} from 'node:fs'; +import { availableParallelism } from 'node:os'; +import { join, relative } from 'node:path'; +import { Readable, Writable } from 'node:stream'; +import { afterEach, beforeEach, describe, expect, test, vi } from 'vitest'; + +import { copyfiles } from '../index.js'; +import type { CopyFileOptions } from '../interfaces.js'; + +vi.mock('node:fs', async importOriginal => { + const actual = await importOriginal(); + return { + ...actual, + createReadStream: vi.fn(actual.createReadStream), + createWriteStream: vi.fn(actual.createWriteStream), + existsSync: vi.fn(actual.existsSync), + globSync: vi.fn(actual.globSync), + statSync: vi.fn(actual.statSync), + }; +}); + +vi.mock('node:os', async importOriginal => ({ + ...(await importOriginal()), + availableParallelism: vi.fn(), +})); + +describe('copy resource management', () => { + let root: string; + let input: string; + let output: string; + + beforeEach(async () => { + vi.mocked(availableParallelism).mockReturnValue(12); + const actual = await vi.importActual('node:fs'); + vi.mocked(createReadStream).mockImplementation(actual.createReadStream); + vi.mocked(createWriteStream).mockImplementation(actual.createWriteStream); + mkdirSync('tmp', { recursive: true }); + root = mkdtempSync(join(process.cwd(), 'tmp', 'native-copyfiles-test-')); + input = join(root, 'input'); + output = join(root, 'output'); + mkdirSync(input); + }); + + afterEach(() => { + vi.restoreAllMocks(); + rmSync(root, { recursive: true, force: true }); + }); + + const copy = (sources: string | string[], destination: string, options: CopyFileOptions = {}) => + new Promise((resolve, reject) => { + copyfiles(sources, destination, options, err => (err ? reject(err) : resolve())); + }); + + function assertStreamsClosed() { + for (const stream of [ + ...vi.mocked(createReadStream).mock.results.map(result => result.value), + ...vi.mocked(createWriteStream).mock.results.map(result => result.value), + ]) { + expect(stream.destroyed).toBe(true); + expect(stream.closed).toBe(true); + } + } + + async function assertConcurrency(limit: number, options: CopyFileOptions = {}) { + for (let i = 0; i < 160; i++) { + writeFileSync(join(input, `${i}.txt`), `contents-${i}`); + } + const actual = await vi.importActual('node:fs'); + let active = 0; + let peak = 0; + vi.mocked(createReadStream).mockImplementation((...args) => { + const stream = actual.createReadStream(...args); + active++; + peak = Math.max(peak, active); + stream.once('close', () => active--); + return stream; + }); + + options.flat = true; + await copy(`${input}/*.txt`, output, options); + + expect(peak).toBe(limit); + expect(active).toBe(0); + expect(readdirSync(output)).toHaveLength(160); + for (let i = 0; i < 160; i++) { + expect(readFileSync(join(output, `${i}.txt`), 'utf8')).toBe(`contents-${i}`); + } + assertStreamsClosed(); + } + + test.each([1, 12, 64])('uses available parallelism %i with a maximum default of 32', async parallelism => { + vi.mocked(availableParallelism).mockReturnValue(parallelism); + await assertConcurrency(Math.min(32, parallelism)); + }); + + test.each([1, 3, 48])('honors an explicit concurrency of %i', async concurrency => { + vi.mocked(availableParallelism).mockReturnValue(1); + await assertConcurrency(concurrency, { concurrency }); + expect(availableParallelism).not.toHaveBeenCalled(); + }); + + test.each([0, -1, 1.5, Number.NaN, Number.POSITIVE_INFINITY, Number.MAX_SAFE_INTEGER + 1])( + 'rejects invalid concurrency %s before accessing files', + async concurrency => { + const options = { concurrency }; + expect(() => copyfiles(`${input}/*.txt`, output, options)).toThrow('Concurrency must be a positive safe integer.'); + await expect(copy(`${input}/*.txt`, output, options)).rejects.toThrow('Concurrency must be a positive safe integer.'); + expect(globSync).not.toHaveBeenCalled(); + expect(createReadStream).not.toHaveBeenCalled(); + }, + ); + + test('ignores an inherited concurrency option', async () => { + vi.mocked(availableParallelism).mockReturnValue(3); + await assertConcurrency(3, Object.create({ concurrency: 1 })); + }); + + test('closes both streams after repeated write failures', async () => { + writeFileSync(join(input, 'large.txt'), Buffer.alloc(1024 * 1024)); + mkdirSync(join(output, 'large.txt'), { recursive: true }); + + for (let i = 0; i < 10; i++) { + await expect(copy(join(input, 'large.txt'), output, { flat: true })).rejects.toBeInstanceOf(Error); + assertStreamsClosed(); + } + }); + + test('stops queued copies, aborts active streams and reports the original read error once', async () => { + for (let i = 0; i < 96; i++) { + writeFileSync(join(input, `${i}.txt`), Buffer.alloc(1024 * 1024)); + } + const error = new Error('read failed'); + vi.mocked(createReadStream).mockImplementationOnce(() => { + const stream = new Readable({ read() {} }); + setImmediate(() => stream.destroy(error)); + return stream as ReturnType; + }); + const rename = vi.fn((_src: string, dest: string) => dest); + const callback = vi.fn(); + await new Promise(resolve => { + copyfiles(`${input}/*.txt`, output, { flat: true, rename }, err => { + callback(err); + resolve(); + }); + }); + await new Promise(resolve => setImmediate(resolve)); + + expect(callback).toHaveBeenCalledExactlyOnceWith(error); + expect(rename.mock.calls.length).toBeLessThanOrEqual(12); + assertStreamsClosed(); + }); + + test.each(['read', 'write'])('reports premature %s closure and closes the other stream', async side => { + writeFileSync(join(input, 'file.txt'), Buffer.alloc(1024 * 1024)); + if (side === 'read') { + vi.mocked(createReadStream).mockImplementationOnce( + () => + new Readable({ + read() { + this.destroy(); + }, + }) as ReturnType, + ); + } else { + vi.mocked(createWriteStream).mockImplementationOnce( + () => + new Writable({ + write(_chunk, _encoding, callback) { + this.destroy(); + callback(); + }, + }) as ReturnType, + ); + } + + await expect(copy(join(input, 'file.txt'), output, { flat: true })).rejects.toThrow('closed before copying'); + + assertStreamsClosed(); + }); + + test('does not start queued copies after a synchronous rename error', async () => { + for (let i = 0; i < 64; i++) { + writeFileSync(join(input, `${i}.txt`), 'contents'); + } + const error = new Error('rename failed'); + const rename = vi.fn(() => { + throw error; + }); + + await expect(copy(`${input}/*.txt`, output, { rename })).rejects.toBe(error); + + expect(rename).toHaveBeenCalledTimes(1); + expect(createReadStream).not.toHaveBeenCalled(); + }); + + test('waits for active streams when a later worker encounters a synchronous error', async () => { + const sources = Array.from({ length: 64 }, (_, i) => join(input, `${i}.txt`)); + for (const source of sources) { + writeFileSync(source, Buffer.alloc(8192)); + } + const error = new Error('second rename failed'); + let renamed = 0; + const rename = vi.fn((_src: string, dest: string) => { + if (++renamed === 2) { + throw error; + } + return dest; + }); + const callback = vi.fn(); + await new Promise(resolve => { + copyfiles(sources, output, { flat: true, concurrency: 3, rename }, err => { + callback(err); + resolve(); + }); + }); + await new Promise(resolve => setImmediate(resolve)); + + expect(callback).toHaveBeenCalledExactlyOnceWith(error); + expect(rename).toHaveBeenCalledTimes(2); + expect(createReadStream).toHaveBeenCalledTimes(1); + assertStreamsClosed(); + }); + + test('reports per-file directory creation failures through the callback', async () => { + writeFileSync(join(input, 'file.txt'), 'contents'); + mkdirSync(output); + writeFileSync(join(output, 'blocked'), 'not a directory'); + + await expect( + copy(join(input, 'file.txt'), output, { rename: () => join(output, 'blocked', 'nested', 'file.txt') }), + ).rejects.toBeInstanceOf(Error); + + expect(createReadStream).not.toHaveBeenCalled(); + }); + + test('applies string exclusions', async () => { + writeFileSync(join(input, 'keep.txt'), 'keep'); + writeFileSync(join(input, 'skip.txt'), 'skip'); + + await copy(`${input}/*.txt`, output, { flat: true, exclude: '**/skip.txt' }); + + expect(readdirSync(output)).toEqual(['keep.txt']); + }); + + test('retains default exclusions for an empty exclusion array', async () => { + mkdirSync(join(input, 'node_modules')); + writeFileSync(join(input, 'keep.txt'), 'keep'); + writeFileSync(join(input, 'node_modules', 'skip.txt'), 'skip'); + + await copy(`${input}/**/*.txt`, output, { flat: true, exclude: [] }); + + expect(readdirSync(output)).toEqual(['keep.txt']); + }); + + test('caches discovery but preserves ordered negation and re-inclusion', async () => { + writeFileSync(join(input, 'keep.txt'), 'keep'); + writeFileSync(join(input, 'skip.txt'), 'skip'); + const pattern = `${input}/*.txt`; + + await copy([pattern, `!${pattern}`, pattern, `!${input}/skip.txt`], output, { flat: true }); + + expect(readdirSync(output)).toEqual(['keep.txt']); + expect(globSync).toHaveBeenCalledTimes(2); + // Only source-pattern checks need stat calls, not each matching file. + expect(statSync).toHaveBeenCalledTimes(2); + }); + + test('checks a shared destination directory once per operation', async () => { + for (let i = 0; i < 80; i++) { + writeFileSync(join(input, `${i}.txt`), 'contents'); + } + + await copy(`${input}/*.txt`, output, { flat: true }); + await copy(`${input}/*.txt`, output, { flat: true }); + + expect(vi.mocked(existsSync).mock.calls.filter(([path]) => path === output)).toHaveLength(2); + }); + + test('retains relative paths and symlinked-file copying while filtering linked directories', async () => { + if (process.platform === 'win32') { + return; + } + writeFileSync(join(input, 'file.txt'), 'contents'); + mkdirSync(join(input, 'dir.txt')); + symlinkSync('file.txt', join(input, 'link.txt')); + symlinkSync('dir.txt', join(input, 'dir-link.txt')); + const source = relative(process.cwd(), input).replaceAll('\\', '/'); + const renamedSources: string[] = []; + + await copy(`${source}/*.txt`, output, { + flat: true, + rename: (src, dest) => { + renamedSources.push(src); + return dest; + }, + }); + + expect(readdirSync(output).sort()).toEqual(['file.txt', 'link.txt']); + expect(readFileSync(join(output, 'link.txt'), 'utf8')).toBe('contents'); + expect(renamedSources.sort()).toEqual([`${source}/file.txt`, `${source}/link.txt`].sort()); + }); +}); diff --git a/src/cli.ts b/src/cli.ts index c926163..0f4f82e 100644 --- a/src/cli.ts +++ b/src/cli.ts @@ -32,6 +32,11 @@ try { type: 'boolean', describe: 'Include files & directories begining with a dot (.)', }, + concurrency: { + alias: 'c', + type: 'number', + describe: 'Maximum simultaneous file copies (default: available parallelism, capped at 32)', + }, dryRun: { alias: 'd', type: 'boolean', diff --git a/src/index.ts b/src/index.ts index fd6248c..9a3e30b 100644 --- a/src/index.ts +++ b/src/index.ts @@ -1,4 +1,5 @@ -import { createReadStream, createWriteStream, existsSync, globSync, mkdirSync, type PathLike, statSync } from 'node:fs'; +import { createReadStream, createWriteStream, existsSync, globSync, mkdirSync, type ReadStream, statSync, type WriteStream } from 'node:fs'; +import { availableParallelism } from 'node:os'; import { basename, dirname, extname, join, normalize, posix, sep } from 'node:path'; import untildify from 'untildify'; import type { CopyFileOptions } from './interfaces.js'; @@ -121,7 +122,7 @@ export function filterDotFiles(paths: string[], dot: boolean): string[] { }); } -function tryCreatingDir(path: PathLike, defaultReturn: any) { +function tryCreatingDir(path: string, defaultReturn: T): string | T { try { if (statSync(path).isDirectory()) { return `${path}/**`; @@ -143,37 +144,41 @@ function getMatchedFiles( options: CopyFileOptions, ): Set { const allFilesSet = new Set(); + const filesByPattern = new Map(); + const isSingleFileRename = isSingleFile && isDestFile; for (const pattern of sources) { const isNegated = typeof pattern === 'string' && pattern.startsWith('!'); const dirPart = isNegated ? pattern.slice(1) : pattern; - const adjustedPattern = tryCreatingDir(dirPart, dirPart); - let files = globSync(adjustedPattern, { exclude: excludeGlobs }) || []; - if (options.all && adjustedPattern.includes('*') && !adjustedPattern.startsWith('.')) { - const dotPattern = adjustedPattern.replace(/(\*\.[^/]+$|\*$)/, '.$1'); - if (dotPattern !== adjustedPattern) { - files = files.concat(globSync(dotPattern, { exclude: excludeGlobs })); + let files = filesByPattern.get(dirPart); + if (!files) { + const adjustedPattern = tryCreatingDir(dirPart, dirPart); + let entries = globSync(adjustedPattern, { exclude: excludeGlobs, withFileTypes: true }); + if (options.all && adjustedPattern.includes('*') && !adjustedPattern.startsWith('.')) { + const dotPattern = adjustedPattern.replace(/(\*\.[^/]+$|\*$)/, '.$1'); + if (dotPattern !== adjustedPattern) { + entries = entries.concat(globSync(dotPattern, { exclude: excludeGlobs, withFileTypes: true })); + } } - } - files = arrify(files); - files = files.map(f => f.replaceAll('\\', '/')); - files = files.filter(f => !tryCreatingDir(f, false)); - - // Special case: single file rename to a file path - if (isSingleFile && isDestFile) { - for (const f of files) { - allFilesSet.add(f); + files = []; + for (const entry of entries) { + if (entry.isDirectory()) { + continue; + } + const filePath = join(entry.parentPath, entry.name).replaceAll('\\', '/'); + // Dirents identify regular files without an additional stat; symlinks + // still need their target checked to preserve directory filtering. + if (!entry.isSymbolicLink() || !tryCreatingDir(filePath, false)) { + files.push(filePath); + } } - continue; + filesByPattern.set(dirPart, files); } - // Use globSync results directly, filter dotfiles if needed - const finalFiles = options.all ? files : filterDotFiles(files, false); - if (isNegated) { - for (const f of finalFiles) { + const finalFiles = options.all || isSingleFileRename ? files : filterDotFiles(files, false); + for (const f of finalFiles) { + if (isNegated && !isSingleFileRename) { allFilesSet.delete(f); - } - } else { - for (const f of finalFiles) { + } else { allFilesSet.add(f); } } @@ -194,6 +199,7 @@ export function copyfiles(sources: string | string[], outPath: string, options: options = createSafeOptions(options); const cb = callback || options.callback; sources = arrify(sources); + const concurrency = options.concurrency === undefined ? Math.min(32, availableParallelism()) : options.concurrency; if (options.verbose || options.stat) { console.time('Execution time'); @@ -204,6 +210,8 @@ export function copyfiles(sources: string | string[], outPath: string, options: errorMsg = 'Please make sure to provide both and , i.e.: "copyfiles "'; } else if (options.flat && options.up) { errorMsg = 'Cannot use --flat in conjunction with --up option.'; + } else if (!Number.isSafeInteger(concurrency) || concurrency < 1) { + errorMsg = 'Concurrency must be a positive safe integer.'; } if (errorMsg) { throwOrCallback(new Error(errorMsg), cb); @@ -236,8 +244,8 @@ export function copyfiles(sources: string | string[], outPath: string, options: } // Set default excludeGlobs only if not provided by user - const excludeGlobs = - Array.isArray(options.exclude) && options.exclude.length > 0 ? options.exclude : ['**/.git/**', '**/node_modules/**']; + const exclude = options.exclude === undefined ? [] : arrify(options.exclude); + const excludeGlobs = exclude.length > 0 ? exclude : ['**/.git/**', '**/node_modules/**']; // Use a Set for deduplication from the start const allFilesSet = getMatchedFiles(sources, excludeGlobs, isSingleFile, isDestFile, options); @@ -247,23 +255,12 @@ export function copyfiles(sources: string | string[], outPath: string, options: } if (options.error && allFilesSet.size < 1) { - const err = new Error('nothing copied'); - if (typeof cb === 'function') { - cb(err); - } else { - throw err; - } + throwOrCallback(new Error('nothing copied'), cb); return; } - let completed = 0; - let hasError = false; - if (allFilesSet.size === 0) { - if (options.verbose || options.stat) { - console.log(`Files copied: 0`); - console.timeEnd('Execution time'); - } + displayStatWhenEnabled(options, 0); if (typeof cb === 'function') { cb(); } @@ -286,32 +283,47 @@ export function copyfiles(sources: string | string[], outPath: string, options: return; } - for (const inFile of allFilesSet) { + const files = allFilesSet.values(); + const createdDirs = new Set(); + const activeStreams = new Set(); + const workerCount = Math.min(concurrency, allFilesSet.size); + let remainingWorkers = workerCount; + let firstError: Error | undefined; + + const copyNext = () => { + const nextFile = files.next(); + if (firstError || nextFile.done) { + // Each worker finishes once, after its current streams have closed. + if (--remainingWorkers === 0) { + if (!firstError) { + displayStatWhenEnabled(options, allFilesSet.size); + } + if (typeof cb === 'function') { + cb(firstError); + } + } + return; + } copyFileStream( - inFile, + nextFile.value, outPath, options, err => { - if (hasError) { - return; - } - if (err) { - hasError = true; - if (typeof cb === 'function') { - cb(err); - } - return; - } - completed++; - if (completed === allFilesSet.size) { - displayStatWhenEnabled(options, allFilesSet.size); - if (typeof cb === 'function') { - cb(); + if (err && !firstError) { + firstError = err; + for (const stream of activeStreams) { + stream.destroy(); } } + copyNext(); }, - isSingleFile && isDestFile, // pass as single rename mode + isSingleFile && isDestFile, + createdDirs, + activeStreams, ); + }; + for (let i = 0; i < workerCount; i++) { + copyNext(); } } @@ -323,42 +335,57 @@ export function copyfiles(sources: string | string[], outPath: string, options: * @param {(e?: Error) => void} cb * @param {Boolean} isSingleFileRename - whether the operation is a single file rename (no glob, dest is not a directory, no *) */ -function copyFileStream(inFile: string, outDir: string, options: CopyFileOptions, cb: (e?: Error) => void, isSingleFileRename = false) { - outDir = outDir.startsWith('~') ? untildify(outDir) : outDir; +function copyFileStream( + inFile: string, + outDir: string, + options: CopyFileOptions, + cb: (e?: Error) => void, + isSingleFileRename: boolean, + createdDirs: Set, + activeStreams: Set, +) { let dest: string; try { dest = getDestinationPath(inFile, outDir, options, isSingleFileRename); + const destDir = dirname(dest); + if (!createdDirs.has(destDir)) { + createDir(destDir); + createdDirs.add(destDir); + } } catch (err) { cb(err as Error); return; } - createDir(dirname(dest)); - if (options.verbose) { console.log('copy:', { from: convertToPosix(inFile), to: convertToPosix(dest) }); } const readStream = createReadStream(inFile); const writeStream = createWriteStream(dest); - - let called = false; - const onceCallback = (err?: Error) => { - if (!called) { - called = true; - cb(err); - } + activeStreams.add(readStream); + activeStreams.add(writeStream); + let error: Error | undefined; + let remaining = 2; + const fail = (err: Error) => { + error ??= err; + readStream.destroy(); + writeStream.destroy(); }; - - readStream.on('error', onceCallback); - writeStream.on('error', onceCallback); - writeStream.on('close', () => { - // Only execute callback if not already called by an error - if (!called) { - onceCallback(); + const close = (stream: ReadStream | WriteStream, completed: boolean) => { + if (!completed && !error) { + const side = stream === readStream ? 'Read' : 'Write'; + fail(new Error(`${side} stream closed before copying ${inFile}`)); } - }); - + activeStreams.delete(stream); + if (--remaining === 0) { + cb(error); + } + }; + readStream.once('error', fail); + writeStream.once('error', fail); + readStream.once('close', () => close(readStream, readStream.readableEnded)); + writeStream.once('close', () => close(writeStream, writeStream.writableFinished)); readStream.pipe(writeStream); } diff --git a/src/interfaces.ts b/src/interfaces.ts index 771a06e..8ce1ba3 100644 --- a/src/interfaces.ts +++ b/src/interfaces.ts @@ -2,6 +2,9 @@ export interface CopyFileOptions { /** Include files & directories beginning with a dot (.) */ all?: boolean; + /** Maximum simultaneous file copies; defaults to os.availableParallelism(), capped at 32 */ + concurrency?: number; + /** Show what would be copied, but do not actually copy any files */ dryRun?: boolean;