Skip to content
Closed
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
25 changes: 18 additions & 7 deletions src/io/PipeWriter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -887,10 +887,13 @@ impl<Parent: PosixStreamingWriterParent> PosixStreamingWriter<Parent> {
return WriteResult::Wrote(buf_len);
}

self.try_write_newly_buffered_data()
self.try_write_newly_buffered_data(buf_len)
}

fn try_write_newly_buffered_data(&mut self) -> WriteResult {
/// `chunk_len` is the caller's chunk, already appended to `outgoing`.
/// `Wrote(n)` reports bytes of *that chunk* accepted, not bytes drained
/// from the buffer (older chunks included) — that is `flush()`'s accounting.
fn try_write_newly_buffered_data(&mut self, chunk_len: usize) -> WriteResult {
debug_assert!(!self.is_done);

// Borrow `self.outgoing` only for the syscall. `try_write` takes `&self`
Expand All @@ -903,20 +906,28 @@ impl<Parent: PosixStreamingWriterParent> PosixStreamingWriter<Parent> {

match rc {
WriteResult::Wrote(amt) => {
if amt == self.outgoing.size() {
self.outgoing.reset();
self.parent_on_write(amt, WriteStatus::Drained);
} else {
if amt < self.outgoing.size() {
self.outgoing.wrote(amt);
self.parent_on_write(amt, WriteStatus::Pending);
Self::register_poll(self);
return WriteResult::Pending(amt);
}

self.outgoing.reset();
self.parent_on_write(amt, WriteStatus::Drained);
return WriteResult::Wrote(chunk_len);
}
WriteResult::Done(amt) => {
// `amt` drains the whole buffer (older chunks first); report
// only the portion that belongs to this chunk, like `Wrote`.
let old_buffered = self.outgoing.size().saturating_sub(chunk_len);
self.outgoing.reset();
Comment thread
robobun marked this conversation as resolved.
self.parent_on_write(amt, WriteStatus::EndOfFile);
return WriteResult::Done(amt.saturating_sub(old_buffered));
}
// `Pending` returns a promise whose resolution value is accounted
// for separately (see `FileSink::bytes_accepted`), so it is left as
// the drained count here.
WriteResult::Pending(amt) => {
self.outgoing.wrote(amt);
self.parent_on_write(amt, WriteStatus::Pending);
Expand Down Expand Up @@ -957,7 +968,7 @@ impl<Parent: PosixStreamingWriterParent> PosixStreamingWriter<Parent> {
return WriteResult::Err(sys::Error::oom());
}

return self.try_write_newly_buffered_data();
return self.try_write_newly_buffered_data(buf.len());
}

let rc = self.try_write(self.force_sync, buf);
Expand Down
37 changes: 36 additions & 1 deletion test/js/bun/util/filesink.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { createSocketPair, fileSinkInternals } from "bun:internal-for-testing";
import { describe, expect, it } from "bun:test";
import { bunEnv, bunExe, fileDescriptorLeakChecker, isPosix, isWindows, tmpdirSync } from "harness";
import { bunEnv, bunExe, fileDescriptorLeakChecker, isPosix, isWindows, tempDir, tmpdirSync } from "harness";
import { mkfifo } from "mkfifo";
import { join } from "node:path";

Expand Down Expand Up @@ -273,6 +273,41 @@ it.skipIf(!isPosix)("a backpressured string write() resolves to its encoded byte
expect(received).toBe(size);
});

// A chunk big enough to flush the bytes buffered by earlier writes must still
// report its own size. Writes are issued without `await` in between so the
// auto-flush microtask cannot drain the buffer first. Skipped on Windows: there
// every write to a file-backed sink returns a (shared) promise instead.
it.skipIf(!isPosix)("write result is not cumulative when the chunk flushes buffered bytes", async () => {
using dir = tempDir("filesink-cumulative", {});
const filename = path.join(String(dir), "test.bin");
const writer = Bun.file(filename).writer();

const ascii = Buffer.alloc(505, "a").toString();
const bytes = Buffer.alloc(64 * 1024, "c");
const latin1 = Buffer.alloc(40_000, "é").toString();
const utf16 = Buffer.alloc(20_000, "😀").toString();

const results = [
writer.write(ascii), // buffered
writer.write(bytes), // flushes `ascii` + itself
writer.write(Buffer.alloc(10, "d")), // buffered
writer.write(latin1), // flushes the 10 bytes + itself
writer.write("e"), // buffered
writer.write(utf16), // flushes the 1 byte + itself
];
await writer.end();

expect(results).toEqual([
Buffer.byteLength(ascii),
bytes.byteLength,
10,
Buffer.byteLength(latin1),
1,
Buffer.byteLength(utf16),
]);
expect(results.reduce((sum, n) => sum + n, 0)).toBe(fs.statSync(filename).size);
});

if (isWindows) {
it("ENOENT, Windows", () => {
expect(() => Bun.file("A:\\this-does-not-exist.txt").writer()).toThrow(
Expand Down
Loading