Skip to content

Commit db16496

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 789c7fd commit db16496

8 files changed

Lines changed: 320 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 {
@@ -2947,6 +2948,23 @@ function writeFile(path, data, options, callback) {
29472948
if (checkAborted(options.signal, callback))
29482949
return;
29492950

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

‎lib/internal/fs/promises.js‎

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

21202120
checkAborted(options.signal);
21212121

2122+
if (!flush && !isCustomIterable(data) && data.byteLength <= kWriteFileMaxChunkSize) {
2123+
path = getValidatedPath(path);
2124+
await writeFileInOneRoundTrip(path, stringToFlags(flag, 'options.flag'),
2125+
parseFileMode(options.mode, 'mode', 0o666), data);
2126+
checkAborted(options.signal); // An abort during the write still wins.
2127+
return;
2128+
}
2129+
21222130
const fd = await open(path, flag, options.mode);
21232131
let writeOp = writeFileHandle(fd, data, options.signal, options.encoding);
21242132

@@ -2129,6 +2137,33 @@ async function writeFile(path, data, options) {
21292137
return handleFdClose(writeOp, fd.close);
21302138
}
21312139

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

‎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;
@@ -3220,6 +3221,7 @@ class ReadFileJob final : public AsyncWrap, public ThreadPoolWork {
32203221
SET_SELF_SIZE(ReadFileJob)
32213222

32223223
private:
3224+
friend class WriteFileJob;
32233225
static constexpr size_t kUnknownSizeChunk = 64 * 1024;
32243226
static constexpr size_t kMaxReadChunk = 256 * 1024 * 1024;
32253227

@@ -3303,6 +3305,162 @@ class ReadFileJob final : public AsyncWrap, public ThreadPoolWork {
33033305
int fd_ = -1;
33043306
};
33053307

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

4780+
Local<FunctionTemplate> wfj = NewFunctionTemplate(isolate, WriteFileJob::New);
4781+
wfj->InstanceTemplate()->SetInternalFieldCount(
4782+
WriteFileJob::kInternalFieldCount);
4783+
wfj->Inherit(AsyncWrap::GetConstructorTemplate(isolate_data));
4784+
SetProtoMethod(isolate, wfj, "run", WriteFileJob::Run);
4785+
SetConstructorFunction(isolate, target, "WriteFileJob", wfj);
4786+
46224787
// Create FunctionTemplate for FSReqCallback
46234788
Local<FunctionTemplate> fst = NewFunctionTemplate(isolate, NewFSReqCallback);
46244789
fst->InstanceTemplate()->SetInternalFieldCount(
@@ -4693,6 +4858,8 @@ void RegisterExternalReferences(ExternalReferenceRegistry* registry) {
46934858
registry->Register(Open);
46944859
registry->Register(ReadFileJob::New);
46954860
registry->Register(ReadFileJob::Run);
4861+
registry->Register(WriteFileJob::New);
4862+
registry->Register(WriteFileJob::Run);
46964863
registry->Register(OpenFileHandle);
46974864
registry->Register(Read);
46984865
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)