Skip to content

Commit 85ce7e1

Browse files
committed
fs: write files in one thread pool round trip
fs.writeFile(path, data) took three libuv thread pool round trips (open, write, close), each its own request with its own queue wait, completion callback and JS/C++ crossing, and fs.promises.writeFile() did the same through a FileHandle. For the small files applications write most, the round trips are the cost, and each occupies a pool slot that concurrent fs, dns.lookup() and crypto work is also queueing for. Add WriteFileJob next to ReadFileJob: an AsyncWrap + ThreadPoolWork that opens, writes the whole buffer (looping on short writes) and closes as one pool task, keeping the buffer alive until it is done. fs.writeFile() uses it for path arguments without flush; fs.promises.writeFile() additionally keeps data above one write chunk (and iterables) on the FileHandle path, so large writes stay abortable between chunks as before. File descriptors, FileHandles, flush: true and an active VFS keep their existing paths. Behavior is otherwise kept: open failures report syscall 'open' with the path, write failures 'write'; permission errors are delivered through the callback/promise; an abort signalled while the write is in flight is still reported as an AbortError; the job is an FSREQCALLBACK resource for async_hooks and emits the 'write' fs trace event. Tests that used fs.writeFile() as a proxy for open/close trace events, or injected FileHandle faults for path-based writes, are adjusted to keep testing what they test. The job holds the buffer's backing store, so the memory stays valid if the buffer is detached or collected before the write finishes; a resizable ArrayBuffer could still have its pages decommitted by a shrink, so its contents are copied when the job is created. Signed-off-by: Shelley Vohr <shelley.vohr@gmail.com>
1 parent abb365a commit 85ce7e1

8 files changed

Lines changed: 319 additions & 6 deletions

