authorgravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2025-04-21 14:09:47-07:00
committergravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2025-07-01 16:35:27-07:00
log19b82ca7ab7c44c3a1f4d0c6cbaf43bc0bded43a
tree7ffeb0cf448f9ed7a43989b8363d633c62566976
parentd62f22cc4d2f17db0877cd9316bf8be9470d04a5

std.http.Server: implement chunked request parsing


5 files changed, 328 insertions(+), 196 deletions(-)

lib/compiler/std-docs.zig+1-1
...@@ -428,7 +428,7 @@ fn receiveWasmMessage(...@@ -428,7 +428,7 @@ fn receiveWasmMessage(
428 },428 },
429 else => {429 else => {
430 // Ignore other messages.430 // Ignore other messages.
431 try br.discard(header.bytes_len);431 try br.discardAll(header.bytes_len);
432 },432 },
433 }433 }
434}434}
lib/std/http/Server.zig+250-168
...@@ -22,9 +22,6 @@ out: *std.io.BufferedWriter,...@@ -22,9 +22,6 @@ out: *std.io.BufferedWriter,
22state: State,22state: State,
23head_parse_err: Request.Head.ParseError,23head_parse_err: Request.Head.ParseError,
2424
25/// being deleted...
26next_request_start: usize = 0,
27
28pub const State = enum {25pub const State = enum {
29 /// The connection is available to be used for the first time, or reused.26 /// The connection is available to be used for the first time, or reused.
30 ready,27 ready,
...@@ -95,6 +92,8 @@ pub fn receiveHead(s: *Server) ReceiveHeadError!Request {...@@ -95,6 +92,8 @@ pub fn receiveHead(s: *Server) ReceiveHeadError!Request {
95 if (hp.state == .finished) return .{92 if (hp.state == .finished) return .{
96 .server = s,93 .server = s,
97 .head_end = head_end,94 .head_end = head_end,
95 .trailers_len = 0,
96 .read_err = null,
98 .head = Request.Head.parse(buf[0..head_end]) catch |err| {97 .head = Request.Head.parse(buf[0..head_end]) catch |err| {
99 s.head_parse_err = err;98 s.head_parse_err = err;
100 return error.HttpHeadersInvalid;99 return error.HttpHeadersInvalid;
...@@ -108,11 +107,38 @@ pub const Request = struct {...@@ -108,11 +107,38 @@ pub const Request = struct {
108 server: *Server,107 server: *Server,
109 /// Index into `Server.in` internal buffer.108 /// Index into `Server.in` internal buffer.
110 head_end: usize,109 head_end: usize,
110 /// Number of bytes of HTTP trailers. These are at the end of a
111 /// transfer-encoding: chunked message.
112 trailers_len: usize,
111 head: Head,113 head: Head,
112 reader_state: union {114 reader_state: union {
113 remaining_content_length: u64,115 remaining_content_length: u64,
114 chunk_parser: http.ChunkParser,116 remaining_chunk_len: RemainingChunkLen,
115 },117 },
118 read_err: ?ReadError,
119
120 pub const ReadError = error{
121 HttpChunkInvalid,
122 HttpHeadersOversize,
123 };
124
125 pub const max_chunk_header_len = 22;
126
127 pub const RemainingChunkLen = enum(u64) {
128 head = 0,
129 n = 1,
130 rn = 2,
131 done = std.math.maxInt(u64),
132 _,
133
134 pub fn init(integer: u64) RemainingChunkLen {
135 return @enumFromInt(integer);
136 }
137
138 pub fn int(rcl: RemainingChunkLen) u64 {
139 return @intFromEnum(rcl);
140 }
141 };
116142
117 pub const Compression = union(enum) {143 pub const Compression = union(enum) {
118 deflate: std.compress.zlib.Decompressor,144 deflate: std.compress.zlib.Decompressor,
...@@ -559,177 +585,240 @@ pub const Request = struct {...@@ -559,177 +585,240 @@ pub const Request = struct {
559 };585 };
560 }586 }
561587
562 pub const ReadError = net.Stream.ReadError || error{588 fn contentLengthRead(
563 HttpChunkInvalid,
564 HttpHeadersOversize,
565 };
566
567 fn contentLengthReader_read(
568 ctx: ?*anyopaque,589 ctx: ?*anyopaque,
569 bw: *std.io.BufferedWriter,590 bw: *std.io.BufferedWriter,
570 limit: std.io.Reader.Limit,591 limit: std.io.Reader.Limit,
571 ) std.io.Reader.Error!usize {592 ) std.io.Reader.RwError!usize {
572 const request: *Request = @alignCast(@ptrCast(ctx));593 const request: *Request = @alignCast(@ptrCast(ctx));
573 _ = request;594 const remaining_content_length = &request.reader_state.remaining_content_length;
574 _ = bw;595 const remaining = remaining_content_length.*;
575 _ = limit;596 const server = request.server;
576 @panic("TODO");597 if (remaining == 0) {
598 server.state = .ready;
599 return error.EndOfStream;
600 }
601 const n = try server.in.read(bw, limit.min(.limited(remaining)));
602 const new_remaining = remaining - n;
603 remaining_content_length.* = new_remaining;
604 return n;
577 }605 }
578606
579 fn contentLengthReader_readVec(ctx: ?*anyopaque, data: []const []u8) std.io.Reader.Error!usize {607 fn contentLengthReadVec(context: ?*anyopaque, data: []const []u8) std.io.Reader.Error!usize {
580 const request: *Request = @alignCast(@ptrCast(ctx));608 const request: *Request = @alignCast(@ptrCast(context));
581 _ = request;609 const remaining_content_length = &request.reader_state.remaining_content_length;
582 _ = data;610 const server = request.server;
583 @panic("TODO");611 const remaining = remaining_content_length.*;
612 if (remaining == 0) {
613 server.state = .ready;
614 return error.EndOfStream;
615 }
616 const n = try server.in.readVecLimit(data, .limited(remaining));
617 const new_remaining = remaining - n;
618 remaining_content_length.* = new_remaining;
619 return n;
584 }620 }
585621
586 fn contentLengthReader_discard(ctx: ?*anyopaque, limit: std.io.Reader.Limit) std.io.Reader.Error!usize {622 fn contentLengthDiscard(ctx: ?*anyopaque, limit: std.io.Reader.Limit) std.io.Reader.Error!usize {
587 const request: *Request = @alignCast(@ptrCast(ctx));623 const request: *Request = @alignCast(@ptrCast(ctx));
588 _ = request;624 const remaining_content_length = &request.reader_state.remaining_content_length;
589 _ = limit;625 const server = request.server;
590 @panic("TODO");626 const remaining = remaining_content_length.*;
627 if (remaining == 0) {
628 server.state = .ready;
629 return error.EndOfStream;
630 }
631 const n = try server.in.discard(limit.min(.limited(remaining)));
632 const new_remaining = remaining - n;
633 remaining_content_length.* = new_remaining;
634 return n;
591 }635 }
592636
593 fn chunkedReader_read(637 fn chunkedRead(
594 ctx: ?*anyopaque,638 ctx: ?*anyopaque,
595 bw: *std.io.BufferedWriter,639 bw: *std.io.BufferedWriter,
596 limit: std.io.Reader.Limit,640 limit: std.io.Reader.Limit,
597 ) std.io.Reader.Error!usize {641 ) std.io.Reader.RwError!usize {
598 const request: *Request = @alignCast(@ptrCast(ctx));642 const request: *Request = @alignCast(@ptrCast(ctx));
599 _ = request;643 const chunk_len_ptr = &request.reader_state.remaining_chunk_len;
600 _ = bw;644 const in = request.server.in;
601 _ = limit;645 len: switch (chunk_len_ptr.*) {
602 @panic("TODO");646 .head => {
647 var cp: http.ChunkParser = .init;
648 const i = cp.feed(in.bufferContents());
649 switch (cp.state) {
650 .invalid => return request.failRead(error.HttpChunkInvalid),
651 .data => {
652 if (i > max_chunk_header_len) return request.failRead(error.HttpChunkInvalid);
653 in.toss(i);
654 },
655 else => {
656 try in.fill(max_chunk_header_len);
657 const next_i = cp.feed(in.bufferContents()[i..]);
658 if (cp.state != .data) return request.failRead(error.HttpChunkInvalid);
659 const header_len = i + next_i;
660 if (header_len > max_chunk_header_len) return request.failRead(error.HttpChunkInvalid);
661 in.toss(header_len);
662 },
663 }
664 if (cp.chunk_len == 0) return parseTrailers(request, 0);
665 const n = try in.read(bw, limit.min(.limited(cp.chunk_len)));
666 chunk_len_ptr.* = .init(cp.chunk_len + 2 - n);
667 return n;
668 },
669 .n => {
670 if ((try in.peekByte()) != '\n') return request.failRead(error.HttpChunkInvalid);
671 in.toss(1);
672 continue :len .head;
673 },
674 .rn => {
675 const rn = try in.peekArray(2);
676 if (rn[0] != '\r' or rn[1] != '\n') return request.failRead(error.HttpChunkInvalid);
677 in.toss(2);
678 continue :len .head;
679 },
680 else => |remaining_chunk_len| {
681 const n = try in.read(bw, limit.min(.limited(@intFromEnum(remaining_chunk_len) - 2)));
682 chunk_len_ptr.* = .init(@intFromEnum(remaining_chunk_len) - n);
683 return n;
684 },
685 .done => return error.EndOfStream,
686 }
603 }687 }
604688
605 fn chunkedReader_readVec(ctx: ?*anyopaque, data: []const []u8) std.io.Reader.Error!usize {689 fn chunkedReadVec(ctx: ?*anyopaque, data: []const []u8) std.io.Reader.Error!usize {
606 const request: *Request = @alignCast(@ptrCast(ctx));690 const request: *Request = @alignCast(@ptrCast(ctx));
607 _ = request;691 const chunk_len_ptr = &request.reader_state.remaining_chunk_len;
608 _ = data;692 const in = request.server.in;
609 @panic("TODO");693 var already_requested_more = false;
694 var amt_read: usize = 0;
695 data: for (data) |d| {
696 len: switch (chunk_len_ptr.*) {
697 .head => {
698 var cp: http.ChunkParser = .init;
699 const available_buffer = in.bufferContents();
700 const i = cp.feed(available_buffer);
701 if (cp.state == .invalid) return request.failRead(error.HttpChunkInvalid);
702 if (i == available_buffer.len) {
703 if (already_requested_more) {
704 chunk_len_ptr.* = .head;
705 return amt_read;
706 }
707 already_requested_more = true;
708 try in.fill(max_chunk_header_len);
709 const next_i = cp.feed(in.bufferContents()[i..]);
710 if (cp.state != .data) return request.failRead(error.HttpChunkInvalid);
711 const header_len = i + next_i;
712 if (header_len > max_chunk_header_len) return request.failRead(error.HttpChunkInvalid);
713 in.toss(header_len);
714 } else {
715 if (i > max_chunk_header_len) return request.failRead(error.HttpChunkInvalid);
716 in.toss(i);
717 }
718 if (cp.chunk_len == 0) return parseTrailers(request, amt_read);
719 continue :len .init(cp.chunk_len + 2);
720 },
721 .n => {
722 if (in.bufferContents().len < 1) already_requested_more = true;
723 if ((try in.takeByte()) != '\n') return request.failRead(error.HttpChunkInvalid);
724 continue :len .head;
725 },
726 .rn => {
727 if (in.bufferContents().len < 2) already_requested_more = true;
728 const rn = try in.takeArray(2);
729 if (rn[0] != '\r' or rn[1] != '\n') return request.failRead(error.HttpChunkInvalid);
730 continue :len .head;
731 },
732 else => |remaining_chunk_len| {
733 const available_buffer = in.bufferContents();
734 const copy_len = @min(available_buffer.len, d.len, remaining_chunk_len.int() - 2);
735 @memcpy(d[0..copy_len], available_buffer[0..copy_len]);
736 amt_read += copy_len;
737 in.toss(copy_len);
738 const next_chunk_len: RemainingChunkLen = .init(remaining_chunk_len.int() - copy_len);
739 if (copy_len == d.len) {
740 chunk_len_ptr.* = next_chunk_len;
741 continue :data;
742 }
743 if (already_requested_more) {
744 chunk_len_ptr.* = next_chunk_len;
745 return amt_read;
746 }
747 already_requested_more = true;
748 try in.fill(3);
749 continue :len next_chunk_len;
750 },
751 .done => return error.EndOfStream,
752 }
753 }
754 return amt_read;
610 }755 }
611756
612 fn chunkedReader_discard(ctx: ?*anyopaque, limit: std.io.Reader.Limit) std.io.Reader.Error!usize {757 fn chunkedDiscard(ctx: ?*anyopaque, limit: std.io.Reader.Limit) std.io.Reader.Error!usize {
613 const request: *Request = @alignCast(@ptrCast(ctx));758 const request: *Request = @alignCast(@ptrCast(ctx));
614 _ = request;759 const chunk_len_ptr = &request.reader_state.remaining_chunk_len;
615 _ = limit;760 const in = request.server.in;
616 @panic("TODO");761 len: switch (chunk_len_ptr.*) {
617 }762 .head => {
618763 var cp: http.ChunkParser = .init;
619 fn read_cl(context: *const anyopaque, buffer: []u8) ReadError!usize {764 const i = cp.feed(in.bufferContents());
620 const request: *Request = @alignCast(@ptrCast(context));765 switch (cp.state) {
621 const s = request.server;766 .invalid => return request.failRead(error.HttpChunkInvalid),
622767 .data => {
623 const remaining_content_length = &request.reader_state.remaining_content_length;768 if (i > max_chunk_header_len) return request.failRead(error.HttpChunkInvalid);
624 if (remaining_content_length.* == 0) {769 in.toss(i);
625 s.state = .ready;770 },
626 return 0;771 else => {
772 try in.fill(max_chunk_header_len);
773 const next_i = cp.feed(in.bufferContents()[i..]);
774 if (cp.state != .data) return request.failRead(error.HttpChunkInvalid);
775 const header_len = i + next_i;
776 if (header_len > max_chunk_header_len) return request.failRead(error.HttpChunkInvalid);
777 in.toss(header_len);
778 },
779 }
780 if (cp.chunk_len == 0) return parseTrailers(request, 0);
781 const n = try in.discard(limit.min(.limited(cp.chunk_len)));
782 chunk_len_ptr.* = .init(cp.chunk_len + 2 - n);
783 return n;
784 },
785 .n => {
786 if ((try in.peekByte()) != '\n') return request.failRead(error.HttpChunkInvalid);
787 in.toss(1);
788 continue :len .head;
789 },
790 .rn => {
791 const rn = try in.peekArray(2);
792 if (rn[0] != '\r' or rn[1] != '\n') return request.failRead(error.HttpChunkInvalid);
793 in.toss(2);
794 continue :len .head;
795 },
796 else => |remaining_chunk_len| {
797 const n = try in.discard(limit.min(.limited(remaining_chunk_len.int() - 2)));
798 chunk_len_ptr.* = .init(remaining_chunk_len.int() - n);
799 return n;
800 },
801 .done => return error.EndOfStream,
627 }802 }
628 assert(s.state == .receiving_body);
629 const available = try fill(s, request.head_end);
630 const len = @min(remaining_content_length.*, available.len, buffer.len);
631 @memcpy(buffer[0..len], available[0..len]);
632 remaining_content_length.* -= len;
633 s.next_request_start += len;
634 if (remaining_content_length.* == 0)
635 s.state = .ready;
636 return len;
637 }803 }
638804
639 fn fill(s: *Server, head_end: usize) ReadError![]u8 {805 /// Called when next bytes in the stream are trailers, or "\r\n" to indicate
640 const available = s.read_buffer[s.next_request_start..s.read_buffer_len];806 /// end of chunked body.
641 if (available.len > 0) return available;807 fn parseTrailers(request: *Request, amt_read: usize) std.io.Reader.Error!usize {
642 s.next_request_start = head_end;808 const in = request.server.in;
643 s.read_buffer_len = head_end + try s.connection.stream.read(s.read_buffer[head_end..]);809 var hp: http.HeadParser = .{};
644 return s.read_buffer[head_end..s.read_buffer_len];810 var trailers_len: usize = 0;
645 }811 while (true) {
646812 if (trailers_len >= in.buffer.len) return request.failRead(error.HttpHeadersOversize);
647 fn read_chunked(context: *const anyopaque, buffer: []u8) ReadError!usize {813 try in.fill(trailers_len + 1);
648 const request: *Request = @alignCast(@ptrCast(context));814 trailers_len += hp.feed(in.bufferContents()[trailers_len..]);
649 const s = request.server;815 if (hp.state == .finished) {
650816 request.reader_state.remaining_chunk_len = .done;
651 const cp = &request.reader_state.chunk_parser;817 request.server.state = .ready;
652 const head_end = request.head_end;818 request.trailers_len = trailers_len;
653819 return amt_read;
654 // Protect against returning 0 before the end of stream.
655 var out_end: usize = 0;
656 while (out_end == 0) {
657 switch (cp.state) {
658 .invalid => return 0,
659 .data => {
660 assert(s.state == .receiving_body);
661 const available = try fill(s, head_end);
662 const len = @min(cp.chunk_len, available.len, buffer.len);
663 @memcpy(buffer[0..len], available[0..len]);
664 cp.chunk_len -= len;
665 if (cp.chunk_len == 0)
666 cp.state = .data_suffix;
667 out_end += len;
668 s.next_request_start += len;
669 continue;
670 },
671 else => {
672 assert(s.state == .receiving_body);
673 const available = try fill(s, head_end);
674 const n = cp.feed(available);
675 switch (cp.state) {
676 .invalid => return error.HttpChunkInvalid,
677 .data => {
678 if (cp.chunk_len == 0) {
679 // The next bytes in the stream are trailers,
680 // or \r\n to indicate end of chunked body.
681 //
682 // This function must append the trailers at
683 // head_end so that headers and trailers are
684 // together.
685 //
686 // Since returning 0 would indicate end of
687 // stream, this function must read all the
688 // trailers before returning.
689 if (s.next_request_start > head_end) rebase(s, head_end);
690 var hp: http.HeadParser = .{};
691 {
692 const bytes = s.read_buffer[head_end..s.read_buffer_len];
693 const end = hp.feed(bytes);
694 if (hp.state == .finished) {
695 cp.state = .invalid;
696 s.state = .ready;
697 s.next_request_start = s.read_buffer_len - bytes.len + end;
698 return out_end;
699 }
700 }
701 while (true) {
702 const buf = s.read_buffer[s.read_buffer_len..];
703 if (buf.len == 0)
704 return error.HttpHeadersOversize;
705 const read_n = try s.connection.stream.read(buf);
706 s.read_buffer_len += read_n;
707 const bytes = buf[0..read_n];
708 const end = hp.feed(bytes);
709 if (hp.state == .finished) {
710 cp.state = .invalid;
711 s.state = .ready;
712 s.next_request_start = s.read_buffer_len - bytes.len + end;
713 return out_end;
714 }
715 }
716 }
717 const data = available[n..];
718 const len = @min(cp.chunk_len, data.len, buffer.len);
719 @memcpy(buffer[0..len], data[0..len]);
720 cp.chunk_len -= len;
721 if (cp.chunk_len == 0)
722 cp.state = .data_suffix;
723 out_end += len;
724 s.next_request_start += n + len;
725 continue;
726 },
727 else => continue,
728 }
729 },
730 }820 }
731 }821 }
732 return out_end;
733 }822 }
734823
735 pub const ReaderError = error{824 pub const ReaderError = error{
...@@ -752,7 +841,6 @@ pub const Request = struct {...@@ -752,7 +841,6 @@ pub const Request = struct {
752 const s = request.server;841 const s = request.server;
753 assert(s.state == .received_head);842 assert(s.state == .received_head);
754 s.state = .receiving_body;843 s.state = .receiving_body;
755 s.next_request_start = request.head_end;
756844
757 if (request.head.expect) |expect| {845 if (request.head.expect) |expect| {
758 if (mem.eql(u8, expect, "100-continue")) {846 if (mem.eql(u8, expect, "100-continue")) {
...@@ -765,13 +853,13 @@ pub const Request = struct {...@@ -765,13 +853,13 @@ pub const Request = struct {
765853
766 switch (request.head.transfer_encoding) {854 switch (request.head.transfer_encoding) {
767 .chunked => {855 .chunked => {
768 request.reader_state = .{ .chunk_parser = http.ChunkParser.init };856 request.reader_state = .{ .remaining_chunk_len = .head };
769 return .{857 return .{
770 .context = request,858 .context = request,
771 .vtable = &.{859 .vtable = &.{
772 .read = &chunkedReader_read,860 .read = &chunkedRead,
773 .readVec = &chunkedReader_readVec,861 .readVec = &chunkedReadVec,
774 .discard = &chunkedReader_discard,862 .discard = &chunkedDiscard,
775 },863 },
776 };864 };
777 },865 },
...@@ -782,9 +870,9 @@ pub const Request = struct {...@@ -782,9 +870,9 @@ pub const Request = struct {
782 return .{870 return .{
783 .context = request,871 .context = request,
784 .vtable = &.{872 .vtable = &.{
785 .read = &contentLengthReader_read,873 .read = &contentLengthRead,
786 .readVec = &contentLengthReader_readVec,874 .readVec = &contentLengthReadVec,
787 .discard = &contentLengthReader_discard,875 .discard = &contentLengthDiscard,
788 },876 },
789 };877 };
790 },878 },
...@@ -822,6 +910,11 @@ pub const Request = struct {...@@ -822,6 +910,11 @@ pub const Request = struct {
822 }910 }
823 return false;911 return false;
824 }912 }
913
914 fn failRead(r: *Request, err: ReadError) error{ReadFailed} {
915 r.read_err = err;
916 return error.ReadFailed;
917 }
825};918};
826919
827pub const Response = struct {920pub const Response = struct {
...@@ -1165,14 +1258,3 @@ pub const Response = struct {...@@ -1165,14 +1258,3 @@ pub const Response = struct {
1165 };1258 };
1166 }1259 }
1167};1260};
1168
1169fn rebase(s: *Server, index: usize) void {
1170 const leftover = s.read_buffer[s.next_request_start..s.read_buffer_len];
1171 const dest = s.read_buffer[index..][0..leftover.len];
1172 if (leftover.len <= s.next_request_start - index) {
1173 @memcpy(dest, leftover);
1174 } else {
1175 mem.copyBackwards(u8, dest, leftover);
1176 }
1177 s.read_buffer_len = index + leftover.len;
1178}
lib/std/io/BufferedReader.zig+69-23
...@@ -43,14 +43,26 @@ pub fn reader(br: *BufferedReader) Reader {...@@ -43,14 +43,26 @@ pub fn reader(br: *BufferedReader) Reader {
43 .vtable = &.{43 .vtable = &.{
44 .read = passthruRead,44 .read = passthruRead,
45 .readVec = passthruReadVec,45 .readVec = passthruReadVec,
46 .discard = passthruDiscard,
46 },47 },
47 };48 };
48}49}
4950
51/// Equivalent semantics to `std.io.Reader.VTable.readVec`.
50pub fn readVec(br: *BufferedReader, data: []const []u8) Reader.Error!usize {52pub fn readVec(br: *BufferedReader, data: []const []u8) Reader.Error!usize {
51 return passthruReadVec(br, data);53 return passthruReadVec(br, data);
52}54}
5355
56/// Equivalent semantics to `std.io.Reader.VTable.read`.
57pub fn read(br: *BufferedReader, bw: *BufferedWriter, limit: Reader.Limit) Reader.RwError!usize {
58 return passthruRead(br, bw, limit);
59}
60
61/// Equivalent semantics to `std.io.Reader.VTable.discard`.
62pub fn discard(br: *BufferedReader, limit: Reader.Limit) Reader.Error!usize {
63 return passthruDiscard(br, limit);
64}
65
54pub fn readVecAll(br: *BufferedReader, data: [][]u8) Reader.Error!void {66pub fn readVecAll(br: *BufferedReader, data: [][]u8) Reader.Error!void {
55 var index: usize = 0;67 var index: usize = 0;
56 var truncate: usize = 0;68 var truncate: usize = 0;
...@@ -68,10 +80,6 @@ pub fn readVecAll(br: *BufferedReader, data: [][]u8) Reader.Error!void {...@@ -68,10 +80,6 @@ pub fn readVecAll(br: *BufferedReader, data: [][]u8) Reader.Error!void {
68 }80 }
69}81}
7082
71pub fn read(br: *BufferedReader, bw: *BufferedWriter, limit: Reader.Limit) Reader.RwError!usize {
72 return passthruRead(br, bw, limit);
73}
74
75/// "Pump" data from the reader to the writer.83/// "Pump" data from the reader to the writer.
76pub fn readAll(br: *BufferedReader, bw: *BufferedWriter, limit: Reader.Limit) Reader.RwError!void {84pub fn readAll(br: *BufferedReader, bw: *BufferedWriter, limit: Reader.Limit) Reader.RwError!void {
77 var remaining = limit;85 var remaining = limit;
...@@ -81,21 +89,46 @@ pub fn readAll(br: *BufferedReader, bw: *BufferedWriter, limit: Reader.Limit) Re...@@ -81,21 +89,46 @@ pub fn readAll(br: *BufferedReader, bw: *BufferedWriter, limit: Reader.Limit) Re
81 }89 }
82}90}
8391
84fn passthruRead(ctx: ?*anyopaque, bw: *BufferedWriter, limit: Reader.Limit) Reader.RwError!usize {92/// Equivalent to `readVec` but reads at most `limit` bytes.
85 const br: *BufferedReader = @alignCast(@ptrCast(ctx));93pub fn readVecLimit(br: *BufferedReader, data: []const []u8, limit: Reader.Limit) Reader.Error!usize {
86 const buffer = br.buffer[0..br.end];94 _ = br;
87 const buffered = buffer[br.seek..];95 _ = data;
88 const limited = buffered[0..limit.minInt(buffered.len)];96 _ = limit;
89 if (limited.len > 0) {97 @panic("TODO");
90 const n = try bw.write(limited);98}
99
100fn passthruRead(context: ?*anyopaque, bw: *BufferedWriter, limit: Reader.Limit) Reader.RwError!usize {
101 const br: *BufferedReader = @alignCast(@ptrCast(context));
102 const buffer = limit.slice(br.buffer[br.end..br.seek]);
103 if (buffer.len > 0) {
104 const n = try bw.write(buffer);
91 br.seek += n;105 br.seek += n;
92 return n;106 return n;
93 }107 }
94 return br.unbuffered_reader.read(bw, limit);108 return br.unbuffered_reader.read(bw, limit);
95}109}
96110
97fn passthruReadVec(ctx: ?*anyopaque, data: []const []u8) Reader.Error!usize {111fn passthruDiscard(context: ?*anyopaque, limit: Reader.Limit) Reader.Error!usize {
98 const br: *BufferedReader = @alignCast(@ptrCast(ctx));112 const br: *BufferedReader = @alignCast(@ptrCast(context));
113 const buffered_len = br.end - br.seek;
114 if (limit.toInt()) |n| {
115 if (buffered_len >= n) {
116 br.seek += n;
117 return n;
118 }
119 br.seek = 0;
120 br.end = 0;
121 const additional = try br.unbuffered_reader.discard(.limited(n - buffered_len));
122 return n + additional;
123 }
124 const n = try br.unbuffered_reader.discard(.unlimited);
125 br.seek = 0;
126 br.end = 0;
127 return buffered_len + n;
128}
129
130fn passthruReadVec(context: ?*anyopaque, data: []const []u8) Reader.Error!usize {
131 const br: *BufferedReader = @alignCast(@ptrCast(context));
99 var total: usize = 0;132 var total: usize = 0;
100 for (data, 0..) |buf, i| {133 for (data, 0..) |buf, i| {
101 const buffered = br.buffer[br.seek..br.end];134 const buffered = br.buffer[br.seek..br.end];
...@@ -171,14 +204,14 @@ pub fn peek(br: *BufferedReader, n: usize) Reader.Error![]u8 {...@@ -171,14 +204,14 @@ pub fn peek(br: *BufferedReader, n: usize) Reader.Error![]u8 {
171}204}
172205
173/// Returns all the next buffered bytes from `unbuffered_reader`, after filling206/// Returns all the next buffered bytes from `unbuffered_reader`, after filling
174/// the buffer to ensure it contains at least `min_len` bytes.207/// the buffer to ensure it contains at least `n` bytes.
175///208///
176/// Invalidates previously returned values from `peek` and `peekGreedy`.209/// Invalidates previously returned values from `peek` and `peekGreedy`.
177///210///
178/// Asserts that the `BufferedReader` was initialized with a buffer capacity at211/// Asserts that the `BufferedReader` was initialized with a buffer capacity at
179/// least as big as `min_len`.212/// least as big as `n`.
180///213///
181/// If there are fewer than `min_len` bytes left in the stream, `error.EndOfStream`214/// If there are fewer than `n` bytes left in the stream, `error.EndOfStream`
182/// is returned instead.215/// is returned instead.
183///216///
184/// See also:217/// See also:
...@@ -253,7 +286,8 @@ pub fn peekArray(br: *BufferedReader, comptime n: usize) Reader.Error!*[n]u8 {...@@ -253,7 +286,8 @@ pub fn peekArray(br: *BufferedReader, comptime n: usize) Reader.Error!*[n]u8 {
253/// * `toss`286/// * `toss`
254/// * `discardRemaining`287/// * `discardRemaining`
255/// * `discardShort`288/// * `discardShort`
256pub fn discard(br: *BufferedReader, n: usize) Reader.Error!void {289/// * `discard`
290pub fn discardAll(br: *BufferedReader, n: usize) Reader.Error!void {
257 if ((try br.discardShort(n)) != n) return error.EndOfStream;291 if ((try br.discardShort(n)) != n) return error.EndOfStream;
258}292}
259293
...@@ -265,9 +299,9 @@ pub fn discard(br: *BufferedReader, n: usize) Reader.Error!void {...@@ -265,9 +299,9 @@ pub fn discard(br: *BufferedReader, n: usize) Reader.Error!void {
265/// if the stream reached the end.299/// if the stream reached the end.
266///300///
267/// See also:301/// See also:
268/// * `discard`302/// * `discardAll`
269/// * `toss`
270/// * `discardRemaining`303/// * `discardRemaining`
304/// * `discard`
271pub fn discardShort(br: *BufferedReader, n: usize) Reader.ShortError!usize {305pub fn discardShort(br: *BufferedReader, n: usize) Reader.ShortError!usize {
272 const proposed_seek = br.seek + n;306 const proposed_seek = br.seek + n;
273 if (proposed_seek <= br.end) {307 if (proposed_seek <= br.end) {
...@@ -609,18 +643,30 @@ pub fn fill(br: *BufferedReader, n: usize) Reader.Error!void {...@@ -609,18 +643,30 @@ pub fn fill(br: *BufferedReader, n: usize) Reader.Error!void {
609 }643 }
610}644}
611645
612/// Reads 1 byte from the stream or returns `error.EndOfStream`.646/// Returns the next byte from the stream or returns `error.EndOfStream`.
613pub fn takeByte(br: *BufferedReader) Reader.Error!u8 {647///
648/// Does not advance the seek position.
649///
650/// Asserts the buffer capacity is nonzero.
651pub fn peekByte(br: *BufferedReader) Reader.Error!u8 {
614 const buffer = br.buffer[0..br.end];652 const buffer = br.buffer[0..br.end];
615 const seek = br.seek;653 const seek = br.seek;
616 if (seek >= buffer.len) {654 if (seek >= buffer.len) {
617 @branchHint(.unlikely);655 @branchHint(.unlikely);
618 try fill(br, 1);656 try fill(br, 1);
619 }657 }
620 br.seek = seek + 1;
621 return buffer[seek];658 return buffer[seek];
622}659}
623660
661/// Reads 1 byte from the stream or returns `error.EndOfStream`.
662///
663/// Asserts the buffer capacity is nonzero.
664pub fn takeByte(br: *BufferedReader) Reader.Error!u8 {
665 const result = try peekByte(br);
666 br.seek += 1;
667 return result;
668}
669
624/// Same as `takeByte` except the returned byte is signed.670/// Same as `takeByte` except the returned byte is signed.
625pub fn takeByteSigned(br: *BufferedReader) Reader.Error!i8 {671pub fn takeByteSigned(br: *BufferedReader) Reader.Error!i8 {
626 return @bitCast(try br.takeByte());672 return @bitCast(try br.takeByte());
...@@ -813,7 +859,7 @@ test peekArray {...@@ -813,7 +859,7 @@ test peekArray {
813 return error.Unimplemented;859 return error.Unimplemented;
814}860}
815861
816test discard {862test discardAll {
817 var br: BufferedReader = undefined;863 var br: BufferedReader = undefined;
818 br.initFixed("foobar");864 br.initFixed("foobar");
819 try br.discard(3);865 try br.discard(3);
lib/std/io/BufferedWriter.zig+2-2
...@@ -558,8 +558,8 @@ fn passthruWriteFile(...@@ -558,8 +558,8 @@ fn passthruWriteFile(
558 const remaining_buffers = buffers[1..];558 const remaining_buffers = buffers[1..];
559 const send_trailers_len: usize = @min(trailers.len, remaining_buffers.len);559 const send_trailers_len: usize = @min(trailers.len, remaining_buffers.len);
560 @memcpy(remaining_buffers[0..send_trailers_len], trailers[0..send_trailers_len]);560 @memcpy(remaining_buffers[0..send_trailers_len], trailers[0..send_trailers_len]);
561 const send_headers_len = 1;561 const send_headers_len = @intFromBool(end != 0);
562 const send_buffers = buffers[0 .. send_headers_len + send_trailers_len];562 const send_buffers = buffers[1 - send_headers_len .. 1 + send_trailers_len];
563 const n = try bw.unbuffered_writer.writeFile(file, offset, limit, send_buffers, send_headers_len);563 const n = try bw.unbuffered_writer.writeFile(file, offset, limit, send_buffers, send_headers_len);
564 if (n < end) {564 if (n < end) {
565 @branchHint(.unlikely);565 @branchHint(.unlikely);
lib/std/io/Reader.zig+6-2
...@@ -126,7 +126,9 @@ pub const Limit = enum(usize) {...@@ -126,7 +126,9 @@ pub const Limit = enum(usize) {
126};126};
127127
128pub fn read(r: Reader, bw: *BufferedWriter, limit: Limit) RwError!usize {128pub fn read(r: Reader, bw: *BufferedWriter, limit: Limit) RwError!usize {
129 return r.vtable.read(r.context, bw, limit);129 const n = try r.vtable.read(r.context, bw, limit);
130 assert(n <= @intFromEnum(limit));
131 return n;
130}132}
131133
132pub fn readVec(r: Reader, data: []const []u8) Error!usize {134pub fn readVec(r: Reader, data: []const []u8) Error!usize {
...@@ -134,7 +136,9 @@ pub fn readVec(r: Reader, data: []const []u8) Error!usize {...@@ -134,7 +136,9 @@ pub fn readVec(r: Reader, data: []const []u8) Error!usize {
134}136}
135137
136pub fn discard(r: Reader, limit: Limit) Error!usize {138pub fn discard(r: Reader, limit: Limit) Error!usize {
137 return r.vtable.discard(r.context, limit);139 const n = try r.vtable.discard(r.context, limit);
140 assert(n <= @intFromEnum(limit));
141 return n;
138}142}
139143
140/// Returns total number of bytes written to `bw`.144/// Returns total number of bytes written to `bw`.