Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
116 changes: 99 additions & 17 deletions src/runtime/webcore/Blob.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4675,7 +4675,7 @@ pub(crate) fn write_file_internal(
// except if you're on Windows. Windows I/O is slower. Let's not even try.
#[cfg(not(windows))]
{
let mut needs_async = false;
let mut deferred: Option<Deferred> = None;
let fast_path_ok = matches!(*path_or_blob, PathOrBlob::Path(_))
|| (matches!(*path_or_blob, PathOrBlob::Blob(ref b)
if b.offset.get() == 0 && !b.is_s3()
Expand All @@ -4702,18 +4702,27 @@ pub(crate) fn write_file_internal(
cx.global(),
pathlike,
&str,
&mut needs_async,
&mut deferred,
)
} else {
write_string_to_file_fast::<false>(
cx.global(),
pathlike,
&str,
&mut needs_async,
&mut deferred,
)
};
if !needs_async {
return Ok(result);
match deferred {
None => return Ok(result),
Some(Deferred::Rest(resume)) => {
return write_file_after_would_block(
cx,
path_or_blob,
resume,
&options,
);
}
Some(Deferred::Whole) => {}
}
}
} else if let Some(buffer_view) = data.as_array_buffer(cx.global()) {
Expand All @@ -4733,18 +4742,27 @@ pub(crate) fn write_file_internal(
cx.global(),
pathlike,
buffer_view.byte_slice(),
&mut needs_async,
&mut deferred,
)
} else {
write_bytes_to_file_fast::<false>(
cx.global(),
pathlike,
buffer_view.byte_slice(),
&mut needs_async,
&mut deferred,
)
};
if !needs_async {
return Ok(result);
match deferred {
None => return Ok(result),
Some(Deferred::Rest(resume)) => {
return write_file_after_would_block(
cx,
path_or_blob,
resume,
&options,
);
}
Some(Deferred::Whole) => {}
}
}
}
Expand Down Expand Up @@ -5135,12 +5153,62 @@ pub(crate) fn write_file(global_this: &JSGlobalObject, callframe: &CallFrame) ->

const WRITE_PERMISSIONS: bun_sys::Mode = 0o664;

/// Why the synchronous attempt of `Bun.write` left the write to the async path.
#[cfg(not(windows))]
enum Deferred {
/// Nothing is written and nothing is open: the async path does the whole write.
Whole,
/// A write returned `EAGAIN` after some bytes: the async path does the rest.
Rest(Resume),
}

/// Where the synchronous attempt stopped. The async path gets only what is not written yet.
#[cfg(not(windows))]
struct Resume {
/// How many bytes are written. More than 0.
written: usize,
/// The bytes that are not written yet.
tail: Vec<u8>,
/// The fd the attempt opened for a path. It stays open: a close ends a FIFO reader's data.
fd: Option<bun_sys::CloseOnDrop>,
}

/// Schedules a `WriteFile` for `resume.tail`. Its promise resolves with the whole write's count.
#[cfg(not(windows))]
fn write_file_after_would_block(
cx: &bun_jsc::JsThread<'_>,
path_or_blob: &mut PathOrBlob,
resume: Resume,
options: &WriteFileOptions,
) -> JsResult<JSValue> {
let destination_blob = match path_or_blob {
// A local file: the synchronous attempt opened it.
PathOrBlob::Path(path) => Blob::find_or_create_file_from_path(path, cx.global(), false),
PathOrBlob::Blob(blob) => blob.dupe(),
};
let mut write_file = write_file_mod::WriteFile::create(
destination_blob,
Blob::init(resume.tail, cx.global()),
options.mkdirp_if_not_exists.unwrap_or(true),
)
.expect("unreachable");
write_file.resume_from(resume.written, resume.fd);
let promise = Box::new(WriteFilePromise {
promise: jsc::JSPromiseStrong::init(cx.global()),
global_this: cx.global(),
});
let promise_value = promise.promise.value();
promise_value.ensure_still_alive();
write_file_mod::WriteFile::schedule(write_file, promise, cx);
Ok(promise_value)
}

#[cfg(not(windows))]
fn write_string_to_file_fast<const NEEDS_OPEN: bool>(
global_this: &JSGlobalObject,
pathlike: &PathOrFileDescriptor,
str: &BunString,
needs_async: &mut bool,
deferred: &mut Option<Deferred>,
) -> JSValue {
let fd: Fd = if !NEEDS_OPEN {
pathlike.fd()
Expand All @@ -5156,7 +5224,7 @@ fn write_string_to_file_fast<const NEEDS_OPEN: bool>(
bun_sys::Result::Ok(result) => result,
bun_sys::Result::Err(err) => {
if err.get_errno() == bun_sys::E::ENOENT {
*needs_async = true;
*deferred = Some(Deferred::Whole);
return JSValue::ZERO;
}
return JSPromise::rejected_promise(
Expand All @@ -5169,7 +5237,7 @@ fn write_string_to_file_fast<const NEEDS_OPEN: bool>(
};

// Declared before the truncate guard so it drops *after* it (close runs last).
let _close = NEEDS_OPEN.then(|| bun_sys::CloseOnDrop::new(fd));
let mut close = NEEDS_OPEN.then(|| bun_sys::CloseOnDrop::new(fd));

// scopeguard's closure captures borrows at construction, conflicting
// with later `written += ...` / `truncate = false`. Route through `Cell`
Expand Down Expand Up @@ -5200,7 +5268,14 @@ fn write_string_to_file_fast<const NEEDS_OPEN: bool>(
bun_sys::Result::Err(err) => {
truncate.set(false);
if err.get_errno() == bun_sys::E::EAGAIN {
*needs_async = true;
*deferred = Some(match written.get() {
0 => Deferred::Whole,
written => Deferred::Rest(Resume {
written,
tail: remain.to_vec(),
fd: close.take(),
}),
});
return JSValue::ZERO;
}
let err_js = if !NEEDS_OPEN {
Expand All @@ -5222,7 +5297,7 @@ fn write_bytes_to_file_fast<const NEEDS_OPEN: bool>(
global_this: &JSGlobalObject,
pathlike: &PathOrFileDescriptor,
bytes: &[u8],
needs_async: &mut bool,
deferred: &mut Option<Deferred>,
) -> JSValue {
let fd: Fd = if !NEEDS_OPEN {
pathlike.fd()
Expand All @@ -5241,7 +5316,7 @@ fn write_bytes_to_file_fast<const NEEDS_OPEN: bool>(
bun_sys::Result::Ok(result) => result,
bun_sys::Result::Err(err) => {
if err.get_errno() == bun_sys::E::ENOENT {
*needs_async = true;
*deferred = Some(Deferred::Whole);
return JSValue::ZERO;
}
return JSPromise::rejected_promise(
Expand All @@ -5257,7 +5332,7 @@ fn write_bytes_to_file_fast<const NEEDS_OPEN: bool>(

let truncate = NEEDS_OPEN || bytes.is_empty();
let mut written: usize = 0;
let _close = NEEDS_OPEN.then(|| bun_sys::CloseOnDrop::new(fd));
let mut close = NEEDS_OPEN.then(|| bun_sys::CloseOnDrop::new(fd));

let mut remain = bytes;
while !remain.is_empty() {
Expand All @@ -5271,7 +5346,14 @@ fn write_bytes_to_file_fast<const NEEDS_OPEN: bool>(
}
bun_sys::Result::Err(err) => {
if err.get_errno() == bun_sys::E::EAGAIN {
*needs_async = true;
*deferred = Some(match written {
0 => Deferred::Whole,
written => Deferred::Rest(Resume {
written,
tail: remain.to_vec(),
fd: close.take(),
}),
});
return JSValue::ZERO;
}
let err_js = if !NEEDS_OPEN {
Expand Down
30 changes: 28 additions & 2 deletions src/runtime/webcore/blob/write_file.rs
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,14 @@ pub(crate) struct WriteFile {
pub(crate) state: AtomicU8, // ClosingState

pub(crate) total_written: usize,
/// Bytes of this write that the synchronous attempt wrote before the job (`resume_from`).
pub(crate) base_written: usize,
/// The fd the synchronous attempt opened, until `run` takes it. Closes on drop if never run.
#[cfg(not(windows))]
pub(crate) adopted_fd: Option<sys::CloseOnDrop>,
/// `opened_fd` was opened without `O_TRUNC`: cut the file at the last byte written.
#[cfg(not(windows))]
pub(crate) truncate_on_finish: bool,

#[cfg(not(windows))]
pub(crate) could_block: bool,
Expand Down Expand Up @@ -306,13 +314,23 @@ impl WriteFile {
io_parking: super::IoParking::new(),
state: AtomicU8::new(ClosingState::Running as u8),
total_written: 0,
base_written: 0,
adopted_fd: None,
truncate_on_finish: false,
could_block: false,
close_after_io: false,
mkdirp_if_not_exists,
};
Ok(write_file)
}

/// The rest of a write whose synchronous attempt got `EAGAIN` after `written` bytes, on `fd`.
#[cfg(not(windows))]
pub(crate) fn resume_from(&mut self, written: usize, fd: Option<sys::CloseOnDrop>) {
self.base_written = written;
self.adopted_fd = fd;
}

// reshaped for borrowck — take (off, len) here and re-derive the slice
// internally so callers don't hold a borrow of self across the &mut self call.
#[cfg(not(windows))]
Expand Down Expand Up @@ -348,7 +366,7 @@ impl WriteFile {
let cb: WriteFileOnWriteFileCallback = WriteFilePromise::run;
let cb_ctx = bun_core::heap::into_raw(promise).cast::<c_void>();
let system_error = this.system_error.take();
let total_written = this.total_written;
let total_written = this.base_written + this.total_written;
drop(this);

if let Some(err) = system_error {
Expand All @@ -373,6 +391,10 @@ impl WriteFile {
#[cfg(not(windows))]
{
self.io_task = Some(task);
if let Some(fd) = self.adopted_fd.take() {
self.opened_fd = fd.release();
self.truncate_on_finish = true;
}
self.run_async();
}
}
Expand Down Expand Up @@ -400,6 +422,10 @@ impl WriteFile {
bun_output::scoped_log!(WriteFile, "WriteFile.onFinish()");

let close_after_io = self.close_after_io;
if core::mem::take(&mut self.truncate_on_finish) {
let len = self.base_written + self.total_written;
let _ = sys::ftruncate(self.opened_fd, i64::try_from(len).expect("int cast"));
}
if self.do_close(self.is_allowed_to_close()) {
return;
}
Expand Down Expand Up @@ -460,7 +486,7 @@ impl WriteFile {
if !self.could_block && self.bytes_blob.shared_view().len() > 1024 {
let _ = sys::preallocate_file(
fd.native(),
0,
i64::try_from(self.base_written).expect("int cast"),
i64::try_from(self.bytes_blob.shared_view().len()).expect("int cast"),
); // we don't care if it fails.
}
Expand Down
5 changes: 5 additions & 0 deletions src/sys/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8503,6 +8503,11 @@ impl CloseOnDrop {
pub fn new(fd: Fd) -> Self {
Self(fd)
}
/// The fd, still open. The caller closes it from here on.
#[inline]
pub fn release(self) -> Fd {
core::mem::ManuallyDrop::new(self).0
}
}
impl Drop for CloseOnDrop {
#[inline]
Expand Down
Loading
Loading