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
4 changes: 3 additions & 1 deletion docs/runtime/file-io.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -192,7 +192,7 @@ To flush the buffer and close the file:
writer.end();
```

By default, the `bun` process stays alive until this `FileSink` is explicitly closed with `.end()`. To opt out of this behavior, "unref" the instance.
When the destination is a FIFO or pipe, the `bun` process stays alive while a write to it is still in flight, so buffered data isn't silently dropped on exit. To opt out of that, "unref" the instance — the process is then free to exit with the write unfinished.

```ts
writer.unref();
Expand All @@ -201,6 +201,8 @@ writer.unref();
writer.ref();
```

For regular files, `.ref()` and `.unref()` do nothing: those writes never hold the event loop open. A `FileSink` on its own does not keep `bun` running indefinitely; if you need the process to stay alive while you wait for something to write, keep that source of work (a server, a socket, a timer) referenced.

---

## Directories
Expand Down
18 changes: 11 additions & 7 deletions packages/bun-types/s3.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -42,13 +42,13 @@ declare module "bun" {
}): void;

/**
* For FIFOs & pipes, keep Bun's process alive until the pipe closes.
* For FIFOs & pipes, restore the automatic keep-alive management that
* {@link unref} turned off.
*
* By default, this is managed automatically: while the stream is open the
* process stays alive, and once the other end hangs up or the stream
* closes, the process exits.
*
* If you previously called {@link unref}, call this to re-enable automatic management.
* By default keep-alive is managed automatically: Bun's process stays alive
* while a write to the pipe is still in flight, and stops holding the event
* loop open once that write drains. This does not pin an idle stream —
* calling it is only useful after {@link unref}.
*
* Internally, calls to {@link ref} and {@link unref} are reference counted. The count starts at 1.
*
Expand All @@ -58,7 +58,11 @@ declare module "bun" {
ref(): void;

/**
* For FIFOs & pipes, allow Bun's process to exit while the stream is open.
* For FIFOs & pipes, allow Bun's process to exit while the stream is open,
* even if a write to it has not drained yet.
*
* This survives later writes: automatic keep-alive management stays off
* until {@link ref} is called.
*
* If the file is not a FIFO or pipe, {@link ref} and {@link unref} do
* nothing. If the pipe is already closed, this does nothing.
Expand Down
43 changes: 43 additions & 0 deletions src/io/keep_alive.rs
Original file line number Diff line number Diff line change
Expand Up @@ -164,3 +164,46 @@ impl KeepAlive {
self.unref_concurrently(loop_);
}
}

/// Gate for a keep-alive that is both automatically managed ("a write is in
/// flight") and exposed to JS as `.ref()`/`.unref()`.
///
/// `allowed` is what the user asked for (starts `true`; `.unref()` clears it,
/// `.ref()` restores it). `wanted` is what the automatic management asks for
/// right now. The owner holds the actual handle (a `FilePoll` or libuv handle)
/// and applies [`is_active`](Self::is_active) to it after each set call, so
/// `.unref()` survives the automatic path re-asserting `wanted` and `.ref()`
/// restores it without pinning an idle handle.
#[derive(Default)]
pub struct UserKeepAlive {
allowed: bool,
wanted: bool,
}

impl UserKeepAlive {
pub fn init() -> Self {
Self {
allowed: true,
wanted: false,
}
}

/// JS-facing `.ref()`/`.unref()` (the 0↔1 crossing). Returns
/// [`is_active`](Self::is_active) for the owner to apply.
pub fn set_allowed(&mut self, value: bool) -> bool {
self.allowed = value;
self.is_active()
}

/// Automatic management path. Returns [`is_active`](Self::is_active) for
/// the owner to apply.
pub fn set_wanted(&mut self, value: bool) -> bool {
self.wanted = value;
self.is_active()
}

#[inline]
pub fn is_active(&self) -> bool {
self.allowed && self.wanted
}
}
2 changes: 1 addition & 1 deletion src/io/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ pub mod windows_event_loop;
// `#[cfg(unix)]`-gated so the module still compiles on Windows.
mod keep_alive;
pub mod posix_event_loop;
pub use keep_alive::KeepAlive;
pub use keep_alive::{KeepAlive, UserKeepAlive};