‎lib/fs.js‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -83,6 +83,7 @@ const {
8383
const {
8484
FSReqCallback,
8585
ReadFileJob,
86+
WriteFileJob,
8687
} = binding;
8788
const { toPathIfFileURL } = require('internal/url');
8889
const {
@@ -2946,6 +2947,23 @@ function writeFile(path, data, options, callback) {
29462947
if (checkAborted(options.signal, callback))
29472948
return;
29482949

2950+
if (!flush) {
2951+
// Open + write + close in one thread pool round trip.
2952+
const signal = options.signal;
2953+
path = getValidatedPath(path);
2954+
const job = new WriteFileJob(path, stringToFlags(flag, 'options.flag'),
2955+
parseFileMode(options.mode, 'mode', 0o666), data);
2956+
job.ondone = signal == null ? callback : (err) => {
2957+
// An abort that arrived while the write was in flight still wins.
2958+
callback(signal.aborted && !err ? new AbortError(undefined, { cause: signal.reason }) : err);
2959+
};
2960+
const accessError = job.run(path);
2961+
if (accessError !== undefined) {
2962+
callback(accessError);
2963+
}
2964+
return;
2965+
}
2966+
29492967
fs.open(path, flag, options.mode, (openErr, fd) => {
29502968
if (openErr) {
29512969
callback(openErr);

‎lib/internal/fs/promises.js‎

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2113,6 +2113,14 @@ async function writeFile(path, data, options) {
21132113

21142114
checkAborted(options.signal);
21152115

2116+
if (!flush && !isCustomIterable(data) && data.byteLength <= kWriteFileMaxChunkSize) {
2117+
path = getValidatedPath(path);
2118+
await writeFileInOneRoundTrip(path, stringToFlags(flag, 'options.flag'),
2119+
parseFileMode(options.mode, 'mode', 0o666), data);
2120+
checkAborted(options.signal); // An abort during the write still wins.
2121+
return;
2122+
}
2123+
21162124
const fd = await open(path, flag, options.mode);
21172125
let writeOp = writeFileHandle(fd, data, options.signal, options.encoding);
21182126

@@ -2123,6 +2131,33 @@ async function writeFile(path, data, options) {
21232131
return handleFdClose(writeOp, fd.close);
21242132
}
21252133

2134+
/**
2135+
* Open + write + close as one thread pool round trip.
2136+
* @param {string|Buffer} path Validated path
2137+
* @param {number} flagsNumber
2138+
* @param {number} mode
2139+
* @param {ArrayBufferView} data
2140+
* @returns {Promise<void>}
2141+
*/
2142+
function writeFileInOneRoundTrip(path, flagsNumber, mode, data) {
2143+
return new Promise((resolve, reject) => {
2144+
const job = new binding.WriteFileJob(path, flagsNumber, mode, data);
2145+
job.ondone = (err) => {
2146+
if (err != null) {
2147+
ErrorCaptureStackTrace(err, writeFileInOneRoundTrip);
2148+
reject(err);
2149+
} else {
2150+
resolve();
2151+
}
2152+
};
2153+
const accessError = job.run(path);
2154+
if (accessError !== undefined) {
2155+
ErrorCaptureStackTrace(accessError, writeFileInOneRoundTrip);
2156+
reject(accessError);
2157+
}
2158+
});
2159+
}
2160+
21262161
function isCustomIterable(obj) {
21272162
return isIterable(obj) && !isArrayBufferView(obj) && typeof obj !== 'string';
21282163
}

‎src/node_file.cc‎

Lines changed: 167 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,7 @@ namespace node {
6464
namespace fs {
6565

6666
using v8::Array;
67+
using v8::ArrayBufferView;
6768
using v8::BigInt;
6869
using v8::Context;
6970
using v8::EscapableHandleScope;
@@ -3214,6 +3215,7 @@ class ReadFileJob final : public AsyncWrap, public ThreadPoolWork {
32143215
SET_SELF_SIZE(ReadFileJob)
32153216

32163217
private:
3218+
friend class WriteFileJob;
32173219
static constexpr size_t kUnknownSizeChunk = 64 * 1024;
32183220
static constexpr size_t kMaxReadChunk = 256 * 1024 * 1024;
32193221

@@ -3297,6 +3299,162 @@ class ReadFileJob final : public AsyncWrap, public ThreadPoolWork {
32973299
int fd_ = -1;
32983300
};
32993301

3302+
// Writes a whole buffer to a file in ONE thread pool round trip -- open +
3303+
// write (until everything is written) + close -- for fs.writeFile() and
3304+
// fs.promises.writeFile() with a path, which otherwise pay one round trip per
3305+
// step.
3306+
//
3307+
// JS: const job = new WriteFileJob(path, flags, mode, buffer);
3308+
// job.ondone = (err) => {...}; job.run(path);
3309+
// `err` carries the syscall that failed ('open', 'write' or 'close'); the file
3310+
// descriptor opened here is always closed.
3311+
class WriteFileJob final : public AsyncWrap, public ThreadPoolWork {
3312+
public:
3313+
static void New(const FunctionCallbackInfo<Value>& args) {
3314+
CHECK(args.IsConstructCall());
3315+
Environment* env = Environment::GetCurrent(args);
3316+
CHECK_GE(args.Length(), 4);
3317+
BufferValue path(env->isolate(), args[0]);
3318+
CHECK_NOT_NULL(*path);
3319+
ToNamespacedPath(env, &path);
3320+
CHECK(args[1]->IsInt32());
3321+
CHECK(args[2]->IsInt32());
3322+
CHECK(args[3]->IsArrayBufferView());
3323+
new WriteFileJob(env,
3324+
args.This(),
3325+
path.ToString(),
3326+
args[1].As<Int32>()->Value(),
3327+
args[2].As<Int32>()->Value(),
3328+
args[3].As<ArrayBufferView>());
3329+
}
3330+
3331+
// Returns undefined when the job was scheduled, or the ERR_ACCESS_DENIED
3332+
// error the asynchronous open() would have delivered (nothing is scheduled).
3333+
static void Run(const FunctionCallbackInfo<Value>& args) {
3334+
WriteFileJob* job;
3335+
ASSIGN_OR_RETURN_UNWRAP(&job, args.This());
3336+
Environment* env = job->AsyncWrap::env();
3337+
CHECK(!job->scheduled_);
3338+
BufferValue path(env->isolate(), args[0]);
3339+
CHECK_NOT_NULL(*path);
3340+
ToNamespacedPath(env, &path);
3341+
Local<Value> access_error;
3342+
if (ReadFileJob::OpenPermissionError(env, path, job->flags_)
3343+
.ToLocal(&access_error)) {
3344+
args.GetReturnValue().Set(access_error);
3345+
return;
3346+
}
3347+
job->scheduled_ = true;
3348+
job->ClearWeak();
3349+
FS_ASYNC_TRACE_BEGIN0(UV_FS_WRITE, job)
3350+
job->ScheduleWork();
3351+
}
3352+
3353+
void DoThreadPoolWork() override {
3354+
uv_fs_t req;
3355+
int fd = uv_fs_open(nullptr, &req, path_.c_str(), flags_, mode_, nullptr);
3356+
uv_fs_req_cleanup(&req);
3357+
if (fd < 0) return Fail("open", fd);
3358+
3359+
size_t written = 0;
3360+
while (written < length_) {
3361+
uv_buf_t buf = uv_buf_init(data_ + written,
3362+
static_cast<unsigned int>(std::min<size_t>(
3363+
length_ - written, kMaxWriteChunk)));
3364+
int r = uv_fs_write(nullptr, &req, fd, &buf, 1, -1, nullptr);
3365+
uv_fs_req_cleanup(&req);
3366+
if (r < 0) {
3367+
Fail("write", r);
3368+
break;
3369+
}
3370+
written += static_cast<size_t>(r);
3371+
}
3372+
3373+
int rc = uv_fs_close(nullptr, &req, fd, nullptr);
3374+
uv_fs_req_cleanup(&req);
3375+
if (rc < 0 && error_ == 0) Fail("close", rc);
3376+
}
3377+
3378+
void AfterThreadPoolWork(int status) override {
3379+
Environment* env = AsyncWrap::env();
3380+
std::unique_ptr<WriteFileJob> self(this);
3381+
CHECK(status == 0 || status == UV_ECANCELED);
3382+
FS_ASYNC_TRACE_END0(UV_FS_WRITE, this)
3383+
if (status == UV_ECANCELED || !env->can_call_into_js()) return;
3384+
HandleScope handle_scope(env->isolate());
3385+
Context::Scope context_scope(env->context());
3386+
Isolate* isolate = env->isolate();
3387+
Local<Value> argv[1] = {Null(isolate)};
3388+
if (error_ != 0) {
3389+
argv[0] = UVException(isolate,
3390+
error_,
3391+
syscall_,
3392+
nullptr,
3393+
syscall_ == kOpen ? path_.c_str() : nullptr);
3394+
}
3395+
MakeCallback(env->ondone_string(), arraysize(argv), argv);
3396+
}
3397+
3398+
bool IsNotIndicativeOfMemoryLeakAtExit() const override { return true; }
3399+
void MemoryInfo(MemoryTracker* tracker) const override {
3400+
tracker->TrackField("buffer", buffer_);
3401+
if (copy_) tracker->TrackFieldWithSize("copy", length_);
3402+
}
3403+
SET_MEMORY_INFO_NAME(WriteFileJob)
3404+
SET_SELF_SIZE(WriteFileJob)
3405+
3406+
private:
3407+
static constexpr size_t kMaxWriteChunk = 256 * 1024 * 1024;
3408+
static constexpr const char* kOpen = "open";
3409+
3410+
WriteFileJob(Environment* env,
3411+
Local<Object> object,
3412+
std::string&& path,
3413+
int flags,
3414+
int mode,
3415+
Local<ArrayBufferView> view)
3416+
: AsyncWrap(env, object, AsyncWrap::PROVIDER_FSREQCALLBACK),
3417+
ThreadPoolWork(env, "fs.writefile"),
3418+
path_(std::move(path)),
3419+
flags_(flags),
3420+
mode_(mode) {
3421+
// Holding the backing store keeps the memory valid even if the buffer is
3422+
// detached or collected meanwhile; a resizable buffer can still have its
3423+
// pages decommitted by a shrink, so its contents are copied instead.
3424+
length_ = view->ByteLength();
3425+
backing_store_ = view->Buffer()->GetBackingStore();
3426+
if (backing_store_->IsResizableByUserJavaScript()) {
3427+
copy_.reset(new char[length_]);
3428+
memcpy(copy_.get(),
3429+
static_cast<char*>(backing_store_->Data()) + view->ByteOffset(),
3430+
length_);
3431+
data_ = copy_.get();
3432+
backing_store_.reset();
3433+
} else {
3434+
buffer_.Reset(env->isolate(), view);
3435+
data_ = static_cast<char*>(backing_store_->Data()) + view->ByteOffset();
3436+
}
3437+
MakeWeak();
3438+
}
3439+
3440+
void Fail(const char* syscall, int error) {
3441+
syscall_ = syscall;
3442+
error_ = error;
3443+
}
3444+
3445+
const std::string path_;
3446+
v8::Global<v8::ArrayBufferView> buffer_;
3447+
std::shared_ptr<v8::BackingStore> backing_store_;
3448+
std::unique_ptr<char[]> copy_;
3449+
char* data_ = nullptr;
3450+
size_t length_ = 0;
3451+
const int flags_;
3452+
const int mode_;
3453+
bool scheduled_ = false;
3454+
int error_ = 0;
3455+
const char* syscall_ = nullptr;
3456+
};
3457+
33003458
// Wrapper for readv(2).
33013459
//
33023460
// bytesRead = fs.readv(fd, buffers[, position], callback)
@@ -4613,6 +4771,13 @@ static void CreatePerIsolateProperties(IsolateData* isolate_data,
46134771
SetProtoMethod(isolate, rfj, "run", ReadFileJob::Run);
46144772
SetConstructorFunction(isolate, target, "ReadFileJob", rfj);
46154773

4774+
Local<FunctionTemplate> wfj = NewFunctionTemplate(isolate, WriteFileJob::New);
4775+
wfj->InstanceTemplate()->SetInternalFieldCount(
4776+
WriteFileJob::kInternalFieldCount);
4777+
wfj->Inherit(AsyncWrap::GetConstructorTemplate(isolate_data));
4778+
SetProtoMethod(isolate, wfj, "run", WriteFileJob::Run);
4779+
SetConstructorFunction(isolate, target, "WriteFileJob", wfj);
4780+
46164781
// Create FunctionTemplate for FSReqCallback
46174782
Local<FunctionTemplate> fst = NewFunctionTemplate(isolate, NewFSReqCallback);
46184783
fst->InstanceTemplate()->SetInternalFieldCount(
@@ -4686,6 +4851,8 @@ void RegisterExternalReferences(ExternalReferenceRegistry* registry) {
46864851
registry->Register(Open);
46874852
registry->Register(ReadFileJob::New);
46884853
registry->Register(ReadFileJob::Run);
4854+
registry->Register(WriteFileJob::New);
4855+
registry->Register(WriteFileJob::Run);
46894856
registry->Register(OpenFileHandle);
46904857
registry->Register(Read);
46914858
registry->Register(ReadFileUtf8);

‎test/parallel/test-fs-promises-file-handle-aggregate-errors.js‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -67,7 +67,9 @@ async function checkAggregateError(op) {
6767
tmpdir.refresh();
6868
await checkAggregateError((filePath) => truncate(filePath));
6969
await checkAggregateError((filePath) => readFile(filePath));
70-
await checkAggregateError((filePath) => writeFile(filePath, '123'));
70+
// More than one write chunk (512 KiB), so that writeFile(path) goes through
71+
// a FileHandle as well.
72+
await checkAggregateError((filePath) => writeFile(filePath, '123'.repeat(200_000)));
7173
if (common.isMacOS) {
7274
await checkAggregateError((filePath) => lchmod(filePath, 0o777));
7375
}

‎test/parallel/test-fs-promises-file-handle-close-errors.js‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -62,7 +62,9 @@ async function checkCloseError(op) {
6262
tmpdir.refresh();
6363
await checkCloseError((filePath) => truncate(filePath));
6464
await checkCloseError((filePath) => readFile(filePath));
65-
await checkCloseError((filePath) => writeFile(filePath, '123'));
65+
// More than one write chunk (512 KiB), so that writeFile(path) goes through
66+
// a FileHandle as well.
67+
await checkCloseError((filePath) => writeFile(filePath, '123'.repeat(200_000)));
6668
if (common.isMacOS) {
6769
await checkCloseError((filePath) => lchmod(filePath, 0o777));
6870
}

‎test/parallel/test-fs-promises-file-handle-op-errors.js‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -56,7 +56,9 @@ async function checkOperationError(op) {
5656
tmpdir.refresh();
5757
await checkOperationError((filePath) => truncate(filePath));
5858
await checkOperationError((filePath) => readFile(filePath));
59-
await checkOperationError((filePath) => writeFile(filePath, '123'));
59+
// More than one write chunk (512 KiB), so that writeFile(path) goes through
60+
// a FileHandle as well.
61+
await checkOperationError((filePath) => writeFile(filePath, '123'.repeat(200_000)));
6062
if (common.isMacOS) {
6163
await checkOperationError((filePath) => lchmod(filePath, 0o777));
6264
}

0 commit comments

Comments
 (0)