use core::ffi::c_void;
use core::mem;
use bun_collections::ByteVecExt;
use bun_core::OOM;
use bun_ptr::LaunderedSelf; #[cfg(windows)]
use bun_sys::ReturnCodeExt as _;
#[cfg(windows)]
use bun_sys::windows::libuv as uv;
#[cfg(windows)]
use bun_sys::windows::libuv::UvHandle as _;
use bun_sys::{self as sys, Fd};
use crate::{EventLoopHandle, FilePollFlag, FilePollKind, FilePollRef, Owner, PollTag};
use crate::pipes::{FileType, PollOrFd};
#[cfg(windows)]
use crate::source::Source;
bun_core::define_scoped_log!(log, PipeWriter, hidden);
#[inline]
fn oom_err() -> sys::Error {
sys::Error::from_code(sys::E::ENOMEM, sys::Tag::write)
}
#[derive(Clone)]
pub enum WriteResult {
Done(usize),
Wrote(usize),
Pending(usize),
Err(sys::Error),
}
#[derive(Copy, Clone, Eq, PartialEq)]
pub enum WriteStatus {
EndOfFile,
Drained,
Pending,
}
pub trait PosixPipeWriter {
fn get_fd(&self) -> Fd;
fn get_buffer(&self) -> &[u8];
fn on_write(&mut self, written: usize, status: WriteStatus);
fn register_poll(&mut self);
const HAS_REGISTER_POLL: bool = true;
fn on_error(&mut self, err: sys::Error);
fn get_file_type(&self) -> FileType;
fn get_force_sync(&self) -> bool;
fn handle(&self) -> &PollOrFd;
fn try_write(&self, force_sync: bool, buf: &[u8]) -> WriteResult {
let ft = if !force_sync {
self.get_file_type()
} else {
FileType::File
};
match ft {
FileType::NonblockingPipe | FileType::File => {
self.try_write_with_write_fn(buf, sys::write)
}
FileType::Pipe => self.try_write_with_write_fn(buf, write_to_blocking_pipe),
FileType::Socket => self.try_write_with_write_fn(buf, sys::send_non_block),
}
}
fn try_write_with_write_fn(
&self,
buf: &[u8],
write_fn: fn(Fd, &[u8]) -> sys::Result<usize>,
) -> WriteResult {
let fd = self.get_fd();
if fd == Fd::INVALID {
return WriteResult::Done(0);
}
let mut offset: usize = 0;
while offset < buf.len() {
match write_fn(fd, &buf[offset..]) {
sys::Result::Err(err) => {
if err.is_retry() {
return WriteResult::Pending(offset);
}
return WriteResult::Err(err);
}
sys::Result::Ok(wrote) => {
offset += wrote;
if wrote == 0 {
return WriteResult::Done(offset);
}
}
}
}
WriteResult::Wrote(offset)
}
fn on_poll(&mut self, size_hint: isize, received_hup: bool) {
let buffer_len = self.get_buffer().len();
log!("onPoll({})", buffer_len);
if buffer_len == 0 && !received_hup {
let self_addr = std::ptr::from_ref(self).cast::<()>() as usize;
log!(
"PosixPipeWriter(0x{:x}) handle={}",
self_addr,
self.handle().tag_name()
);
if let PollOrFd::Poll(poll) = self.handle() {
log!(
"PosixPipeWriter(0x{:x}) got 0, registered state = {}",
self_addr,
poll.is_registered()
);
}
return;
}
let max_write = if size_hint > 0 && self.get_file_type().is_blocking() {
usize::try_from(size_hint).expect("int cast")
} else {
usize::MAX
};
match self.drain_buffered_data(max_write, received_hup) {
WriteResult::Pending(wrote) => {
if wrote > 0 {
self.on_write(wrote, WriteStatus::Pending);
}
if Self::HAS_REGISTER_POLL {
self.register_poll();
}
}
WriteResult::Wrote(amt) => {
self.on_write(amt, WriteStatus::Drained);
}
WriteResult::Err(err) => {
self.on_error(err);
}
WriteResult::Done(amt) => {
self.on_write(amt, WriteStatus::EndOfFile);
}
}
}
fn drain_buffered_data(&self, max_write_size: usize, received_hup: bool) -> WriteResult {
let _ = received_hup;
let buf_len = self.get_buffer().len();
let limit = if max_write_size < buf_len && max_write_size > 0 {
max_write_size
} else {
buf_len
};
let mut drained: usize = 0;
while drained < limit {
let force_sync = self.get_force_sync();
match self.try_write(force_sync, &self.get_buffer()[drained..limit]) {
WriteResult::Pending(pending) => {
drained += pending;
return WriteResult::Pending(drained);
}
WriteResult::Wrote(amt) => {
drained += amt;
}
WriteResult::Err(err) => {
return WriteResult::Err(err);
}
WriteResult::Done(amt) => {
drained += amt;
return WriteResult::Done(drained);
}
}
}
WriteResult::Wrote(drained)
}
}
fn write_to_blocking_pipe(fd: Fd, buf: &[u8]) -> sys::Result<usize> {
#[cfg(any(target_os = "linux", target_os = "android"))]
{
if bun_sys::linux::RWFFlagSupport::is_maybe_supported() {
return sys::write_nonblocking(fd, buf);
}
}
match bun_core::is_writable(fd) {
bun_core::Pollable::Ready | bun_core::Pollable::Hup => sys::write(fd, buf),
bun_core::Pollable::NotReady => sys::Result::Err(sys::Error::retry()),
}
}
pub trait PosixBufferedWriterParent {
const POLL_OWNER_TAG: PollTag;
unsafe fn on_write(this: *mut Self, amount: usize, status: WriteStatus);
unsafe fn on_error(this: *mut Self, err: sys::Error);
const HAS_ON_CLOSE: bool;
unsafe fn on_close(_this: *mut Self) {}
unsafe fn get_buffer<'a>(this: *mut Self) -> &'a [u8];
const HAS_ON_WRITABLE: bool;
unsafe fn on_writable(_this: *mut Self) {}
unsafe fn event_loop(this: *mut Self) -> EventLoopHandle;
}
pub struct PosixBufferedWriter<Parent: PosixBufferedWriterParent> {
pub handle: PollOrFd,
pub parent: Option<bun_ptr::ParentRef<Parent>>,
pub is_done: bool,
pub pollable: bool,
pub closed_without_reporting: bool,
pub close_fd: bool,
}
impl<Parent: PosixBufferedWriterParent> Default for PosixBufferedWriter<Parent> {
fn default() -> Self {
Self {
handle: PollOrFd::Closed,
parent: None, is_done: false,
pollable: false,
closed_without_reporting: false,
close_fd: true,
}
}
}
impl<Parent: PosixBufferedWriterParent> PosixPipeWriter for PosixBufferedWriter<Parent> {
fn get_fd(&self) -> Fd {
self.handle.get_fd()
}
fn get_buffer(&self) -> &[u8] {
self.get_buffer_internal()
}
fn on_write(&mut self, written: usize, status: WriteStatus) {
self._on_write(written, status);
}
fn register_poll(&mut self) {
Self::register_poll(self);
}
fn on_error(&mut self, err: sys::Error) {
self._on_error(err);
}
fn get_file_type(&self) -> FileType {
Self::get_file_type(self)
}
fn get_force_sync(&self) -> bool {
false
}
fn handle(&self) -> &PollOrFd {
&self.handle
}
}
unsafe impl<Parent: PosixBufferedWriterParent> bun_ptr::LaunderedSelf
for PosixBufferedWriter<Parent>
{
}
impl<Parent: PosixBufferedWriterParent> PosixBufferedWriter<Parent> {
#[inline]
fn parent(&self) -> *mut Parent {
self.parent
.map_or(core::ptr::null_mut(), bun_ptr::ParentRef::as_mut_ptr)
}
#[inline]
fn parent_event_loop(&self) -> EventLoopHandle {
unsafe { Parent::event_loop(self.parent()) }
}
#[inline]
fn parent_on_error(&self, err: sys::Error) {
unsafe { Parent::on_error(self.parent(), err) }
}
pub fn memory_cost(&self) -> usize {
mem::size_of::<Self>()
}
pub fn create_poll(&mut self, fd: Fd) -> FilePollRef {
FilePollRef::init(
self.parent_event_loop(),
fd,
Owner::new(Parent::POLL_OWNER_TAG, std::ptr::from_mut(self).cast()),
)
}
pub fn get_poll(&self) -> Option<FilePollRef> {
self.handle.get_poll()
}
pub fn get_file_type(&self) -> FileType {
let Some(poll) = self.get_poll() else {
return FileType::File;
};
poll.file_type()
}
pub fn get_fd(&self) -> Fd {
self.handle.get_fd()
}
fn _on_error(&mut self, err: sys::Error) {
debug_assert!(!err.is_retry());
self.parent_on_error(err);
self.close();
}
pub fn get_force_sync(&self) -> bool {
false
}
fn _on_write(&mut self, written: usize, status: WriteStatus) {
let this: *mut Self = core::hint::black_box(core::ptr::from_mut(self));
let was_done = Self::r(this).is_done;
let parent = Self::r(this).parent();
if status == WriteStatus::EndOfFile && !was_done {
Self::r(this).close_without_reporting();
}
unsafe { Parent::on_write(parent, written, status) };
core::hint::black_box(this);
if status == WriteStatus::EndOfFile && !was_done {
Self::r(this).close();
}
}
fn _on_writable(&mut self) {
if self.is_done {
return;
}
if Parent::HAS_ON_WRITABLE {
unsafe { Parent::on_writable(self.parent()) };
}
}
pub fn register_poll(&mut self) {
let Some(poll) = self.get_poll() else { return };
let loop_ = self.parent_event_loop().loop_();
match poll.register_with_fd(loop_, FilePollKind::Writable, poll.fd()) {
sys::Result::Err(err) => {
self._on_error(err);
}
sys::Result::Ok(()) => {}
}
}
pub fn has_ref(&self) -> bool {
if self.is_done {
return false;
}
let Some(poll) = self.get_poll() else {
return false;
};
poll.can_enable_keeping_process_alive()
}
pub fn enable_keeping_process_alive(&self, event_loop: EventLoopHandle) {
self.update_ref(event_loop, true);
}
pub fn disable_keeping_process_alive(&self, event_loop: EventLoopHandle) {
self.update_ref(event_loop, false);
}
fn get_buffer_internal(&self) -> &[u8] {
unsafe { Parent::get_buffer(self.parent()) }
}
pub fn end(&mut self) {
if self.is_done {
return;
}
self.is_done = true;
self.close();
}
fn close_without_reporting(&mut self) {
if self.get_fd() != Fd::INVALID {
debug_assert!(!self.closed_without_reporting);
self.closed_without_reporting = true;
if self.close_fd {
self.handle.close(None, None::<fn(*mut c_void)>);
}
}
}
pub fn close(&mut self) {
if Parent::HAS_ON_CLOSE {
if self.closed_without_reporting {
self.closed_without_reporting = false;
unsafe { Parent::on_close(self.parent()) };
} else {
let parent = self.parent();
self.handle.close_impl(
Some(parent.cast()),
Some(|ctx: *mut c_void| unsafe { Parent::on_close(ctx.cast::<Parent>()) }),
self.close_fd,
);
}
}
}
pub fn update_ref(&self, event_loop: EventLoopHandle, value: bool) {
let Some(poll) = self.get_poll() else { return };
poll.set_keeping_process_alive(event_loop, value);
}
pub fn set_parent(&mut self, parent: *mut Parent) {
self.parent = Some(bun_ptr::ParentRef::from(
core::ptr::NonNull::new(parent).expect("set_parent: parent must not be null"),
));
let owner = std::ptr::from_mut(self).cast::<c_void>();
self.handle
.set_owner(Owner::new(Parent::POLL_OWNER_TAG, owner.cast()));
}
pub fn write(&mut self) {
self.on_poll(0, false);
}
pub fn watch(&mut self) {
if self.pollable {
if matches!(self.handle, PollOrFd::Fd(_)) {
let fd = self.get_fd();
self.handle = PollOrFd::Poll(self.create_poll(fd));
}
Self::register_poll(self);
}
}
pub fn start(&mut self, rawfd: Fd, pollable: bool) -> sys::Result<()> {
let fd = rawfd;
self.pollable = pollable;
if !pollable {
debug_assert!(!matches!(self.handle, PollOrFd::Poll(_)));
self.handle = PollOrFd::Fd(fd);
return sys::Result::Ok(());
}
let poll = match self.get_poll() {
Some(p) => p,
None => {
let p = self.create_poll(fd);
self.handle = PollOrFd::Poll(p);
p
}
};
let loop_ = self.parent_event_loop().loop_();
match poll.register_with_fd(loop_, FilePollKind::Writable, fd) {
sys::Result::Err(err) => {
return sys::Result::Err(err);
}
sys::Result::Ok(()) => {
let event_loop = self.parent_event_loop();
self.enable_keeping_process_alive(event_loop);
}
}
sys::Result::Ok(())
}
}
pub trait PosixStreamingWriterParent {
const POLL_OWNER_TAG: PollTag;
unsafe fn on_write(this: *mut Self, amount: usize, status: WriteStatus);
unsafe fn on_error(this: *mut Self, err: sys::Error);
const HAS_ON_READY: bool;
unsafe fn on_ready(_this: *mut Self) {}
unsafe fn on_close(this: *mut Self);
unsafe fn event_loop(this: *mut Self) -> EventLoopHandle;
unsafe fn loop_(this: *mut Self) -> *mut bun_uws_sys::Loop;
}
pub struct PosixStreamingWriter<Parent: PosixStreamingWriterParent> {
pub outgoing: StreamBuffer,
pub handle: PollOrFd,
pub parent: *mut Parent,
pub is_done: bool,
pub closed_without_reporting: bool,
pub force_sync: bool,
}
impl<Parent: PosixStreamingWriterParent> Default for PosixStreamingWriter<Parent> {
fn default() -> Self {
Self {
outgoing: StreamBuffer::default(),
handle: PollOrFd::Closed,
parent: core::ptr::null_mut(), is_done: false,
closed_without_reporting: false,
force_sync: false,
}
}
}
impl<Parent: PosixStreamingWriterParent> PosixPipeWriter for PosixStreamingWriter<Parent> {
fn get_fd(&self) -> Fd {
self.handle.get_fd()
}
fn get_buffer(&self) -> &[u8] {
self.outgoing.slice()
}
fn on_write(&mut self, written: usize, status: WriteStatus) {
self._on_write(written, status);
}
fn register_poll(&mut self) {
Self::register_poll(self);
}
fn on_error(&mut self, err: sys::Error) {
self._on_error(err);
}
fn get_file_type(&self) -> FileType {
Self::get_file_type(self)
}
fn get_force_sync(&self) -> bool {
self.force_sync
}
fn handle(&self) -> &PollOrFd {
&self.handle
}
}
unsafe impl<Parent: PosixStreamingWriterParent> bun_ptr::LaunderedSelf
for PosixStreamingWriter<Parent>
{
}
impl<Parent: PosixStreamingWriterParent> PosixStreamingWriter<Parent> {
const CHUNK_SIZE: usize = 4096;
#[inline]
fn parent(&self) -> *mut Parent {
self.parent
}
#[inline]
fn parent_on_write(&self, amount: usize, status: WriteStatus) {
unsafe { Parent::on_write(self.parent(), amount, status) }
}
pub fn get_force_sync(&self) -> bool {
self.force_sync
}
pub fn memory_cost(&self) -> usize {
mem::size_of::<Self>() + self.outgoing.memory_cost()
}
pub fn get_poll(&self) -> Option<FilePollRef> {
self.handle.get_poll()
}
pub fn get_fd(&self) -> Fd {
self.handle.get_fd()
}
pub fn get_file_type(&self) -> FileType {
let Some(poll) = self.get_poll() else {
return FileType::File;
};
poll.file_type()
}
pub fn has_pending_data(&self) -> bool {
self.outgoing.is_not_empty()
}
pub fn should_buffer(&self, addition: usize) -> bool {
!self.force_sync && self.outgoing.size() + addition < Self::CHUNK_SIZE
}
pub fn get_buffer(&self) -> &[u8] {
self.outgoing.slice()
}
fn _on_error(&mut self, err: sys::Error) {
debug_assert!(!err.is_retry());
self.close_without_reporting();
self.is_done = true;
self.outgoing.reset();
unsafe { Parent::on_error(self.parent(), err) };
self.close();
}
fn _on_write(&mut self, written: usize, status: WriteStatus) {
self.outgoing.wrote(written);
if status == WriteStatus::EndOfFile && !self.is_done {
self.close_without_reporting();
}
if self.outgoing.is_empty() {
self.outgoing.cursor = 0;
if status != WriteStatus::EndOfFile {
self.outgoing.maybe_shrink();
}
self.outgoing.list.clear();
}
self.parent_on_write(written, status);
}
pub fn set_parent(&mut self, parent: *mut Parent) {
self.parent = parent;
let owner = std::ptr::from_mut(self).cast::<c_void>();
self.handle
.set_owner(Owner::new(Parent::POLL_OWNER_TAG, owner.cast()));
}
fn _on_writable(&mut self) {
if self.is_done || self.closed_without_reporting {
return;
}
self.outgoing.reset();
if Parent::HAS_ON_READY {
unsafe { Parent::on_ready(self.parent()) };
}
}
fn close_without_reporting(&mut self) {
if self.get_fd() != Fd::INVALID {
debug_assert!(!self.closed_without_reporting);
self.closed_without_reporting = true;
self.handle.close(None, None::<fn(*mut c_void)>);
}
}
fn register_poll(&mut self) {
let Some(poll) = self.get_poll() else { return };
let loop_ = unsafe { Parent::loop_(self.parent()) }.cast();
match poll.register_with_fd(loop_, FilePollKind::Writable, poll.fd()) {
sys::Result::Err(err) => {
let this: *mut Self = core::hint::black_box(core::ptr::from_mut(self));
unsafe { Parent::on_error(Self::r(this).parent(), err) };
Self::r(this).close();
}
sys::Result::Ok(()) => {}
}
}
pub fn write_utf16(&mut self, buf: &[u16]) -> WriteResult {
if self.is_done || self.closed_without_reporting {
return WriteResult::Done(0);
}
let before_len = self.outgoing.size();
if self.outgoing.write_utf16(buf).is_err() {
return WriteResult::Err(oom_err());
}
let buf_len = self.outgoing.size() - before_len;
self.maybe_write_newly_buffered_data(buf_len)
}
pub fn write_latin1(&mut self, buf: &[u8]) -> WriteResult {
if self.is_done || self.closed_without_reporting {
return WriteResult::Done(0);
}
if bun_core::strings::is_all_ascii(buf) {
return self.write(buf);
}
let before_len = self.outgoing.size();
const CHECK_ASCII: bool = false;
if self.outgoing.write_latin1::<CHECK_ASCII>(buf).is_err() {
return WriteResult::Err(oom_err());
}
let buf_len = self.outgoing.size() - before_len;
self.maybe_write_newly_buffered_data(buf_len)
}
fn maybe_write_newly_buffered_data(&mut self, buf_len: usize) -> WriteResult {
debug_assert!(!self.is_done);
if self.should_buffer(0) {
self.parent_on_write(buf_len, WriteStatus::Drained);
Self::register_poll(self);
return WriteResult::Wrote(buf_len);
}
self.try_write_newly_buffered_data()
}
fn try_write_newly_buffered_data(&mut self) -> WriteResult {
debug_assert!(!self.is_done);
let rc = self.try_write(self.force_sync, self.outgoing.slice());
match rc {
WriteResult::Wrote(amt) => {
if amt == self.outgoing.size() {
self.outgoing.reset();
self.parent_on_write(amt, WriteStatus::Drained);
} else {
self.outgoing.wrote(amt);
self.parent_on_write(amt, WriteStatus::Pending);
Self::register_poll(self);
return WriteResult::Pending(amt);
}
}
WriteResult::Done(amt) => {
self.outgoing.reset();
self.parent_on_write(amt, WriteStatus::EndOfFile);
}
WriteResult::Pending(amt) => {
self.outgoing.wrote(amt);
self.parent_on_write(amt, WriteStatus::Pending);
Self::register_poll(self);
}
WriteResult::Err(e) => return WriteResult::Err(e),
}
rc
}
pub fn write(&mut self, buf: &[u8]) -> WriteResult {
if self.is_done || self.closed_without_reporting {
return WriteResult::Done(0);
}
if self.should_buffer(buf.len()) {
if self.outgoing.write(buf).is_err() {
return WriteResult::Err(oom_err());
}
self.parent_on_write(buf.len(), WriteStatus::Drained);
Self::register_poll(self);
return WriteResult::Wrote(buf.len());
}
if self.outgoing.size() > 0 {
if self.outgoing.write(buf).is_err() {
return WriteResult::Err(oom_err());
}
return self.try_write_newly_buffered_data();
}
let rc = self.try_write(self.force_sync, buf);
match rc {
WriteResult::Pending(amt) => {
if self.outgoing.write(&buf[amt..]).is_err() {
return WriteResult::Err(oom_err());
}
self.parent_on_write(amt, WriteStatus::Pending);
Self::register_poll(self);
}
WriteResult::Wrote(amt) => {
if amt < buf.len() {
if self.outgoing.write(&buf[amt..]).is_err() {
return WriteResult::Err(oom_err());
}
self.parent_on_write(amt, WriteStatus::Pending);
Self::register_poll(self);
} else {
self.outgoing.reset();
self.parent_on_write(amt, WriteStatus::Drained);
}
}
WriteResult::Done(amt) => {
self.outgoing.reset();
self.parent_on_write(amt, WriteStatus::EndOfFile);
return WriteResult::Done(amt);
}
_ => {}
}
rc
}
pub fn flush(&mut self) -> WriteResult {
if self.closed_without_reporting || self.is_done {
return WriteResult::Done(0);
}
let buffer_len = self.get_buffer().len();
if buffer_len == 0 {
self.outgoing.reset();
return WriteResult::Wrote(0);
}
let received_hup = 'brk: {
if let Some(poll) = self.get_poll() {
break 'brk poll.has_flag(FilePollFlag::Hup);
}
false
};
let rc = self.drain_buffered_data(usize::MAX, received_hup);
match rc {
WriteResult::Pending(written) => {
self.outgoing.wrote(written);
if self.outgoing.is_empty() {
self.outgoing.reset();
}
}
WriteResult::Wrote(written) => {
self.outgoing.wrote(written);
if self.outgoing.is_empty() {
self.outgoing.reset();
}
}
_ => {
self.outgoing.reset();
}
}
rc
}
pub fn has_ref(&self) -> bool {
let Some(poll) = self.get_poll() else {
return false;
};
!self.is_done && poll.can_enable_keeping_process_alive()
}
pub fn enable_keeping_process_alive(&self, event_loop: EventLoopHandle) {
if self.is_done {
return;
}
let Some(poll) = self.get_poll() else { return };
poll.enable_keeping_process_alive(event_loop);
}
pub fn disable_keeping_process_alive(&self, event_loop: EventLoopHandle) {
let Some(poll) = self.get_poll() else { return };
poll.disable_keeping_process_alive(event_loop);
}
pub fn update_ref(&self, event_loop: EventLoopHandle, value: bool) {
if value {
self.enable_keeping_process_alive(event_loop);
} else {
self.disable_keeping_process_alive(event_loop);
}
}
pub fn end(&mut self) {
if self.is_done {
return;
}
self.is_done = true;
self.close();
}
pub fn close(&mut self) {
if self.closed_without_reporting {
self.closed_without_reporting = false;
debug_assert!(self.get_fd() == Fd::INVALID);
unsafe { Parent::on_close(self.parent()) };
return;
}
let parent = self.parent;
self.handle.close(
Some(parent.cast()),
Some(|ctx: *mut c_void| unsafe { Parent::on_close(ctx.cast::<Parent>()) }),
);
}
pub fn start(&mut self, fd: Fd, is_pollable: bool) -> sys::Result<()> {
if !is_pollable {
self.close();
self.handle = PollOrFd::Fd(fd);
return sys::Result::Ok(());
}
let loop_ = unsafe { Parent::event_loop(self.parent()) };
let poll = match self.get_poll() {
Some(p) => p,
None => {
let p = FilePollRef::init(
loop_,
fd,
Owner::new(Parent::POLL_OWNER_TAG, std::ptr::from_mut(self).cast()),
);
self.handle = PollOrFd::Poll(p);
p
}
};
match poll.register_with_fd(loop_.loop_(), FilePollKind::Writable, fd) {
sys::Result::Err(err) => {
return sys::Result::Err(err);
}
sys::Result::Ok(()) => {}
}
sys::Result::Ok(())
}
}
impl<Parent: PosixStreamingWriterParent> Drop for PosixStreamingWriter<Parent> {
fn drop(&mut self) {
self.close_without_reporting();
}
}
#[cfg(windows)]
pub trait BaseWindowsPipeWriter {
type Parent: WindowsWriterParent;
const HAS_CURRENT_PAYLOAD: bool;
fn source(&self) -> &Option<Source>;
fn source_mut(&mut self) -> &mut Option<Source>;
fn parent_ptr(&self) -> *mut Self::Parent;
fn set_parent_ptr(&mut self, p: *mut Self::Parent);
fn is_done(&self) -> bool;
fn set_is_done(&mut self, v: bool);
fn owns_fd(&self) -> bool;
fn start_with_current_pipe(&mut self) -> sys::Result<()>;
fn on_close_source(&mut self);
fn get_fd(&self) -> Fd {
let Some(pipe) = self.source() else {
return Fd::INVALID;
};
pipe.get_fd()
}
fn has_ref(&self) -> bool {
if self.is_done() {
return false;
}
if let Some(pipe) = self.source() {
return pipe.has_ref();
}
false
}
fn enable_keeping_process_alive(&mut self, event_loop: EventLoopHandle) {
self.update_ref(event_loop, true);
}
fn disable_keeping_process_alive(&mut self, event_loop: EventLoopHandle) {
self.update_ref(event_loop, false);
}
fn close(&mut self) {
self.set_is_done(true);
let Some(source) = self.source_mut().take() else {
return;
};
let has_inflight_write = if Self::HAS_CURRENT_PAYLOAD {
match &source {
Source::SyncFile(file) | Source::File(file) => {
file.state == crate::source::FileState::Operating
|| file.state == crate::source::FileState::Canceling
}
_ => false,
}
} else {
false
};
match source {
Source::SyncFile(file) | Source::File(file) => {
let raw = bun_core::heap::into_raw(file);
unsafe {
if self.owns_fd() {
(*raw).detach();
} else {
(*raw).stop();
(*raw).fs.data = core::ptr::null_mut();
if (*raw).state == crate::source::FileState::Deinitialized {
drop(bun_core::heap::take(raw));
}
}
}
}
Source::Pipe(pipe) => {
let raw = bun_core::heap::into_raw(pipe);
unsafe {
(*raw).data = raw.cast::<c_void>();
(*raw).close(on_pipe_close);
}
}
Source::Tty(tty) => {
let p = tty.as_ptr();
unsafe { (*p).data = p.cast::<c_void>() };
unsafe { (*p).close(on_tty_close) };
}
}
*self.source_mut() = None;
self.on_close_source();
if has_inflight_write {
unsafe { Self::Parent::deref(self.parent_ptr()) };
}
}
fn update_ref(&mut self, _event_loop: EventLoopHandle, value: bool) {
if let Some(pipe) = self.source_mut().as_mut() {
if value {
pipe.ref_();
} else {
pipe.unref();
}
}
}
fn set_parent(&mut self, parent: *mut Self::Parent) {
self.set_parent_ptr(parent);
if !self.is_done() {
let self_ptr = core::ptr::from_mut(self).cast::<c_void>();
if let Some(pipe) = self.source_mut().as_mut() {
pipe.set_data(self_ptr);
}
}
}
fn watch(&mut self) {
}
unsafe fn start_with_pipe(&mut self, pipe: *mut uv::Pipe) -> sys::Result<()> {
debug_assert!(self.source().is_none());
*self.source_mut() = Some(Source::Pipe(unsafe { bun_core::heap::take(pipe) }));
let p = self.parent_ptr();
self.set_parent(p);
self.start_with_current_pipe()
}
fn start_sync(&mut self, fd: Fd, _pollable: bool) -> sys::Result<()> {
debug_assert!(self.source().is_none());
let mut source = Source::SyncFile(Source::open_file(fd));
source.set_data(core::ptr::from_mut(self).cast::<c_void>());
*self.source_mut() = Some(source);
let p = self.parent_ptr();
self.set_parent(p);
self.start_with_current_pipe()
}
fn start_with_file(&mut self, fd: Fd) -> sys::Result<()> {
debug_assert!(self.source().is_none());
let mut source = Source::File(Source::open_file(fd));
source.set_data(core::ptr::from_mut(self).cast::<c_void>());
*self.source_mut() = Some(source);
let p = self.parent_ptr();
self.set_parent(p);
self.start_with_current_pipe()
}
fn start(&mut self, rawfd: Fd, _pollable: bool) -> sys::Result<()> {
let fd = rawfd;
debug_assert!(self.source().is_none());
let loop_ = unsafe { Self::Parent::loop_(self.parent_ptr()) };
let mut source = match Source::open(loop_, fd) {
sys::Result::Ok(source) => source,
sys::Result::Err(err) => return sys::Result::Err(err),
};
let _ = matches!(source, Source::Pipe(_) | Source::Tty(_));
source.set_data(core::ptr::from_mut(self).cast::<c_void>());
*self.source_mut() = Some(source);
let p = self.parent_ptr();
self.set_parent(p);
self.start_with_current_pipe()
}
unsafe fn set_pipe(&mut self, pipe: *mut uv::Pipe) {
debug_assert!(self.source().is_none());
*self.source_mut() = Some(Source::Pipe(unsafe { bun_core::heap::take(pipe) }));
let p = self.parent_ptr();
self.set_parent(p);
}
fn get_stream(&mut self) -> Option<*mut uv::uv_stream_t> {
let source = self.source_mut().as_mut()?;
if matches!(source, Source::File(_) | Source::SyncFile(_)) {
return None;
}
Some(source.to_stream())
}
}
#[cfg(windows)]
extern "C" fn on_pipe_close(handle: *mut uv::Pipe) {
drop(unsafe { bun_core::heap::take(handle) });
}
#[cfg(windows)]
extern "C" fn on_tty_close(handle: *mut uv::uv_tty_t) {
if !crate::source::stdin_tty::is_stdin_tty(handle) {
drop(unsafe { bun_core::heap::take(handle) });
}
}
#[cfg(windows)]
pub trait WindowsWriterParent {
unsafe fn loop_(this: *mut Self) -> *mut uv::Loop;
unsafe fn ref_(this: *mut Self);
unsafe fn deref(this: *mut Self);
}
#[cfg(windows)]
pub trait WindowsBufferedWriterParent: WindowsWriterParent {
unsafe fn on_write(this: *mut Self, amount: usize, status: WriteStatus);
unsafe fn on_error(this: *mut Self, err: sys::Error);
const HAS_ON_CLOSE: bool;
unsafe fn on_close(_this: *mut Self) {}
unsafe fn get_buffer<'a>(this: *mut Self) -> &'a [u8];
const HAS_ON_WRITABLE: bool;
unsafe fn on_writable(_this: *mut Self) {}
}
#[cfg(windows)]
pub struct WindowsBufferedWriter<Parent: WindowsBufferedWriterParent> {
pub source: Option<Source>,
pub owns_fd: bool,
pub parent: *mut Parent,
pub is_done: bool,
pub write_req: uv::uv_write_t,
pub write_buffer: uv::uv_buf_t,
pub pending_payload_size: usize,
}
#[cfg(windows)]
impl<Parent: WindowsBufferedWriterParent> Default for WindowsBufferedWriter<Parent> {
fn default() -> Self {
Self {
source: None,
owns_fd: true,
parent: core::ptr::null_mut(), is_done: false,
write_req: bun_core::ffi::zeroed(),
write_buffer: uv::uv_buf_t::init(b""),
pending_payload_size: 0,
}
}
}
#[cfg(windows)]
impl<Parent: WindowsBufferedWriterParent> BaseWindowsPipeWriter for WindowsBufferedWriter<Parent> {
type Parent = Parent;
const HAS_CURRENT_PAYLOAD: bool = false;
fn source(&self) -> &Option<Source> {
&self.source
}
fn source_mut(&mut self) -> &mut Option<Source> {
&mut self.source
}
fn parent_ptr(&self) -> *mut Parent {
self.parent
}
fn set_parent_ptr(&mut self, p: *mut Parent) {
self.parent = p;
}
fn is_done(&self) -> bool {
self.is_done
}
fn set_is_done(&mut self, v: bool) {
self.is_done = v;
}
fn owns_fd(&self) -> bool {
self.owns_fd
}
fn on_close_source(&mut self) {
if Parent::HAS_ON_CLOSE {
unsafe { Parent::on_close(self.parent) };
}
}
fn start_with_current_pipe(&mut self) -> sys::Result<()> {
debug_assert!(self.source.is_some());
self.is_done = false;
self.write();
sys::Result::Ok(())
}
}
#[cfg(windows)]
unsafe impl<Parent: WindowsBufferedWriterParent> bun_ptr::LaunderedSelf
for WindowsBufferedWriter<Parent>
{
}
#[cfg(windows)]
impl<Parent: WindowsBufferedWriterParent> WindowsBufferedWriter<Parent> {
#[inline]
fn parent(&self) -> *mut Parent {
self.parent
}
#[inline]
fn parent_on_error(&self, err: sys::Error) {
unsafe { Parent::on_error(self.parent(), err) }
}
#[inline(always)]
fn r_on_error(this: *mut Self, err: sys::Error) {
let parent = Self::r(this).parent;
unsafe { Parent::on_error(parent, err) }
}
pub fn memory_cost(&self) -> usize {
mem::size_of::<Self>() + self.write_buffer.len as usize
}
fn on_write_complete(&mut self, status: uv::ReturnCode) {
let this: *mut Self = core::hint::black_box(core::ptr::from_mut(self));
let written = Self::r(this).pending_payload_size;
Self::r(this).pending_payload_size = 0;
if let Some(err) = status.to_error(sys::Tag::write) {
Self::r(this).close();
Self::r_on_error(this, err);
return;
}
let pending = Self::r(this).get_buffer_internal();
let has_pending_data = (pending.len() - written) != 0;
let is_done_before = Self::r(this).is_done;
unsafe {
Parent::on_write(
Self::r(this).parent(),
written,
if is_done_before && !has_pending_data {
WriteStatus::Drained
} else {
WriteStatus::Pending
},
)
};
core::hint::black_box(this);
if Self::r(this).is_done && !has_pending_data {
Self::r(this).close();
return;
}
if Parent::HAS_ON_WRITABLE {
unsafe { Parent::on_writable(Self::r(this).parent()) };
}
}
extern "C" fn on_fs_write_complete(fs: *mut uv::fs_t) {
let (file, result, parent_ptr) = unsafe { crate::source::File::from_fs_callback(fs) };
let was_canceled = result.int() == uv::UV_ECANCELED as i64;
file.complete(was_canceled);
if parent_ptr.is_null() {
if file.state == crate::source::FileState::Deinitialized {
drop(unsafe { bun_core::heap::take(core::ptr::from_mut(file)) });
}
return;
}
let this: *mut Self = core::hint::black_box(parent_ptr.cast::<Self>());
if was_canceled {
Self::r(this).pending_payload_size = 0;
return;
}
if let Some(err) = result.to_error(sys::Tag::write) {
Self::r(this).close();
core::hint::black_box(this);
Self::r_on_error(this, err);
return;
}
Self::r(this).on_write_complete(uv::ReturnCode::zero());
}
pub fn write(&mut self) {
let buffer = self.get_buffer_internal();
if self.is_done || self.pending_payload_size > 0 || buffer.len() == 0 {
return;
}
let buffer_len = buffer.len();
let write_buf = uv::uv_buf_t::init(buffer);
let (file_raw, stream_raw): (*mut crate::source::File, *mut uv::uv_stream_t) =
match self.source.as_mut() {
None => return,
Some(Source::SyncFile(_)) => {
panic!("This code path shouldn't be reached - sync_file in PipeWriter.zig");
}
Some(Source::File(f)) => (f.as_mut() as *mut _, core::ptr::null_mut()),
Some(s) => (core::ptr::null_mut(), s.to_stream()),
};
if !file_raw.is_null() {
let file = unsafe { &mut *file_raw };
debug_assert!(file.can_start());
self.pending_payload_size = buffer_len;
file.fs.data = core::ptr::from_mut(self).cast::<c_void>();
file.prepare();
self.write_buffer = write_buf;
if let Some(err) = unsafe {
uv::uv_fs_write(
Parent::loop_(self.parent()),
&mut file.fs,
file.file,
&self.write_buffer,
1,
-1,
Some(Self::on_fs_write_complete),
)
}
.to_error(sys::Tag::write)
{
file.complete(false);
self.close();
self.parent_on_error(err);
}
} else {
self.pending_payload_size = buffer_len;
self.write_buffer = write_buf;
let self_ptr = self as *mut Self;
if let Some(write_err) = self
.write_req
.write(stream_raw, &self.write_buffer, self_ptr, |p, s| unsafe {
(*p).on_write_complete(s)
})
.to_error(sys::Tag::write)
{
self.close();
self.parent_on_error(write_err);
}
}
}
fn get_buffer_internal(&self) -> &[u8] {
unsafe { Parent::get_buffer(self.parent()) }
}
pub fn end(&mut self) {
if self.is_done {
return;
}
self.is_done = true;
if self.pending_payload_size == 0 {
self.close();
}
}
}
#[derive(Default)]
pub struct StreamBuffer {
pub list: Vec<u8>,
pub cursor: usize,
}
impl StreamBuffer {
pub fn reset(&mut self) {
self.cursor = 0;
self.maybe_shrink();
self.list.clear();
}
pub fn maybe_shrink(&mut self) {
let page = 4096usize;
if self.list.capacity() > page {
self.list.truncate(page);
self.list.shrink_to(page);
}
}
pub fn memory_cost(&self) -> usize {
self.list.capacity()
}
pub fn size(&self) -> usize {
self.list.len() - self.cursor
}
pub fn is_empty(&self) -> bool {
self.size() == 0
}
pub fn is_not_empty(&self) -> bool {
self.size() > 0
}
pub fn write(&mut self, buffer: &[u8]) -> Result<(), OOM> {
self.list.extend_from_slice(buffer);
Ok(())
}
pub fn wrote(&mut self, amount: usize) {
self.cursor += amount;
}
pub fn write_assume_capacity(&mut self, buffer: &[u8]) {
self.list.extend_from_slice(buffer);
}
pub fn ensure_unused_capacity(&mut self, capacity: usize) -> Result<(), OOM> {
self.list.reserve(capacity);
Ok(())
}
pub fn write_type_as_bytes<T: bun_core::NoUninit>(&mut self, data: &T) -> Result<(), OOM> {
self.write(bun_core::bytes_of(data))
}
pub fn write_type_as_bytes_assume_capacity<T: bun_core::NoUninit>(&mut self, data: T) {
self.list.extend_from_slice(bun_core::bytes_of(&data));
}
pub fn write_or_fallback<'a>(
&'a mut self,
buffer_u8: Option<&'a [u8]>,
buffer_u16: Option<&[u16]>,
kind: WriteKind,
) -> Result<&'a [u8], OOM> {
match kind {
WriteKind::Latin1 => {
let buffer = buffer_u8.unwrap();
if bun_core::strings::is_all_ascii(buffer) {
return Ok(buffer);
}
self.write_latin1::<false>(buffer)?;
Ok(&self.list[self.cursor..])
}
WriteKind::Utf16 => {
let buffer = buffer_u16.unwrap();
self.write_utf16(buffer)?;
Ok(&self.list[self.cursor..])
}
WriteKind::Bytes => Ok(buffer_u8.unwrap()),
}
}
pub fn write_latin1<const CHECK_ASCII: bool>(&mut self, buffer: &[u8]) -> Result<(), OOM> {
if CHECK_ASCII {
if bun_core::strings::is_all_ascii(buffer) {
return self.write(buffer);
}
}
let len = self.list.len();
let list = mem::take(&mut self.list);
self.list = bun_core::strings::allocate_latin1_into_utf8_with_list(list, len, buffer);
Ok(())
}
pub fn write_utf16(&mut self, buffer: &[u16]) -> Result<(), OOM> {
ByteVecExt::write_utf16(&mut self.list, buffer)?;
Ok(())
}
pub fn slice(&self) -> &[u8] {
&self.list[self.cursor..]
}
}
#[derive(Copy, Clone, Eq, PartialEq)]
pub enum WriteKind {
Bytes,
Latin1,
Utf16,
}
#[cfg(windows)]
pub trait WindowsStreamingWriterParent: WindowsWriterParent {
unsafe fn on_write(this: *mut Self, amount: usize, status: WriteStatus);
unsafe fn on_error(this: *mut Self, err: sys::Error);
const HAS_ON_WRITABLE: bool;
unsafe fn on_writable(_this: *mut Self) {}
unsafe fn on_close(this: *mut Self);
}
#[cfg(windows)]
pub struct WindowsStreamingWriter<Parent: WindowsStreamingWriterParent> {
pub source: Option<Source>,
pub owns_fd: bool,
pub parent: *mut Parent,
pub is_done: bool,
#[allow(dead_code)]
pub write_req: uv::uv_write_t,
#[allow(dead_code)]
pub write_buffer: uv::uv_buf_t,
#[allow(dead_code)]
pub outgoing: StreamBuffer,
#[allow(dead_code)]
pub current_payload: StreamBuffer,
#[allow(dead_code)]
pub last_write_result: WriteResult,
pub closed_without_reporting: bool,
}
#[cfg(windows)]
impl<Parent: WindowsStreamingWriterParent> Default for WindowsStreamingWriter<Parent> {
fn default() -> Self {
Self {
source: None,
owns_fd: true,
parent: core::ptr::null_mut(), is_done: false,
write_req: bun_core::ffi::zeroed(),
write_buffer: uv::uv_buf_t::init(b""),
outgoing: StreamBuffer::default(),
current_payload: StreamBuffer::default(),
last_write_result: WriteResult::Wrote(0),
closed_without_reporting: false,
}
}
}
#[cfg(windows)]
impl<Parent: WindowsStreamingWriterParent> BaseWindowsPipeWriter
for WindowsStreamingWriter<Parent>
{
type Parent = Parent;
const HAS_CURRENT_PAYLOAD: bool = true;
fn source(&self) -> &Option<Source> {
&self.source
}
fn source_mut(&mut self) -> &mut Option<Source> {
&mut self.source
}
fn parent_ptr(&self) -> *mut Parent {
self.parent
}
fn set_parent_ptr(&mut self, p: *mut Parent) {
self.parent = p;
}
fn is_done(&self) -> bool {
self.is_done
}
fn set_is_done(&mut self, v: bool) {
self.is_done = v;
}
fn owns_fd(&self) -> bool {
self.owns_fd
}
fn on_close_source(&mut self) {
self.source = None;
if self.closed_without_reporting {
self.closed_without_reporting = false;
return;
}
unsafe { Parent::on_close(self.parent) };
}
fn start_with_current_pipe(&mut self) -> sys::Result<()> {
debug_assert!(self.source.is_some());
self.is_done = false;
sys::Result::Ok(())
}
}
#[cfg(windows)]
unsafe impl<Parent: WindowsStreamingWriterParent> bun_ptr::LaunderedSelf
for WindowsStreamingWriter<Parent>
{
}
#[cfg(windows)]
impl<Parent: WindowsStreamingWriterParent> WindowsStreamingWriter<Parent> {
#[inline]
#[allow(dead_code)]
fn parent(&self) -> *mut Parent {
self.parent
}
#[inline(always)]
#[allow(dead_code)]
fn r_on_error(this: *mut Self, err: sys::Error) {
let parent = Self::r(this).parent;
unsafe { Parent::on_error(parent, err) }
}
#[inline(always)]
#[allow(dead_code)]
fn r_on_write(this: *mut Self, written: usize, status: WriteStatus) {
let parent = Self::r(this).parent;
unsafe { Parent::on_write(parent, written, status) }
}
#[inline(always)]
#[allow(dead_code)]
fn r_deref(this: *mut Self) {
let parent = Self::r(this).parent;
unsafe { Parent::deref(parent) }
}
#[allow(dead_code)]
pub fn memory_cost(&self) -> usize {
mem::size_of::<Self>() + self.current_payload.memory_cost() + self.outgoing.memory_cost()
}
#[allow(dead_code)]
pub fn has_pending_data(&self) -> bool {
self.outgoing.is_not_empty() || self.current_payload.is_not_empty()
}
#[allow(dead_code)]
fn on_write_complete(&mut self, status: uv::ReturnCode) {
let this: *mut Self = core::hint::black_box(core::ptr::from_mut(self));
let _g = scopeguard::guard(this, |s| Self::r_deref(s));
if let Some(err) = status.to_error(sys::Tag::write) {
log!("onWrite() = {}", bstr::BStr::new(err.name()));
Self::r(this).last_write_result = WriteResult::Err(err.clone());
Self::r_on_error(this, err);
core::hint::black_box(this);
Self::r(this).close_without_reporting();
return;
}
let written = Self::r(this).current_payload.size();
Self::r(this).current_payload.reset();
let done = Self::r(this).outgoing.is_empty();
let was_done = Self::r(this).is_done;
log!(
"onWrite({}) ({} left)",
written,
Self::r(this).outgoing.size()
);
if was_done && done {
Self::r(this).last_write_result = WriteResult::Done(written);
Self::r_on_write(this, written, WriteStatus::EndOfFile);
return;
}
Self::r(this).last_write_result = WriteResult::Wrote(written);
Self::r_on_write(
this,
written,
if done {
WriteStatus::Drained
} else {
WriteStatus::Pending
},
);
core::hint::black_box(this);
Self::r(this).process_send();
if Parent::HAS_ON_WRITABLE {
unsafe { Parent::on_writable(Self::r(this).parent()) };
}
}
extern "C" fn on_fs_write_complete(fs: *mut uv::fs_t) {
let (file, result, parent_ptr) = unsafe { crate::source::File::from_fs_callback(fs) };
let was_canceled = result.int() == uv::UV_ECANCELED as i64;
file.complete(was_canceled);
if parent_ptr.is_null() {
if file.state == crate::source::FileState::Deinitialized {
drop(unsafe { bun_core::heap::take(core::ptr::from_mut(file)) });
}
return;
}
let this: *mut Self = core::hint::black_box(parent_ptr.cast::<Self>());
if was_canceled {
Self::r(this).current_payload.reset();
Self::r_deref(this);
return;
}
if let Some(err) = result.to_error(sys::Tag::write) {
let _g = scopeguard::guard(this, |s| Self::r_deref(s));
Self::r(this).close();
core::hint::black_box(this);
Self::r_on_error(this, err);
return;
}
Self::r(this).on_write_complete(uv::ReturnCode::zero());
}
fn process_send(&mut self) {
log!("processSend");
let this: *mut Self = core::hint::black_box(core::ptr::from_mut(self));
if Self::r(this).current_payload.is_not_empty() {
Self::r(this).last_write_result = WriteResult::Pending(0);
return;
}
let bytes_len = Self::r(this).outgoing.slice().len();
if bytes_len == 0 {
Self::r(this).last_write_result = WriteResult::Wrote(0);
return;
}
let (file_raw, stream_raw): (*mut crate::source::File, *mut uv::uv_stream_t) =
match Self::r(this).source.as_mut() {
None => {
let err = sys::Error::from_code(sys::E::PIPE, sys::Tag::pipe);
Self::r(this).last_write_result = WriteResult::Err(err.clone());
Self::r_on_error(this, err);
core::hint::black_box(this);
Self::r(this).close_without_reporting();
return;
}
Some(Source::SyncFile(_)) => {
panic!("sync_file pipe write should not be reachable");
}
Some(Source::File(f)) => (f.as_mut() as *mut _, core::ptr::null_mut()),
Some(s) => (core::ptr::null_mut(), s.to_stream()),
};
{
let s = Self::r(this);
mem::swap(&mut s.current_payload, &mut s.outgoing);
}
let write_buf = {
let s = Self::r(this);
debug_assert_eq!(s.current_payload.slice().len(), bytes_len);
uv::uv_buf_t::init(s.current_payload.slice())
};
if !file_raw.is_null() {
let file = unsafe { &mut *file_raw };
debug_assert!(file.can_start());
file.fs.data = this.cast::<c_void>();
file.prepare();
Self::r(this).write_buffer = write_buf;
if let Some(err) = unsafe {
uv::uv_fs_write(
Parent::loop_((*this).parent()),
&mut file.fs,
file.file,
&(*this).write_buffer,
1,
-1,
Some(Self::on_fs_write_complete),
)
}
.to_error(sys::Tag::write)
{
file.complete(false);
Self::r(this).last_write_result = WriteResult::Err(err.clone());
Self::r_on_error(this, err);
core::hint::black_box(this);
Self::r(this).close_without_reporting();
return;
}
} else {
Self::r(this).write_buffer = write_buf;
if let Some(err) = unsafe {
(*this)
.write_req
.write(stream_raw, &(*this).write_buffer, this, |p, s| {
(*p).on_write_complete(s)
})
}
.to_error(sys::Tag::write)
{
Self::r(this).last_write_result = WriteResult::Err(err.clone());
Self::r_on_error(this, err);
core::hint::black_box(this);
Self::r(this).close_without_reporting();
return;
}
}
unsafe { Parent::ref_(Self::r(this).parent()) };
Self::r(this).last_write_result = WriteResult::Pending(0);
}
fn close_without_reporting(&mut self) {
if self.get_fd() != Fd::INVALID {
debug_assert!(!self.closed_without_reporting);
self.closed_without_reporting = true;
self.close();
}
}
#[allow(dead_code)]
fn write_internal_u8(&mut self, buffer: &[u8], kind: WriteKind) -> WriteResult {
if self.is_done {
return WriteResult::Done(0);
}
if matches!(self.source, Some(Source::SyncFile(_))) {
let result = (|| {
let remain = match self.outgoing.write_or_fallback(Some(buffer), None, kind) {
Ok(r) => r,
Err(_) => return WriteResult::Err(oom_err()),
};
let initial_len = remain.len();
let mut remain = remain;
let fd = Fd::from_uv(match &self.source {
Some(Source::SyncFile(f)) => f.file,
_ => unreachable!(),
});
while remain.len() > 0 {
match sys::write(fd, remain) {
sys::Result::Err(err) => return WriteResult::Err(err),
sys::Result::Ok(wrote) => {
remain = &remain[wrote..];
if wrote == 0 {
break;
}
}
}
}
let wrote = initial_len - remain.len();
if wrote == 0 {
return WriteResult::Done(wrote);
}
WriteResult::Wrote(wrote)
})();
self.outgoing.reset();
return result;
}
let had_buffered_data = self.outgoing.is_not_empty();
let r = match kind {
WriteKind::Latin1 => self.outgoing.write_latin1::<true>(buffer),
WriteKind::Bytes => self.outgoing.write(buffer),
WriteKind::Utf16 => unreachable!(),
};
if r.is_err() {
return WriteResult::Err(oom_err());
}
if had_buffered_data {
return WriteResult::Pending(0);
}
self.process_send();
self.last_write_result.clone()
}
#[allow(dead_code)]
fn write_internal_u16(&mut self, buffer: &[u16]) -> WriteResult {
if self.is_done {
return WriteResult::Done(0);
}
if matches!(self.source, Some(Source::SyncFile(_))) {
let result = (|| {
let remain =
match self
.outgoing
.write_or_fallback(None, Some(buffer), WriteKind::Utf16)
{
Ok(r) => r,
Err(_) => return WriteResult::Err(oom_err()),
};
let initial_len = remain.len();
let mut remain = remain;
let fd = Fd::from_uv(match &self.source {
Some(Source::SyncFile(f)) => f.file,
_ => unreachable!(),
});
while remain.len() > 0 {
match sys::write(fd, remain) {
sys::Result::Err(err) => return WriteResult::Err(err),
sys::Result::Ok(wrote) => {
remain = &remain[wrote..];
if wrote == 0 {
break;
}
}
}
}
let wrote = initial_len - remain.len();
if wrote == 0 {
return WriteResult::Done(wrote);
}
WriteResult::Wrote(wrote)
})();
self.outgoing.reset();
return result;
}
let had_buffered_data = self.outgoing.is_not_empty();
if self.outgoing.write_utf16(buffer).is_err() {
return WriteResult::Err(oom_err());
}
if had_buffered_data {
return WriteResult::Pending(0);
}
self.process_send();
self.last_write_result.clone()
}
#[allow(dead_code)]
pub fn write_utf16(&mut self, buf: &[u16]) -> WriteResult {
self.write_internal_u16(buf)
}
#[allow(dead_code)]
pub fn write_latin1(&mut self, buffer: &[u8]) -> WriteResult {
self.write_internal_u8(buffer, WriteKind::Latin1)
}
#[allow(dead_code)]
pub fn write(&mut self, buffer: &[u8]) -> WriteResult {
self.write_internal_u8(buffer, WriteKind::Bytes)
}
#[allow(dead_code)]
pub fn flush(&mut self) -> WriteResult {
if self.is_done {
return WriteResult::Done(0);
}
if !self.has_pending_data() {
return WriteResult::Wrote(0);
}
self.process_send();
self.last_write_result.clone()
}
#[allow(dead_code)]
pub fn end(&mut self) {
if self.is_done {
return;
}
self.closed_without_reporting = false;
self.is_done = true;
if !self.has_pending_data() {
if !self.owns_fd {
return;
}
self.close();
}
}
}
#[cfg(windows)]
impl<Parent: WindowsStreamingWriterParent> Drop for WindowsStreamingWriter<Parent> {
fn drop(&mut self) {
self.close_without_reporting();
}
}
#[cfg(unix)]
pub type BufferedWriter<P> = PosixBufferedWriter<P>;
#[cfg(not(unix))]
pub type BufferedWriter<P> = WindowsBufferedWriter<P>;
#[cfg(unix)]
pub type StreamingWriter<P> = PosixStreamingWriter<P>;
#[cfg(not(unix))]
pub type StreamingWriter<P> = WindowsStreamingWriter<P>;
#[doc(hidden)]
pub mod __parent_macro {
pub use ::bun_sys::Error as SysError;
#[cfg(windows)]
pub use ::bun_sys::windows::libuv::Loop as UvLoop;
pub use ::bun_uws_sys::Loop as UwsLoop;
}
#[macro_export]
macro_rules! impl_streaming_writer_parent {
(@call mut $p:expr; $m:ident($($a:tt)*)) => { (&mut *$p).$m($($a)*) };
(@call shared $p:expr; $m:ident($($a:tt)*)) => { (&*$p).$m($($a)*) };
(@call ptr $p:expr; $m:ident($($a:tt)*)) => { <Self>::$m($p, $($a)*) };
(@emit
[$($gen:tt)*] $Ty:ty;
poll_tag = $poll_tag:expr,
borrow = $borrow:tt,
on_write = $on_write:ident,
on_error = $on_error:ident,
on_ready = $on_ready:ident,
on_close = $on_close:ident,
event_loop = |$el_this:ident| $el:expr,
uws_loop = |$uws_this:ident| $uws:expr,
uv_loop = |$uv_this:ident| $uv:expr,
ref_ = |$ref_this:ident| $ref_:expr,
deref = |$deref_this:ident| $deref:expr,
) => {
#[cfg(unix)]
impl $($gen)* $crate::pipe_writer::PosixStreamingWriterParent for $Ty {
const POLL_OWNER_TAG: $crate::PollTag = $poll_tag;
const HAS_ON_READY: bool = true;
#[inline]
unsafe fn on_write(this: *mut Self, amount: usize, status: $crate::WriteStatus) {
unsafe { $crate::impl_streaming_writer_parent!(@call $borrow this; $on_write(amount, status)) }
}
#[inline]
unsafe fn on_error(this: *mut Self, err: $crate::pipe_writer::__parent_macro::SysError) {
unsafe { $crate::impl_streaming_writer_parent!(@call $borrow this; $on_error(err)) }
}
#[inline]
unsafe fn on_ready(this: *mut Self) {
unsafe { $crate::impl_streaming_writer_parent!(@call $borrow this; $on_ready()) }
}
#[inline]
unsafe fn on_close(this: *mut Self) {
unsafe { $crate::impl_streaming_writer_parent!(@call $borrow this; $on_close()) }
}
#[inline]
unsafe fn event_loop(this: *mut Self) -> $crate::EventLoopHandle {
let $el_this = this;
#[allow(unused_unsafe)]
unsafe { $el }
}
#[inline]
unsafe fn loop_(this: *mut Self) -> *mut $crate::pipe_writer::__parent_macro::UwsLoop {
let $uws_this = this;
#[allow(unused_unsafe)]
unsafe { $uws }
}
}
#[cfg(windows)]
impl $($gen)* $crate::pipe_writer::WindowsWriterParent for $Ty {
#[inline]
unsafe fn loop_(this: *mut Self) -> *mut $crate::pipe_writer::__parent_macro::UvLoop {
let $uv_this = this;
#[allow(unused_unsafe)]
unsafe { $uv }
}
#[inline]
unsafe fn ref_(this: *mut Self) {
let $ref_this = this;
#[allow(unused_unsafe)]
unsafe { $ref_ };
}
#[inline]
unsafe fn deref(this: *mut Self) {
let $deref_this = this;
#[allow(unused_unsafe)]
unsafe { $deref };
}
}
#[cfg(windows)]
impl $($gen)* $crate::pipe_writer::WindowsStreamingWriterParent for $Ty {
const HAS_ON_WRITABLE: bool = true;
#[inline]
unsafe fn on_write(this: *mut Self, amount: usize, status: $crate::WriteStatus) {
unsafe { $crate::impl_streaming_writer_parent!(@call $borrow this; $on_write(amount, status)) }
}
#[inline]
unsafe fn on_error(this: *mut Self, err: $crate::pipe_writer::__parent_macro::SysError) {
unsafe { $crate::impl_streaming_writer_parent!(@call $borrow this; $on_error(err)) }
}
#[inline]
unsafe fn on_writable(this: *mut Self) {
unsafe { $crate::impl_streaming_writer_parent!(@call $borrow this; $on_ready()) }
}
#[inline]
unsafe fn on_close(this: *mut Self) {
unsafe { $crate::impl_streaming_writer_parent!(@call $borrow this; $on_close()) }
}
}
};
(
for<$($gp:ident $(: $b0:path)?),+> $Ty:ty;
$($rest:tt)*
) => {
$crate::impl_streaming_writer_parent! {
@emit [<$($gp $(: $b0)?),+>] $Ty; $($rest)*
}
};
(
$Ty:ty;
$($rest:tt)*
) => {
$crate::impl_streaming_writer_parent! { @emit [] $Ty; $($rest)* }
};
}
#[macro_export]
macro_rules! impl_buffered_writer_parent {
(@borrow mut $p:expr) => { &mut *$p };
(@borrow shared $p:expr) => { &*$p };
(@emit
[$($gen:tt)*] $Ty:ty;
poll_tag = $poll_tag:expr,
borrow = $borrow:tt,
on_write = $on_write:ident,
on_error = $on_error:ident,
on_close = $on_close:ident,
get_buffer = |$gb_this:ident| $gb:expr,
event_loop = |$el_this:ident| $el:expr,
uv_loop = |$uv_this:ident| $uv:expr,
ref_ = |$ref_this:ident| $ref_:expr,
deref = |$deref_this:ident| $deref:expr,
win_on_write_guard = |$guard_this:ident| $guard:expr,
) => {
#[cfg(not(windows))]
impl $($gen)* $crate::pipe_writer::PosixBufferedWriterParent for $Ty {
const POLL_OWNER_TAG: $crate::PollTag = $poll_tag;
#[inline]
unsafe fn on_write(this: *mut Self, amount: usize, status: $crate::WriteStatus) {
unsafe { ($crate::impl_buffered_writer_parent!(@borrow $borrow this)).$on_write(amount, status) };
}
#[inline]
unsafe fn on_error(this: *mut Self, err: $crate::pipe_writer::__parent_macro::SysError) {
unsafe { ($crate::impl_buffered_writer_parent!(@borrow $borrow this)).$on_error(&err) };
}
const HAS_ON_CLOSE: bool = true;
#[inline]
unsafe fn on_close(this: *mut Self) {
unsafe { ($crate::impl_buffered_writer_parent!(@borrow $borrow this)).$on_close() };
}
#[inline]
unsafe fn get_buffer<'a>(this: *mut Self) -> &'a [u8] {
let $gb_this = this;
#[allow(unused_unsafe)]
unsafe { $gb }
}
const HAS_ON_WRITABLE: bool = false;
#[inline]
unsafe fn event_loop(this: *mut Self) -> $crate::EventLoopHandle {
let $el_this = this;
#[allow(unused_unsafe)]
unsafe { $el }
}
}
#[cfg(windows)]
impl $($gen)* $crate::pipe_writer::WindowsWriterParent for $Ty {
#[inline]
unsafe fn loop_(this: *mut Self) -> *mut $crate::pipe_writer::__parent_macro::UvLoop {
let $uv_this = this;
#[allow(unused_unsafe)]
unsafe { $uv }
}
#[inline]
unsafe fn ref_(this: *mut Self) {
let $ref_this = this;
#[allow(unused_unsafe)]
unsafe { $ref_ };
}
#[inline]
unsafe fn deref(this: *mut Self) {
let $deref_this = this;
#[allow(unused_unsafe)]
unsafe { $deref };
}
}
#[cfg(windows)]
impl $($gen)* $crate::pipe_writer::WindowsBufferedWriterParent for $Ty {
#[inline]
unsafe fn on_write(this: *mut Self, amount: usize, status: $crate::WriteStatus) {
let $guard_this = this;
#[allow(unused_unsafe, clippy::let_unit_value)]
let _guard = unsafe { $guard };
unsafe { ($crate::impl_buffered_writer_parent!(@borrow $borrow this)).$on_write(amount, status) };
}
#[inline]
unsafe fn on_error(this: *mut Self, err: $crate::pipe_writer::__parent_macro::SysError) {
unsafe { ($crate::impl_buffered_writer_parent!(@borrow $borrow this)).$on_error(&err) };
}
const HAS_ON_CLOSE: bool = true;
#[inline]
unsafe fn on_close(this: *mut Self) {
unsafe { ($crate::impl_buffered_writer_parent!(@borrow $borrow this)).$on_close() };
}
#[inline]
unsafe fn get_buffer<'a>(this: *mut Self) -> &'a [u8] {
let $gb_this = this;
#[allow(unused_unsafe)]
unsafe { $gb }
}
const HAS_ON_WRITABLE: bool = false;
}
};
(
for<$($gp:ident $(: $b0:path)?),+> $Ty:ty;
$($rest:tt)*
) => {
$crate::impl_buffered_writer_parent! {
@emit [<$($gp $(: $b0)?),+>] $Ty; $($rest)*
}
};
(
$Ty:ty;
$($rest:tt)*
) => {
$crate::impl_buffered_writer_parent! { @emit [] $Ty; $($rest)* }
};
}