// ParentDeathWatchdog is POSIX-only (uses `libc::pid_t`, `getppid`, signals);
// Windows handles orphan death via Job Objects in `spawn`. Downstream code
Expand Down
64 changes: 36 additions & 28 deletions src/runtime/webcore/FileSink.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,10 @@ pub struct FileSink {
pub started: Cell<bool>,
pub must_be_kept_alive_until_eof: Cell<bool>,

/// Keep-alive handle. `allowed` mirrors `JSFileSink::m_refCount > 0`
/// (`unref()`/`ref()` from JS), `wanted` is set while a write is in flight.
pub poll_ref: JsCell<bun_io::UserKeepAlive>,

// TODO: these fields are duplicated on writer()
// we should not duplicate these fields...
pub pollable: Cell<bool>,
Expand Down Expand Up @@ -401,12 +405,7 @@ impl FileSink {
let has_pending_data = (*this).writer.get().has_pending_data();
// Only keep the event loop ref'd while there's a pending write in progress.
// If there's no pending write, no need to keep the event loop ref'd.
// `with_mut`: Windows `update_ref` is `&mut self` (posix is `&self`).
// Hoist `io_evtloop()` out of the closure so no raw deref appears inside it.
let evtloop = (*this).io_evtloop();
(*this)
.writer
.with_mut(|w| w.update_ref(evtloop, has_pending_data));
(*this).set_keep_alive(has_pending_data);

if has_pending_data {
if let Some(vm) = (*this).js_vm() {
Expand Down Expand Up @@ -657,8 +656,7 @@ impl FileSink {
return sys::Result::Err(err);
}
sys::Result::Ok(()) => {
self.writer
.with_mut(|w| w.update_ref(self.io_evtloop(), false));
self.set_keep_alive(false);
}
}
return sys::Result::Ok(());
Expand All @@ -672,10 +670,10 @@ impl FileSink {
return sys::Result::Err(err);
}
sys::Result::Ok(()) => {
// Only keep the event loop ref'd while there's a pending write in progress.
// If there's no pending write, no need to keep the event loop ref'd.
self.writer
.with_mut(|w| w.update_ref(self.io_evtloop(), false));
// Only keep the event loop ref'd while there's a pending write
// in progress. If there's no pending write, no need to keep
// the event loop ref'd.
self.set_keep_alive(false);
#[cfg(unix)]
{
if self.nonblocking.get() {
Expand Down Expand Up @@ -821,7 +819,7 @@ impl FileSink {
// SAFETY: caller contract — `this` is live with write+dealloc provenance.
unsafe {
if (*this).done.get() || !(*this).writer.get().has_pending_data() {
(*this).update_ref(false);
(*this).set_keep_alive(false);
(*this).auto_flusher.with_mut(|a| a.registered.set(false));
return false;
}
Expand All @@ -834,12 +832,12 @@ impl FileSink {
// callback it may trigger goes via the stored `*mut FileSink` backref.
match (*this).writer.with_mut(|w| w.flush()) {
WriteResult::Err(_) | WriteResult::Done(_) => {
(*this).update_ref(false);
(*this).set_keep_alive(false);
(*this).run_pending_later();
}
WriteResult::Wrote(amount_drained) => {
if amount_drained == amount_buffered {
(*this).update_ref(false);
(*this).set_keep_alive(false);
(*this).run_pending_later();
}
}
Expand Down Expand Up @@ -1088,7 +1086,7 @@ impl FileSink {

match flush_result {
WriteResult::Done(written) => {
self.update_ref(false);
self.set_keep_alive(false);
self.writer.with_mut(|w| w.end());
sys::Result::Ok(JSValue::js_number(written as f64))
}
Expand Down Expand Up @@ -1132,17 +1130,27 @@ impl FileSink {
crate::webcore::sink::Sink::init(self)
}

/// `JSFileSink::ref()` / `unref()`, called when their `m_refCount` crosses
/// 0↔1. `ref()` restores the automatic management, it does not pin.
pub fn update_ref(&self, value: bool) {
// `with_mut`: the Windows `BaseWindowsPipeWriter` impls take `&mut self`
// (the posix `PosixStreamingWriter` impls are `&self`); `with_mut`
// covers both. No JS re-entry — pure libuv ref/unref.
self.writer.with_mut(|w| {
if value {
w.enable_keeping_process_alive(self.io_evtloop());
} else {
w.disable_keeping_process_alive(self.io_evtloop());
}
});
let active = self.poll_ref.with_mut(|k| k.set_allowed(value));
self.apply_keep_alive(active);
}

/// Automatic management path (`on_write`, `on_auto_flush`, `end_from_js`,
/// `assign_to_stream`). `UserKeepAlive` ANDs this with what JS
/// `ref()`/`unref()` allows.
fn set_keep_alive(&self, wants: bool) {
let active = self.poll_ref.with_mut(|k| k.set_wanted(wants));
self.apply_keep_alive(active);
}

fn apply_keep_alive(&self, active: bool) {
// `with_mut`: the Windows `BaseWindowsPipeWriter` impls take `&mut
// self` (the posix `PosixStreamingWriter` impls are `&self`); it
// covers both. No JS re-entry — pure libuv/poll ref/unref.
let evtloop = self.io_evtloop();
self.writer.with_mut(|w| w.update_ref(evtloop, active));
}
}

Expand Down Expand Up @@ -1302,6 +1310,7 @@ impl FileSink {
done: Cell::new(false),
started: Cell::new(false),
must_be_kept_alive_until_eof: Cell::new(false),
poll_ref: JsCell::new(bun_io::UserKeepAlive::init()),
pollable: Cell::new(false),
nonblocking: Cell::new(false),
force_sync: Cell::new(false),
Expand Down Expand Up @@ -1444,8 +1453,7 @@ impl FileSink {
// SAFETY: `as_any_promise` returned non-null.
match unsafe { (*js_promise).status() } {
bun_jsc::js_promise::Status::Pending => {
self.writer
.with_mut(|w| w.enable_keeping_process_alive(self.io_evtloop()));
self.set_keep_alive(true);
self.ref_();
// TODO: properly propagate exception upwards
// `JSValue::then` takes already-wrapped C-ABI
Expand Down
119 changes: 119 additions & 0 deletions test/js/bun/util/filesink.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -433,3 +433,122 @@ it("fs.promises.writeFile with iterables under GC pressure does not crash", asyn
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
expect({ stdout: stdout.trim(), stderr, exitCode }).toEqual({ stdout: "ok", stderr: "", exitCode: 0 });
});

describe.skipIf(!isPosix)("ref/unref keep-alive", () => {
// A FIFO whose read end nobody drains: the child's write stays pending
// forever, so the process can only exit if unref() actually took effect.
function blockedFifo(label: string) {
const fifo = join(tmpdirSync(), `${label}.fifo`);
mkfifo(fifo, 0o666);
return { fifo, readFd: fs.openSync(fifo, fs.constants.O_RDONLY | fs.constants.O_NONBLOCK) };
}

it("unref() is not re-armed by a later pending write", async () => {
const { fifo, readFd } = blockedFifo("unref-sticky");
try {
await using proc = Bun.spawn({
cmd: [
bunExe(),
"-e",
`
const sink = Bun.file(${JSON.stringify(fifo)}).writer();
sink.unref();
sink.write(Buffer.alloc(4 * 1024 * 1024, 0x61));
process.on("exit", () => console.log("exited"));
`,
],
env: bunEnv,
stderr: "pipe",
});
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
expect({ stdout: stdout.trim(), stderr, exitCode, signalCode: proc.signalCode }).toEqual({
stdout: "exited",
stderr: "",
exitCode: 0,
signalCode: null,
});
} finally {
fs.closeSync(readFd);
}
});

it("unref() after a write already went pending lets the process exit", async () => {
const { fifo, readFd } = blockedFifo("unref-after-write");
try {
await using proc = Bun.spawn({
cmd: [
bunExe(),
"-e",
`
const sink = Bun.file(${JSON.stringify(fifo)}).writer();
sink.write(Buffer.alloc(4 * 1024 * 1024, 0x61));
sink.unref();
process.on("exit", () => console.log("exited"));
`,
],
env: bunEnv,
stderr: "pipe",
});
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
expect({ stdout: stdout.trim(), stderr, exitCode, signalCode: proc.signalCode }).toEqual({
stdout: "exited",
stderr: "",
exitCode: 0,
signalCode: null,
});
} finally {
fs.closeSync(readFd);
}
});

it("ref() restores the keep-alive that unref() turned off", async () => {
const { fifo, readFd } = blockedFifo("ref-restores");
try {
await using proc = Bun.spawn({
cmd: [
bunExe(),
"-e",
`
const sink = Bun.file(${JSON.stringify(fifo)}).writer();
sink.unref();
sink.ref();
sink.write(Buffer.alloc(4 * 1024 * 1024, 0x61));
console.log("writing");
`,
],
env: bunEnv,
stdout: "pipe",
stderr: "pipe",
});

// Wait until the child has issued the write (so the pending keep-alive is
// in play), then confirm it does NOT exit on its own: ref() put the
// in-flight-write keep-alive back that unref() had turned off. The reader
// never drains, so the write stays pending and the only thing that can
// exit the process is a missing keep-alive.
for await (const line of proc.stdout.values()) {
if (new TextDecoder().decode(line).includes("writing")) break;
}
const outcome = await Promise.race([
proc.exited.then(code => ({ exitedOnItsOwn: true, code })),
Bun.sleep(1500).then(() => ({ exitedOnItsOwn: false })),
]);
expect(outcome).toEqual({ exitedOnItsOwn: false });

proc.kill();
await proc.exited;
} finally {
fs.closeSync(readFd);
}
});

it("ref()/unref() on a regular file do nothing and do not break writes", async () => {
const target = join(tmpdirSync(), "ref-noop.txt");
const sink = Bun.file(target).writer();
sink.unref();
sink.ref();
sink.write("hello");
await sink.end();
expect(await Bun.file(target).text()).toBe("hello");
});
});
Loading