bufio: preserve partial flush progress
This commit is contained in:
@@ -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;
|
||||
};
|
||||
|
||||
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user