diff --git a/src/io/PipeWriter.rs b/src/io/PipeWriter.rs index b21b6b70a339..6136309ecae2 100644 --- a/src/io/PipeWriter.rs +++ b/src/io/PipeWriter.rs @@ -887,10 +887,13 @@ impl PosixStreamingWriter { 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` @@ -903,20 +906,28 @@ impl PosixStreamingWriter { 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(); 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); @@ -957,7 +968,7 @@ impl PosixStreamingWriter { 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); diff --git a/test/js/bun/util/filesink.test.ts b/test/js/bun/util/filesink.test.ts index 4fbfe3902743..d9e087e587ad 100644 --- a/test/js/bun/util/filesink.test.ts +++ b/test/js/bun/util/filesink.test.ts @@ -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"; @@ -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(