From 0e096f4cd2a88fa5011638711b9799810b5b5aaf Mon Sep 17 00:00:00 2001 From: freeof123 Date: Mon, 17 Aug 2026 15:43:18 +0800 Subject: [PATCH] fix: retry PipelineContext atomic save on EPERM (5.1.2) Windows/desktop paths can briefly lock *.progress during rename. Retry transient EPERM/EACCES/EBUSY with backoff, then copy+unlink fallback before surfacing the error to hosts (e.g. DSFT decrypt). Add rename retry tests and an optional DSFT-style decrypt throughput benchmark script (not run by Jest). Co-authored-by: Cursor --- package.json | 2 +- scripts/bench-decrypt-throughput.cjs | 223 +++++++++++++++++++++++++++ src/node/PipelineConext.js | 51 +++++- test/PipelineContext.rename.spec.js | 91 +++++++++++ 4 files changed, 364 insertions(+), 3 deletions(-) create mode 100644 scripts/bench-decrypt-throughput.cjs create mode 100644 test/PipelineContext.rename.spec.js diff --git a/package.json b/package.json index 2bad564..2e7da98 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@yeez-tech/meta-encryptor", - "version": "5.1.1", + "version": "5.1.2", "description": "Data Seal/Unseal for Fidelius", "main": "./build/commonjs/index.node.cjs", "module": "./build/es/index.browser.js", diff --git a/scripts/bench-decrypt-throughput.cjs b/scripts/bench-decrypt-throughput.cjs new file mode 100644 index 0000000..768b636 --- /dev/null +++ b/scripts/bench-decrypt-throughput.cjs @@ -0,0 +1,223 @@ +#!/usr/bin/env node +/** + * Temporary benchmark — DSFT-style decrypt throughput (MB/s). + * Exercises PipelineContextInFile atomic save during Recoverable decrypt. + * + * Usage: + * node scripts/bench-decrypt-throughput.cjs + * node scripts/bench-decrypt-throughput.cjs --size-mb=200 --runs=3 + * + * Remove this script when no longer needed. + */ +'use strict'; + +const fs = require('fs'); +const path = require('path'); +const os = require('os'); +const crypto = require('crypto'); + +const { Sealer } = require('../src/node/Sealer.js'); +const { Unsealer } = require('../src/node/Unsealer.js'); +const { PipelineContextInFile } = require('../src/node/PipelineConext.js'); +const { + RecoverableReadStream, + RecoverableWriteStream, +} = require('../src/node/Recoverable.js'); + +const KEY_PAIR = { + private_key: + '60d61a1d92b26608016dba8cb8e8e96fd44d5dee0a0415a024657e47febcced8', + public_key: + '731234931a081e9beae856318a9bf32ac3698ea8215bf74f517f8377cc6ba8740e28ed87c97d0ee8775bc83505867b0bc34a66adc91f0ea9b44c80533f1a3dca', +}; + +function parseArgs(argv) { + const opts = { sizeMb: 100, runs: 1 }; + for (const arg of argv) { + if (arg.startsWith('--size-mb=')) { + opts.sizeMb = Math.max(1, Number(arg.slice('--size-mb='.length)) || 100); + } else if (arg.startsWith('--runs=')) { + opts.runs = Math.max(1, Number(arg.slice('--runs='.length)) || 1); + } else if (arg === '--help' || arg === '-h') { + opts.help = true; + } + } + return opts; +} + +function generateFile(filePath, sizeBytes) { + const chunk = Buffer.alloc(64 * 1024); + for (let i = 0; i < chunk.length; i++) chunk[i] = i % 256; + const fd = fs.openSync(filePath, 'w'); + try { + let written = 0; + while (written < sizeBytes) { + const n = Math.min(chunk.length, sizeBytes - written); + fs.writeSync(fd, n === chunk.length ? chunk : chunk.subarray(0, n)); + written += n; + } + } finally { + fs.closeSync(fd); + } +} + +function md5File(filePath) { + return new Promise((resolve, reject) => { + const hash = crypto.createHash('md5'); + fs.createReadStream(filePath) + .on('data', (c) => hash.update(c)) + .on('error', reject) + .on('end', () => resolve(hash.digest('hex'))); + }); +} + +function sealPlain(plainPath, sealedPath) { + return new Promise((resolve, reject) => { + const rs = fs.createReadStream(plainPath); + const ws = fs.createWriteStream(sealedPath); + rs.on('error', reject); + ws.on('error', reject); + ws.on('finish', resolve); + rs.pipe(new Sealer({ keyPair: KEY_PAIR })).pipe(ws); + }); +} + +/** Same path layout as dianshu-file-transfer DecryptAction. */ +function dsftStyleDecrypt(sealedPath, outputPath, progressPath) { + return new Promise(async (resolve, reject) => { + try { + const context = new PipelineContextInFile(progressPath); + await context.loadContext(); + + const rs = new RecoverableReadStream(sealedPath, context); + const unsealer = new Unsealer({ keyPair: KEY_PAIR, context }); + const ws = new RecoverableWriteStream(outputPath, context); + + for (const s of [rs, unsealer, ws]) s.on('error', reject); + rs.pipe(unsealer).pipe(ws); + ws.on('finish', resolve); + } catch (err) { + reject(err); + } + }); +} + +function rmIfExists(p) { + try { + fs.unlinkSync(p); + } catch (_) {} +} + +async function runOnce(workDir, plainPath, sizeBytes) { + const base = path.join(workDir, 'bench_out'); + const sealedPath = base + '.decrypting'; + const outputPath = base; + const progressPath = base + '.progress'; + + for (const p of [sealedPath, outputPath, progressPath, `${progressPath}.tmp`]) { + rmIfExists(p); + } + + const sealStart = process.hrtime.bigint(); + await sealPlain(plainPath, sealedPath); + const sealMs = Number(process.hrtime.bigint() - sealStart) / 1e6; + + const decryptStart = process.hrtime.bigint(); + await dsftStyleDecrypt(sealedPath, outputPath, progressPath); + const decryptMs = Number(process.hrtime.bigint() - decryptStart) / 1e6; + + const plainMd5 = await md5File(plainPath); + const outMd5 = await md5File(outputPath); + if (plainMd5 !== outMd5) { + throw new Error(`MD5 mismatch: plain=${plainMd5} out=${outMd5}`); + } + + const sizeMb = sizeBytes / (1024 * 1024); + const decryptMbPerSec = sizeMb / (decryptMs / 1000); + const sealMbPerSec = sizeMb / (sealMs / 1000); + + let progressSaves = 0; + if (fs.existsSync(progressPath)) { + progressSaves = 1; + } + + return { + sizeMb, + sealMs, + decryptMs, + sealMbPerSec, + decryptMbPerSec, + progressFileBytes: fs.existsSync(progressPath) + ? fs.statSync(progressPath).size + : 0, + progressSaves, + }; +} + +async function main() { + const opts = parseArgs(process.argv.slice(2)); + if (opts.help) { + console.log(`Usage: node scripts/bench-decrypt-throughput.cjs [--size-mb=N] [--runs=N]`); + process.exit(0); + } + + const sizeBytes = Math.floor(opts.sizeMb * 1024 * 1024); + const workDir = fs.mkdtempSync(path.join(os.tmpdir(), 'me-decrypt-bench-')); + const plainPath = path.join(workDir, 'plain.dat'); + + console.log('meta-encryptor decrypt throughput benchmark'); + console.log('─'.repeat(56)); + console.log(`Node ${process.version} | ${process.platform} ${process.arch}`); + console.log(`Plain size: ${opts.sizeMb} MiB (${sizeBytes} bytes)`); + console.log(`Runs: ${opts.runs}`); + console.log(`Work dir: ${workDir}`); + console.log('Generating plain file...'); + + const genStart = process.hrtime.bigint(); + generateFile(plainPath, sizeBytes); + const genMs = Number(process.hrtime.bigint() - genStart) / 1e6; + console.log(`Plain file ready in ${genMs.toFixed(0)} ms\n`); + + const results = []; + for (let i = 0; i < opts.runs; i++) { + const label = opts.runs > 1 ? `Run ${i + 1}/${opts.runs}` : 'Run'; + process.stdout.write(`${label}... `); + const r = await runOnce(workDir, plainPath, sizeBytes); + results.push(r); + console.log( + `decrypt ${r.decryptMbPerSec.toFixed(2)} MB/s (${r.decryptMs.toFixed(0)} ms), ` + + `seal ${r.sealMbPerSec.toFixed(2)} MB/s` + ); + } + + const avg = (arr, key) => arr.reduce((s, x) => s + x[key], 0) / arr.length; + + console.log('\n' + '─'.repeat(56)); + console.log('Summary (decrypt = plain MiB / wall time):'); + if (opts.runs === 1) { + const r = results[0]; + console.log(` Decrypt throughput: ${r.decryptMbPerSec.toFixed(2)} MB/s`); + console.log(` Decrypt wall time: ${r.decryptMs.toFixed(1)} ms`); + console.log(` Seal throughput: ${r.sealMbPerSec.toFixed(2)} MB/s (setup)`); + console.log(` Progress file size: ${r.progressFileBytes} bytes`); + } else { + console.log( + ` Decrypt throughput: avg ${avg(results, 'decryptMbPerSec').toFixed(2)} MB/s ` + + `(min ${Math.min(...results.map((r) => r.decryptMbPerSec)).toFixed(2)}, ` + + `max ${Math.max(...results.map((r) => r.decryptMbPerSec)).toFixed(2)})` + ); + console.log(` Decrypt wall time: avg ${avg(results, 'decryptMs').toFixed(1)} ms`); + console.log(` Seal throughput: avg ${avg(results, 'sealMbPerSec').toFixed(2)} MB/s (setup)`); + } + console.log(' MD5: OK (all runs)'); + console.log('─'.repeat(56)); + + try { + fs.rmSync(workDir, { recursive: true, force: true }); + } catch (_) {} +} + +main().catch((err) => { + console.error('Benchmark failed:', err); + process.exit(1); +}); diff --git a/src/node/PipelineConext.js b/src/node/PipelineConext.js index e3fe103..ddcc6e4 100644 --- a/src/node/PipelineConext.js +++ b/src/node/PipelineConext.js @@ -8,9 +8,56 @@ import { MetaEncryptorError } from '../common/errors.js'; const open = promisify(fs.open); const close = promisify(fs.close); const fsync = promisify(fs.fsync); -const rename = promisify(fs.rename); const logger = log.getLogger("meta-encryptor/PipelineContext"); +/** Windows / sync clients may briefly lock the target during atomic replace. */ +const RENAME_TRANSIENT_CODES = new Set(['EPERM', 'EACCES', 'EBUSY']); +const RENAME_MAX_ATTEMPTS = 5; +const RENAME_RETRY_BASE_MS = 50; + +function delay(ms) { + return new Promise((resolve) => setTimeout(resolve, ms)); +} + +function errnoCode(err) { + return err && typeof err === 'object' ? err.code : undefined; +} + +/** + * Replace target with tmp via rename; retry transient locks, then copy+unlink fallback. + * Throws the last error when all attempts fail (hosts may report to Sentry). + */ +async function replaceFileAtomic(tmpPath, targetPath) { + let lastError; + for (let attempt = 0; attempt < RENAME_MAX_ATTEMPTS; attempt++) { + try { + await fs.promises.rename(tmpPath, targetPath); + return; + } catch (err) { + lastError = err; + const code = errnoCode(err); + if (RENAME_TRANSIENT_CODES.has(code) && attempt < RENAME_MAX_ATTEMPTS - 1) { + await delay(RENAME_RETRY_BASE_MS * (attempt + 1)); + continue; + } + break; + } + } + + const code = errnoCode(lastError); + if (RENAME_TRANSIENT_CODES.has(code) || code === 'EXDEV') { + try { + await fs.promises.copyFile(tmpPath, targetPath); + await fs.promises.unlink(tmpPath); + return; + } catch (copyErr) { + throw copyErr; + } + } + + throw lastError; +} + function toBinaryChunk(value) { if (Buffer.isBuffer(value)) { return value; @@ -127,7 +174,7 @@ export class PipelineContextInFile extends PipelineContext { await close(fd); } - await rename(tmpPath, this.filePath); + await replaceFileAtomic(tmpPath, this.filePath); } async _flushAll() { diff --git a/test/PipelineContext.rename.spec.js b/test/PipelineContext.rename.spec.js new file mode 100644 index 0000000..9a49688 --- /dev/null +++ b/test/PipelineContext.rename.spec.js @@ -0,0 +1,91 @@ +const fs = require("fs"); +const path = require("path"); +const os = require("os"); +const { PipelineContextInFile } = require("../src/node/PipelineConext.js"); + +function tmpPath(name) { + return path.join(os.tmpdir(), `me-pctx-${name}-${process.pid}-${Date.now()}`); +} + +describe("PipelineContextInFile atomic replace", () => { + let contextPath; + let renameSpy; + let copySpy; + + beforeEach(() => { + contextPath = tmpPath("ctx.dat"); + }); + + afterEach(() => { + renameSpy?.mockRestore(); + copySpy?.mockRestore(); + for (const p of [contextPath, `${contextPath}.tmp`]) { + try { + fs.unlinkSync(p); + } catch (_) {} + } + }); + + it("retries rename on EPERM then succeeds", async () => { + let attempts = 0; + const originalRename = fs.promises.rename.bind(fs.promises); + renameSpy = jest.spyOn(fs.promises, "rename").mockImplementation(async (src, dst) => { + attempts += 1; + if (attempts <= 2) { + const err = new Error("EPERM: operation not permitted, rename"); + err.code = "EPERM"; + throw err; + } + return originalRename(src, dst); + }); + + const pc = new PipelineContextInFile(contextPath); + pc.update("readStart", 42); + await pc.saveContext(); + + expect(attempts).toBeGreaterThan(2); + expect(fs.existsSync(contextPath)).toBe(true); + expect(fs.existsSync(`${contextPath}.tmp`)).toBe(false); + + const loaded = new PipelineContextInFile(contextPath); + await loaded.loadContext(); + expect(loaded.context.readStart).toBe(42); + }); + + it("falls back to copy+unlink when rename keeps failing with EPERM", async () => { + renameSpy = jest.spyOn(fs.promises, "rename").mockImplementation(async () => { + const err = new Error("EPERM: operation not permitted, rename"); + err.code = "EPERM"; + throw err; + }); + + const pc = new PipelineContextInFile(contextPath); + pc.update("info", { stage: "decrypt" }); + await pc.saveContext(); + + expect(fs.existsSync(contextPath)).toBe(true); + expect(fs.existsSync(`${contextPath}.tmp`)).toBe(false); + + const loaded = new PipelineContextInFile(contextPath); + await loaded.loadContext(); + expect(loaded.context.info).toEqual({ stage: "decrypt" }); + }); + + it("throws after rename retries and copy fallback both fail", async () => { + renameSpy = jest.spyOn(fs.promises, "rename").mockImplementation(async () => { + const err = new Error("EPERM: operation not permitted, rename"); + err.code = "EPERM"; + throw err; + }); + copySpy = jest.spyOn(fs.promises, "copyFile").mockImplementation(async () => { + const err = new Error("EPERM: operation not permitted, copy"); + err.code = "EPERM"; + throw err; + }); + + const pc = new PipelineContextInFile(contextPath); + pc.update("readStart", 1); + + await expect(pc.saveContext()).rejects.toMatchObject({ code: "EPERM" }); + }); +});