/* Copyright IBM Corp. 2017 */ import { copy, CopyOptions, ensureDir, remove } from 'fs-extra'; import { createReadStream, createWriteStream, readdir, readFile, stat, Stats, writeFile, PathLike } from 'graceful-fs'; import { join } from 'path'; import { bindNodeCallback, merge, Observable, Observer, of, throwError, UnaryFunction } from 'rxjs'; import { catchError, map, mapTo, mergeMap } from 'rxjs/operators'; import { dir } from 'tmp'; import { VError } from 'verror'; export interface FileDescriptor { root: string; relPath: string; absPath: string; stats: Stats; } export function rxGetFileDescriptor(aFile: string, aRoot: string): Observable { return _rxStats(aFile).pipe( map(stat => ({ root: aRoot, stats: stat, absPath: aFile, relPath: aFile.substring(aRoot.length + 1) }))); } function _doWalkFiles(aRoot: string, aBaseDir: string): Observable { return Observable.create((observer: Observer) => { // check the type const absPath = join(aRoot, aBaseDir); stat(absPath, (statErr, stats) => { // handle the error if (statErr) { if (statErr.code === 'ENOENT') { observer.complete(); } else { observer.error(statErr); } } else { // output the current entry observer.next({ root: aRoot, absPath: absPath, relPath: aBaseDir, stats: stats }); if (stats.isDirectory()) { // read the children readdir(join(aRoot, aBaseDir), (dirErr, files) => { // handle if (dirErr) { observer.error(dirErr); } else { // list this merge(...files.map(f => _doWalkFiles(aRoot, join(aBaseDir, f)))).subscribe(observer); } }); } else { // done observer.complete(); } } }); }); } function _walkFiles(aRoot: string): Observable { return _doWalkFiles(aRoot, ''); } function _writeFileIfChanged(aPath: string, aData: any): Observable { return Observable.create((observer: Observer) => { // test if the file exists readFile(aPath, 'utf-8', (errRead, data) => { // check if the data changes if (!errRead && (data === aData)) { observer.complete(); } else { // override writeFile(aPath, aData, errWrite => { if (errWrite) { observer.error(errWrite); } else { observer.next(aPath); observer.complete(); } }); } }); }); } function _writeFile(aPath: string, aData: any): Observable { return Observable.create((observer: Observer) => { writeFile(aPath, aData, err => { if (err) { observer.error(err); } else { observer.next(aPath); observer.complete(); } }); }); } function _readFile(aPath: string): Observable { return Observable.create((observer: Observer) => { // const t1 = new Date(); setTimeout(() => { readFile(aPath, 'utf-8', (err, data) => { // const t2 = new Date(); // console.log('read', aPath, t2.getTime() - t1.getTime()); if (err) { observer.error(err); } else { observer.next(data); observer.complete(); } }); }, 100); }); // .tap(noop, undefined, () => console.log('readFile', aPath)); } function _readFileBuffer(aPath: string): Observable { return Observable.create((observer: Observer) => { readFile(aPath, (err, data) => { if (err) { observer.error(err); } else { observer.next(data); observer.complete(); } }); }); } // delete const _rxRemove = bindNodeCallback(remove); const _rxDeleteFile = (aPath: string) => _rxRemove(aPath).pipe( mapTo(aPath) ); // mkdir const _rxEnsureDir = bindNodeCallback(ensureDir); const _rxMkdirp = (aPath: string) => _rxEnsureDir(aPath).pipe( mapTo(aPath) ); // stats const _rxStats: UnaryFunction> = bindNodeCallback(stat); // copy const _rxCopy = bindNodeCallback(copy); const _rxCopyFile = (aSrc: string, aDst: string, aOverride?: boolean) => _rxCopy(aSrc, aDst, { overwrite: !!aOverride, preserveTimestamps: true }).pipe( mapTo(aDst) ); export function rxCopyFileStream(aSrc: string, aDst: string): Observable { return Observable.create((aObserver: Observer) => { const onError = (err: any) => aObserver.error(err); const onOk = () => { aObserver.next(aDst); aObserver.complete(); }; const rd = createReadStream(aSrc); rd.on('error', onError); const wr = createWriteStream(aDst); wr.on('error', onError); wr.on('close', onOk); rd.pipe(wr); }); } const _rxReaddir = bindNodeCallback(readdir); function _validateNotExists(aDir: string): Observable { return _rxStats(aDir).pipe( catchError(err => of(err)), mergeMap(data => (data instanceof Error) ? of(aDir) : throwError(`Directory [${aDir}] already exists. Use --override to overwrite it.`)) ); } function _rxTempDir(): Observable { return Observable.create((o: Observer) => { // dispatch dir((tmpError, tmpPath) => { // test for error if (tmpError) { o.error(new VError(tmpError, 'Unable to create temporary directory.')); } else { o.next(tmpPath); o.complete(); } }); }); } export { _walkFiles as walkFiles, _walkFiles as rxWalkFiles, _writeFile as writeFile, _writeFileIfChanged as writeFileIfChanged, _readFile as readFile, _rxMkdirp as rxMkdirp, _rxStats as rxStats, _rxReaddir as rxReaddir, _rxCopyFile as rxCopyFile, _rxDeleteFile as rxDeleteFile, _rxTempDir as rxTempDir, _validateNotExists as rxValidateNotExists };