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
35 changes: 25 additions & 10 deletions src/runtime/node/node_zlib_binding.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,8 +11,8 @@ use bun_io::KeepAlive;
use bun_jsc::ConcurrentTask::{ConcurrentTask, Task};
use bun_jsc::virtual_machine::VirtualMachine;
use bun_jsc::{
self as jsc, CallFrame, ErrorCode, JSGlobalObject, JSValue, JsCell, JsResult, StringJsc as _,
StrongOptional, WorkPoolTask,
self as jsc, CallFrame, ErrorCode, JSGlobalObject, JSValue, JsCell, JsRef, JsResult,
StringJsc as _, WorkPoolTask,
};
use bun_threading::work_pool::WorkPool;
use bun_zlib;
Expand Down Expand Up @@ -237,7 +237,12 @@ pub(crate) trait CompressionStreamImpl: Sized + Taskable + 'static {
}

fn poll_ref(&self) -> &JsCell<CountedKeepAlive>;
fn this_value(&self) -> &JsCell<StrongOptional>;
/// Back-reference to this class's own JS wrapper. Held strong only while
/// an async write is in flight on the work pool (the wrapper has no
/// pending-activity hook, so this is the sole GC root for the wrapper and
/// the cached `pendingInput`/`pendingOutput` buffers across the thread
/// hop). Cleared when the write completes back on the JS thread.
fn this_value(&self) -> &JsCell<JsRef>;
fn task(&self) -> &JsCell<WorkPoolTask>;
fn write_in_progress(&self) -> &Cell<bool>;
fn pending_close(&self) -> &Cell<bool>;
Expand Down Expand Up @@ -437,10 +442,11 @@ impl<T: CompressionStreamImpl> CompressionStream<T> {
s.set_flush(i32::try_from(flush).expect("int cast"));
});

// Only create the strong handle when we have a pending write
// And make sure to clear it when we are done.
// Root the wrapper only while a write is in flight: nothing else keeps
// it (or the pinned pending buffers it caches) alive across the
// work-pool hop. `run_from_js_thread` clears it on completion.
this.this_value()
.with_mut(|v| v.set(global_this, this_value));
.with_mut(|v| v.set_strong(this_value, global_this));

