| ... | @@ -568,52 +568,33 @@ pub fn BufferedOutStreamCustom(comptime buffer_size: usize, comptime OutStreamEr | ... | @@ -568,52 +568,33 @@ pub fn BufferedOutStreamCustom(comptime buffer_size: usize, comptime OutStreamEr |
| 568 | | 568 | |
| 569 | unbuffered_out_stream: *Stream, | 569 | unbuffered_out_stream: *Stream, |
| 570 | | 570 | |
| 571 | buffer: [buffer_size]u8, | 571 | const FifoType = std.fifo.LinearFifo(u8, std.fifo.LinearFifoBufferType{ .Static = buffer_size }); |
| 572 | index: usize, | 572 | fifo: FifoType, |
| 573 | | 573 | |
| 574 | pub fn init(unbuffered_out_stream: *Stream) Self { | 574 | pub fn init(unbuffered_out_stream: *Stream) Self { |
| 575 | return Self{ | 575 | return Self{ |
| 576 | .unbuffered_out_stream = unbuffered_out_stream, | 576 | .unbuffered_out_stream = unbuffered_out_stream, |
| 577 | .buffer = undefined, | 577 | .fifo = FifoType.init(), |
| 578 | .index = 0, | | |
| 579 | .stream = Stream{ .writeFn = writeFn }, | 578 | .stream = Stream{ .writeFn = writeFn }, |
| 580 | }; | 579 | }; |
| 581 | } | 580 | } |
| 582 | | 581 | |
| 583 | pub fn flush(self: *Self) !void { | 582 | pub fn flush(self: *Self) !void { |
| 584 | try self.unbuffered_out_stream.write(self.buffer[0..self.index]); | 583 | while (true) { |
| 585 | self.index = 0; | 584 | const slice = self.fifo.readableSlice(0); |
| | 585 | if (slice.len == 0) break; |
| | 586 | try self.unbuffered_out_stream.write(slice); |
| | 587 | self.fifo.discard(slice.len); |
| | 588 | } |
| 586 | } | 589 | } |
| 587 | | 590 | |
| 588 | fn writeFn(out_stream: *Stream, bytes: []const u8) Error!void { | 591 | fn writeFn(out_stream: *Stream, bytes: []const u8) Error!void { |
| 589 | const self = @fieldParentPtr(Self, "stream", out_stream); | 592 | const self = @fieldParentPtr(Self, "stream", out_stream); |
| 590 | | 593 | if (bytes.len >= self.fifo.writableLength()) { |
| 591 | if (bytes.len == 1) { | | |
| 592 | // This is not required logic but a shorter path | | |
| 593 | // for single byte writes | | |
| 594 | self.buffer[self.index] = bytes[0]; | | |
| 595 | self.index += 1; | | |
| 596 | if (self.index == buffer_size) { | | |
| 597 | try self.flush(); | | |
| 598 | } | | |
| 599 | return; | | |
| 600 | } else if (bytes.len >= self.buffer.len) { | | |
| 601 | try self.flush(); | 594 | try self.flush(); |
| 602 | return self.unbuffered_out_stream.write(bytes); | 595 | return self.unbuffered_out_stream.write(bytes); |
| 603 | } | 596 | } |
| 604 | var src_index: usize = 0; | 597 | self.fifo.writeAssumeCapacity(bytes); |
| 605 | | | |
| 606 | while (src_index < bytes.len) { | | |
| 607 | const dest_space_left = self.buffer.len - self.index; | | |
| 608 | const copy_amt = math.min(dest_space_left, bytes.len - src_index); | | |
| 609 | mem.copy(u8, self.buffer[self.index..], bytes[src_index .. src_index + copy_amt]); | | |
| 610 | self.index += copy_amt; | | |
| 611 | assert(self.index <= self.buffer.len); | | |
| 612 | if (self.index == self.buffer.len) { | | |
| 613 | try self.flush(); | | |
| 614 | } | | |
| 615 | src_index += copy_amt; | | |
| 616 | } | | |
| 617 | } | 598 | } |
| 618 | }; | 599 | }; |
| 619 | } | 600 | } |