From 09b6f22e85e7b2dab3523cf499a65704dc94a96b Mon Sep 17 00:00:00 2001 From: Hojun-Cho Date: Sun, 9 Aug 2026 18:05:21 +0900 Subject: [PATCH] bufio: preserve partial flush progress --- lib/bufio/stream.ww | 63 +++++++++++++++++++----------- lib/bufio/stream_test.ww | 83 ++++++++++++++++++++++++++++++++++++++++ 2 files changed, 123 insertions(+), 23 deletions(-) diff --git a/lib/bufio/stream.ww b/lib/bufio/stream.ww index 0ec6e01c..a1eafe3d 100644 --- a/lib/bufio/stream.ww +++ b/lib/bufio/stream.ww @@ -92,6 +92,30 @@ export fn init(src: io.stream, rbuf: []u8, wbuf: []u8) stream = { return r; }; +fn bufiomove(dst: *u8, src: *u8, n: i32) void = { + if (n <= 0 || dst == src) { return; }; + let i: i32 = 0; + if ((dst: uintptr) > (src: uintptr)) { + i = n; + for (i > 0) { + i -= 1; + dst[i] = src[i]; + }; + } else { + for (i < n) { + dst[i] = src[i]; + i += 1; + }; + }; +}; + +fn flushpending(b: *stream, off: i32) void = { + if (off == 0) { return; }; + let remain: i32 = b.wend - off; + bufiomove(b.wbuf.ptr, b.wbuf.ptr + (off: u64), remain); + b.wend = remain; +}; + // setflush — install a new flush byte-set. Any byte from `bs` appearing // in a write payload triggers an automatic flush after the write copies // into wbuf. Mirrors ref/hare/bufio/stream.ha:128. @@ -102,8 +126,9 @@ export fn setflush(b: *stream, bs: []u8) void = { // flush — drain any pending wbuf data to src. Public API AND the // internal drain called by bwrite / bclose. A 0-byte non-error write // means the sink made no progress (memio.fixed full) — surfaced as a -// nomem-carried io.error rather than spinning, matching what io.writeall -// would do in Hare. Mirrors ref/hare/bufio/stream.ha. +// nomem-carried io.error rather than spinning. Accepted prefixes are removed +// from wbuf before any failure is returned, so a retry cannot duplicate them. +// Mirrors ref/hare/bufio/stream.ha. export fn flush(b: *stream) (void | io.error) = { if (b.wend == 0) { return; }; let off: i32 = 0; @@ -111,14 +136,20 @@ export fn flush(b: *stream) (void | io.error) = { let r: (size | io.error) = io.write(b.src, b.wbuf[off:b.wend]); match (r) { case let n: size => { + assert(n <= (b.wend - off): size, + "bufio.flush: writer returned an oversized count"); if (n == 0: size) { + flushpending(b, off); let nm: nomem; let e: io.error = nm; return e; }; off += n: i32; }; - case let e: io.error => return e; + case let e: io.error => { + flushpending(b, off); + return e; + }; }; }; b.wend = 0; @@ -133,11 +164,8 @@ export fn unread(b: *stream, buf: []u8) void = { if (b.rstart < buf.len) { rtabort("bufio.unread: more data than rbuf has room for"); }; - let i: i32 = 0; - for (i < buf.len) { - b.rbuf[b.rstart - buf.len + i] = buf[i]; - i += 1; - }; + bufiomove(b.rbuf.ptr + ((b.rstart - buf.len): u64), buf.ptr, + buf.len); b.rstart -= buf.len; }; @@ -176,11 +204,7 @@ fn bread(s: io.stream, buf: []u8) (size | io.eof | io.error) = { }; let avail: i32 = b.rend - b.rstart; if (avail < buf.len && avail < b.rbuf.len) { - let i: i32 = 0; - for (i < avail) { - b.rbuf[i] = b.rbuf[b.rstart + i]; - i += 1; - }; + bufiomove(b.rbuf.ptr, b.rbuf.ptr + (b.rstart: u64), avail); b.rstart = 0; b.rend = avail; let r: (size | io.eof | io.error) = io.read(b.src, b.rbuf[b.rend:b.rbuf.len]); @@ -195,11 +219,7 @@ fn bread(s: io.stream, buf: []u8) (size | io.eof | io.error) = { avail = b.rend - b.rstart; let n: i32 = buf.len; if (avail < n) { n = avail; }; - let i: i32 = 0; - for (i < n) { - buf[i] = b.rbuf[b.rstart + i]; - i += 1; - }; + bufiomove(buf.ptr, b.rbuf.ptr + (b.rstart: u64), n); b.rstart += n; return n: size; }; @@ -242,11 +262,8 @@ fn bwrite(s: io.stream, buf: []u8) (size | io.error) = { }; let n: i32 = buf.len - z; if (avail < n) { n = avail; }; - let i: i32 = 0; - for (i < n) { - b.wbuf[b.wend + i] = buf[z + i]; - i += 1; - }; + bufiomove(b.wbuf.ptr + (b.wend: u64), + buf.ptr + (z: u64), n); b.wend += n; z += n; }; diff --git a/lib/bufio/stream_test.ww b/lib/bufio/stream_test.ww index f96596e3..7516fd32 100644 --- a/lib/bufio/stream_test.ww +++ b/lib/bufio/stream_test.ww @@ -36,6 +36,71 @@ fn closesource() io.stream = { return &closevt; }; +type failstream = struct { + vt: io.vtable, + out: [16]u8, + pos: i32, + calls: i32, + zero: bool, +}; + +fn failwrite(s: io.stream, buf: []u8) (size | io.error) = { + let f: *failstream = s: *failstream; + f.calls += 1; + if (f.calls == 2) { + if (f.zero) { return 0: size; }; + let nm: nomem; + let e: io.error = nm; + return e; + }; + let n: i32 = buf.len; + if (f.calls == 1 && n > 2) { n = 2; }; + let i: i32 = 0; + for (i < n) { + f.out[f.pos + i] = buf[i]; + i += 1; + }; + f.pos += n; + return n: size; +}; + +fn checkpartialflush(zero: bool) void = { + let sink: failstream; + sink.vt.writer = (&failwrite): *io.writer; + sink.pos = 0; + sink.calls = 0; + sink.zero = zero; + let rb: [1]u8; + let wb: [4]u8; + let b: bufio.stream = bufio.init(&sink.vt, rb[0:1], wb[0:4]); + let src: [4]u8; + let _: i32 = sputstr("ABCD", src[0:4], 0); + let wr: (size | io.error) = io.write(&b.vt, src[0:4]); + match (wr) { case let n: size => {}; case let e: io.error => abort(); }; + let first: (void | io.error) = bufio.flush(&b); + match (first) { + case void => abort(); + case let e: io.error => assert(e is nomem); + }; + assert(sink.pos == 2); + assert(b.wend == 2); + assert(wb[0] == 'C' && wb[1] == 'D'); + let second: (void | io.error) = bufio.flush(&b); + match (second) { case void => {}; case let e: io.error => abort(); }; + assert(b.wend == 0); + assert(sink.pos == 4); + assert(sink.out[0] == 'A' && sink.out[1] == 'B'); + assert(sink.out[2] == 'C' && sink.out[3] == 'D'); +}; + +@test fn streampartialflusherror() void = { + checkpartialflush(false); +}; + +@test fn streampartialflushzero() void = { + checkpartialflush(true); +}; + @test fn streamsmallwrite() void = { let raw: [16]u8; let mem: memio.stream = memio.fixed(raw[0:16]); @@ -313,6 +378,24 @@ fn closesource() io.stream = { assert(!(out[3] != 71u8)); // 'G' }; +@test fn streamunreadoverlap() void = { + let raw: [1]u8; + let mem: memio.stream = memio.fixed(raw[0:0]); + let rb: [4]u8; + let _: i32 = sputstr("ABCD", rb[0:4], 0); + let wb: [1]u8; + let b: bufio.stream = bufio.init(&mem.vt, rb[0:4], wb[0:1]); + bufio.unread(&b, rb[0:3]); + let out: [3]u8; + let r: (size | io.eof | io.error) = io.read(&b.vt, out[0:3]); + match (r) { + case let n: size => assert(n: i32 == 3); + case io.eof => abort(); + case let e: io.error => abort(); + }; + assert(out[0] == 'A' && out[1] == 'B' && out[2] == 'C'); +}; + @test fn streamscannerunread() void = { let raw: [16]u8; let n: i32 = sputstr("hello\n", raw[0:16], 0);