// SAFETY: `bun_vm()` never returns null for a Bun-owned global.
let vm = global_this.bun_vm();
Expand Down Expand Up @@ -517,8 +523,14 @@ impl<T: CompressionStreamImpl> CompressionStream<T> {

this.write_in_progress().set(false);

// Clear the strong handle before we call any callbacks.
let Some(this_value) = this.this_value().with_mut(|v| v.try_swap()) else {
// Take the rooted wrapper and drop the strong ref before invoking any
// callbacks; `this_value` stays valid for this frame via the native
// stack + `ensure_still_alive` below.
let Some(this_value) = this.this_value().with_mut(|v| {
let value = v.try_get();
*v = JsRef::empty();
value
}) else {
bun_output::scoped_log!(zlib, "this_value is null in runFromJSThread");
this.poll_ref().with_mut(|p| p.unref(vm));
// SAFETY: matching `ref_()` in `write()`; `this_ptr` is the heap
Expand Down Expand Up @@ -749,7 +761,7 @@ impl<T: CompressionStreamImpl> CompressionStream<T> {
}
this.pending_close().set(false);
this.closed().set(true);
this.this_value().with_mut(|v| v.deinit());
this.this_value().with_mut(|v| *v = JsRef::empty());
this.stream().with_mut(|s| s.close());
}

Expand Down Expand Up @@ -855,6 +867,9 @@ impl<T: CompressionStreamImpl> CompressionStream<T> {
}

pub(crate) fn finalize(this: Box<T>) {
// The wrapper is being collected; mark the back-ref terminal so no
// later accessor can observe the dead JSValue.
this.this_value().with_mut(|v| v.finalize());
// Refcounted: release the JS wrapper's +1; allocation may outlive this
// call if other refs remain, so hand ownership back to the raw refcount.
// SAFETY: `this` was the unique GC-owned m_ctx; `deref` frees on count==0.
Expand Down Expand Up @@ -1000,7 +1015,7 @@ macro_rules! __impl_compression_stream {
#[inline] fn global_this(&self) -> &::bun_jsc::JSGlobalObject { self.global_this.get() }
#[inline] fn stream(&self) -> &::bun_jsc::JsCell<Self::Stream> { &self.stream }
#[inline] fn poll_ref(&self) -> &::bun_jsc::JsCell<$crate::node::node_zlib_binding::CountedKeepAlive> { &self.poll_ref }
#[inline] fn this_value(&self) -> &::bun_jsc::JsCell<::bun_jsc::StrongOptional> { &self.this_value }
#[inline] fn this_value(&self) -> &::bun_jsc::JsCell<::bun_jsc::JsRef> { &self.this_value }
#[inline] fn task(&self) -> &::bun_jsc::JsCell<::bun_jsc::WorkPoolTask> { &self.task }
#[inline] fn write_in_progress(&self) -> &::core::cell::Cell<bool> { &self.write_in_progress }
#[inline] fn pending_close(&self) -> &::core::cell::Cell<bool> { &self.pending_close }
Expand Down
11 changes: 5 additions & 6 deletions src/runtime/node/zlib/NativeBrotli.rs
Original file line number Diff line number Diff line change
Expand Up @@ -62,8 +62,8 @@ mod _impl {
use core::ffi::c_uint;

use bun_jsc::{
CallFrame, ErrorCode, JSGlobalObject, JSValue, JsCell, JsResult, RangeErrorOptions,
StrongOptional, WorkPoolTask,
CallFrame, ErrorCode, JSGlobalObject, JSValue, JsCell, JsRef, JsResult, RangeErrorOptions,
WorkPoolTask,
};

use crate::node::node_zlib_binding::{CompressionStream, CountedKeepAlive, Error};
Expand All @@ -87,8 +87,7 @@ mod _impl {
pub global_this: bun_ptr::BackRef<JSGlobalObject>,
pub stream: JsCell<Context>,
pub poll_ref: JsCell<CountedKeepAlive>,
// TODO: Strong self-ref on the wrapper → JsRef per PORTING.md §JSC (Strong back-ref to own wrapper leaks)
pub this_value: JsCell<StrongOptional>, // Strong.Optional — empty-initialised
pub this_value: JsCell<JsRef>,
pub write_in_progress: Cell<bool>,
pub pending_close: Cell<bool>,
pub closed: Cell<bool>,
Expand Down Expand Up @@ -150,7 +149,7 @@ mod _impl {
global_this: bun_ptr::BackRef::new(global_this),
stream: JsCell::new(stream),
poll_ref: JsCell::new(CountedKeepAlive::default()),
this_value: JsCell::new(StrongOptional::empty()),
this_value: JsCell::new(JsRef::empty()),
write_in_progress: Cell::new(false),
pending_close: Cell::new(false),
closed: Cell::new(false),
Expand Down Expand Up @@ -332,7 +331,7 @@ mod _impl {
// ordering. The `stream` close below is load-bearing:
// `Context` has no Drop, so the brotli encoder/decoder state would
// leak without it.
self.this_value.set(StrongOptional::empty());
self.this_value.set(JsRef::empty());
drop(self.poll_ref.replace(CountedKeepAlive::default()));
self.stream.with_mut(|s| match s.mode {
bun_zlib::NodeMode::BROTLI_ENCODE | bun_zlib::NodeMode::BROTLI_DECODE => s.close(),
Expand Down
10 changes: 4 additions & 6 deletions src/runtime/node/zlib/NativeZlib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,7 @@ mod _impl {
use super::*;
use core::cell::Cell;

use bun_jsc::{
CallFrame, JSGlobalObject, JSValue, JsCell, JsResult, StrongOptional, WorkPoolTask,
};
use bun_jsc::{CallFrame, JSGlobalObject, JSValue, JsCell, JsRef, JsResult, WorkPoolTask};

use crate::node::node_zlib_binding::{CompressionStream, CountedKeepAlive};
use crate::node::util::validators;
Expand All @@ -44,7 +42,7 @@ mod _impl {
pub global_this: bun_ptr::BackRef<JSGlobalObject>,
pub stream: JsCell<Context>,
pub poll_ref: JsCell<CountedKeepAlive>,
pub this_value: JsCell<StrongOptional>, // jsc.Strong.Optional
pub this_value: JsCell<JsRef>,
pub write_in_progress: Cell<bool>,
pub pending_close: Cell<bool>,
pub closed: Cell<bool>,
Expand Down Expand Up @@ -95,7 +93,7 @@ mod _impl {
global_this: bun_ptr::BackRef::new(global),
stream: JsCell::new(stream),
poll_ref: JsCell::new(CountedKeepAlive::default()),
this_value: JsCell::new(StrongOptional::empty()),
this_value: JsCell::new(JsRef::empty()),
write_in_progress: Cell::new(false),
pending_close: Cell::new(false),
closed: Cell::new(false),
Expand Down Expand Up @@ -254,7 +252,7 @@ mod _impl {
fn deinit(this: *mut Self) {
// SAFETY: called exactly once by IntrusiveRc when refcount hits 0; `this`
// is the heap::alloc pointer produced at construction. `this_value`
// (Strong) and `poll_ref` (CountedKeepAlive) are Drop types — freed by
// (JsRef) and `poll_ref` (CountedKeepAlive) are Drop types — freed by
// heap::take below.
unsafe {
(*this).stream.with_mut(|s| s.close());
Expand Down
9 changes: 4 additions & 5 deletions src/runtime/node/zlib/NativeZstd.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,7 @@ mod _impl {
use core::ptr;

use bun_jsc::{
self as jsc, CallFrame, JSGlobalObject, JSValue, JsCell, JsResult, StrongOptional,
WorkPoolTask,
self as jsc, CallFrame, JSGlobalObject, JSValue, JsCell, JsRef, JsResult, WorkPoolTask,
};
use bun_zstd::c; // `bun.c` translated-c-headers (ZSTD_* fns/consts live here)

Expand Down Expand Up @@ -43,7 +42,7 @@ mod _impl {
pub global_this: bun_ptr::BackRef<JSGlobalObject>,
pub stream: JsCell<Context>,
pub poll_ref: JsCell<CountedKeepAlive>,
pub this_value: JsCell<StrongOptional>, // jsc.Strong.Optional
pub this_value: JsCell<JsRef>,
pub write_in_progress: Cell<bool>,
pub pending_close: Cell<bool>,
pub closed: Cell<bool>,
Expand Down Expand Up @@ -107,7 +106,7 @@ mod _impl {
global_this: bun_ptr::BackRef::new(global),
stream: JsCell::new(stream),
poll_ref: JsCell::new(CountedKeepAlive::default()),
this_value: JsCell::new(StrongOptional::empty()),
this_value: JsCell::new(JsRef::empty()),
write_in_progress: Cell::new(false),
pending_close: Cell::new(false),
closed: Cell::new(false),
Expand Down Expand Up @@ -281,7 +280,7 @@ mod _impl {
}

// Called by RefCount when the count hits 0. `poll_ref` and `this_value`
// (Strong) cleanup are handled by their own Drop impls; the Box free is
// (JsRef) cleanup are handled by their own Drop impls; the Box free is
// handled by IntrusiveRc dropping the Box.
impl Drop for NativeZstd {
fn drop(&mut self) {
Expand Down
106 changes: 105 additions & 1 deletion test/js/node/zlib/zlib.test.js
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { deflateSync, gunzipSync, gzipSync, inflateSync } from "bun";
import { describe, expect, it } from "bun:test";
import { tmpdirSync } from "harness";
import { bunEnv, bunExe, tmpdirSync } from "harness";
import * as buffer from "node:buffer";
import { randomFillSync } from "node:crypto";
import * as fs from "node:fs";
Expand Down Expand Up @@ -766,3 +766,107 @@ describe("crc32", () => {
expect(zlib.crc32("abc")).toBe(891568578);
});
});

// The native handles hold a back-reference to their own JS wrapper as a JsRef
// that is upgraded to a strong GC root only while an async write is in flight
// on the work pool, and dropped once the write completes. This test drives
// the native handle directly (bypassing the Transform wrapper) so both halves
// of that lifecycle are exercised under forced GC.
describe("native handle wrapper JsRef lifecycle", () => {
it("roots the wrapper across an in-flight write and releases it afterwards", async () => {
await using proc = Bun.spawn({
cmd: [
bunExe(),
"-e",
`
const zlib = require("node:zlib");
const { heapStats } = require("bun:jsc");
const {
DEFLATE, BROTLI_ENCODE, ZSTD_COMPRESS,
Z_FINISH, Z_DEFAULT_WINDOWBITS, Z_DEFAULT_COMPRESSION,
Z_DEFAULT_MEMLEVEL, Z_DEFAULT_STRATEGY,
BROTLI_OPERATION_FINISH, ZSTD_e_end,
} = zlib.constants;

const NativeZlib = zlib.createDeflate()._handle.constructor;
const NativeBrotli = zlib.createBrotliCompress()._handle.constructor;
const NativeZstd = zlib.createZstdCompress()._handle.constructor;

const input = Buffer.alloc(256, "in-flight GC root probe");

// Construct + init + write in its own frame so no stale bytecode register holds
// the wrapper once we return; only the JsRef strong root (set in write()) keeps
// it alive across the forced GC below.
function schedule(name, Ctor, mode, init, flush) {
const { promise, resolve, reject } = Promise.withResolvers();
const ws = new Uint32Array(2);
const out = Buffer.alloc(512);
const h = new Ctor(mode);
init(h, ws, resolve);
h.onerror = (m, e, c) => reject(new Error(name + " " + c + ": " + m));
h.write(flush, input, 0, input.length, out, 0, out.length);
return { weak: new WeakRef(h), promise, ws, out };
}

async function checkRooting(name, Ctor, mode, init, flush) {
for (let i = 0; i < 4; i++) {
const { weak, promise, ws, out } = schedule(name, Ctor, mode, init, flush);
Bun.gc(true);
if (weak.deref() === undefined) {
throw new Error(name + ": wrapper collected while an async write was in flight");
}
await promise;
const have = out.length - ws[0];
if (have <= 0) throw new Error(name + ": no output produced");
}
}

await checkRooting("zlib", NativeZlib, DEFLATE,
(h, ws, cb) => h.init(Z_DEFAULT_WINDOWBITS, Z_DEFAULT_COMPRESSION, Z_DEFAULT_MEMLEVEL, Z_DEFAULT_STRATEGY, ws, cb, undefined),
Z_FINISH);
await checkRooting("brotli", NativeBrotli, BROTLI_ENCODE,
(h, ws, cb) => h.init(new Uint32Array(0), ws, cb),
BROTLI_OPERATION_FINISH);
await checkRooting("zstd", NativeZstd, ZSTD_COMPRESS,
(h, ws, cb) => h.init(new Uint32Array(0), undefined, ws, cb),
ZSTD_e_end);

// Wrappers must be collectable once the write completes and no JS reference
// remains: the strong root is dropped in run_from_js_thread, so the live
// count stays bounded rather than growing by 50 per batch.
async function batch() {
const done = [];
for (let i = 0; i < 50; i++) {
const { promise, resolve } = Promise.withResolvers();
const ws = new Uint32Array(2);
const out = Buffer.alloc(512);
const h = new NativeZlib(DEFLATE);
h.init(Z_DEFAULT_WINDOWBITS, Z_DEFAULT_COMPRESSION, Z_DEFAULT_MEMLEVEL, Z_DEFAULT_STRATEGY, ws, resolve, undefined);
h.write(Z_FINISH, input, 0, input.length, out, 0, out.length);
done.push(promise);
}
await Promise.all(done);
await new Promise(r => setImmediate(r));
}
const counts = [];
for (let i = 0; i < 3; i++) {
await batch();
Bun.gc(true);
await new Promise(r => setImmediate(r));
Bun.gc(true);
counts.push(heapStats().objectTypeCounts?.NativeZlib || 0);
}
const max = Math.max(...counts.slice(1));
if (max > 25) throw new Error("NativeZlib wrappers leaking across GC: " + JSON.stringify(counts));

console.log("OK");
`,
],
env: bunEnv,
stdout: "pipe",
stderr: "pipe",
});
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
expect({ stdout, stderr, exitCode }).toEqual({ stdout: "OK\n", stderr: "", exitCode: 0 });
});
});
Loading