| author | |
| committer | |
| log | 9f8486170c5a5dda56381098cf10ecf8c6364d9b |
| tree | 0e53d8f0eac7314e56733eb24162debf9a89547c |
| parent | d4545f216a6ea86caec208cabc724004b7ef320a |
13 files changed, 290 insertions(+), 237 deletions(-)
lib/std/Build.zig+1-1| ... | @@ -2061,7 +2061,7 @@ pub fn runAllowFail( | ... | @@ -2061,7 +2061,7 @@ pub fn runAllowFail( |
| 2061 | try child.spawn(); | 2061 | try child.spawn(); |
| 2062 | 2062 | ||
| 2063 | var file_reader = child.stdout.?.readerStreaming(); | 2063 | var file_reader = child.stdout.?.readerStreaming(); |
| 2064 | const stdout = try file_reader.interface().readRemainingAlloc(b.allocator, .limited(max_output_size)); | 2064 | const stdout = try file_reader.interface().allocRemaining(b.allocator, .limited(max_output_size)); |
| 2065 | errdefer b.allocator.free(stdout); | 2065 | errdefer b.allocator.free(stdout); |
| 2066 | 2066 | ||
| 2067 | const term = try child.wait(); | 2067 | const term = try child.wait(); |
lib/std/Build/Cache.zig+1-1| ... | @@ -663,7 +663,7 @@ pub const Manifest = struct { | ... | @@ -663,7 +663,7 @@ pub const Manifest = struct { |
| 663 | const input_file_count = self.files.entries.len; | 663 | const input_file_count = self.files.entries.len; |
| 664 | var manifest_reader = self.manifest_file.?.reader(); // Reads positionally from zero. | 664 | var manifest_reader = self.manifest_file.?.reader(); // Reads positionally from zero. |
| 665 | const limit: std.io.Limit = .limited(manifest_file_size_max); | 665 | const limit: std.io.Limit = .limited(manifest_file_size_max); |
| 666 | const file_contents = manifest_reader.interface().readRemainingAlloc(gpa, limit) catch |err| switch (err) { | 666 | const file_contents = manifest_reader.interface().allocRemaining(gpa, limit) catch |err| switch (err) { |
| 667 | error.OutOfMemory => return error.OutOfMemory, | 667 | error.OutOfMemory => return error.OutOfMemory, |
| 668 | error.StreamTooLong => return error.OutOfMemory, | 668 | error.StreamTooLong => return error.OutOfMemory, |
| 669 | error.ReadFailed => { | 669 | error.ReadFailed => { |
lib/std/Build/Step/Run.zig+2-2| ... | @@ -1788,7 +1788,7 @@ fn evalGeneric(run: *Run, child: *std.process.Child) !StdIoResult { | ... | @@ -1788,7 +1788,7 @@ fn evalGeneric(run: *Run, child: *std.process.Child) !StdIoResult { |
| 1788 | stderr_bytes = poller.reader(.stderr).buffered(); | 1788 | stderr_bytes = poller.reader(.stderr).buffered(); |
| 1789 | } else { | 1789 | } else { |
| 1790 | var fr = stdout.readerStreaming(); | 1790 | var fr = stdout.readerStreaming(); |
| 1791 | stdout_bytes = fr.interface().readRemainingAlloc(arena, run.stdio_limit) catch |err| switch (err) { | 1791 | stdout_bytes = fr.interface().allocRemaining(arena, run.stdio_limit) catch |err| switch (err) { |
| 1792 | error.OutOfMemory => return error.OutOfMemory, | 1792 | error.OutOfMemory => return error.OutOfMemory, |
| 1793 | error.ReadFailed => return fr.err.?, | 1793 | error.ReadFailed => return fr.err.?, |
| 1794 | error.StreamTooLong => return error.StdoutStreamTooLong, | 1794 | error.StreamTooLong => return error.StdoutStreamTooLong, |
| ... | @@ -1796,7 +1796,7 @@ fn evalGeneric(run: *Run, child: *std.process.Child) !StdIoResult { | ... | @@ -1796,7 +1796,7 @@ fn evalGeneric(run: *Run, child: *std.process.Child) !StdIoResult { |
| 1796 | } | 1796 | } |
| 1797 | } else if (child.stderr) |stderr| { | 1797 | } else if (child.stderr) |stderr| { |
| 1798 | var fr = stderr.readerStreaming(); | 1798 | var fr = stderr.readerStreaming(); |
| 1799 | stderr_bytes = fr.interface().readRemainingAlloc(arena, run.stdio_limit) catch |err| switch (err) { | 1799 | stderr_bytes = fr.interface().allocRemaining(arena, run.stdio_limit) catch |err| switch (err) { |
| 1800 | error.OutOfMemory => return error.OutOfMemory, | 1800 | error.OutOfMemory => return error.OutOfMemory, |
| 1801 | error.ReadFailed => return fr.err.?, | 1801 | error.ReadFailed => return fr.err.?, |
| 1802 | error.StreamTooLong => return error.StderrStreamTooLong, | 1802 | error.StreamTooLong => return error.StderrStreamTooLong, |
lib/std/Uri.zig+2-2| ... | @@ -465,11 +465,11 @@ fn merge_paths(base: Component, new: []u8, aux_buf: *[]u8) error{NoSpaceLeft}!Co | ... | @@ -465,11 +465,11 @@ fn merge_paths(base: Component, new: []u8, aux_buf: *[]u8) error{NoSpaceLeft}!Co |
| 465 | var aux: Writer = .fixed(aux_buf.*); | 465 | var aux: Writer = .fixed(aux_buf.*); |
| 466 | if (!base.isEmpty()) { | 466 | if (!base.isEmpty()) { |
| 467 | aux.print("{fpath}", .{base}) catch return error.NoSpaceLeft; | 467 | aux.print("{fpath}", .{base}) catch return error.NoSpaceLeft; |
| 468 | aux.end = std.mem.lastIndexOfScalar(u8, aux.getWritten(), '/') orelse | 468 | aux.end = std.mem.lastIndexOfScalar(u8, aux.buffered(), '/') orelse |
| 469 | return remove_dot_segments(new); | 469 | return remove_dot_segments(new); |
| 470 | } | 470 | } |
| 471 | aux.print("/{s}", .{new}) catch return error.NoSpaceLeft; | 471 | aux.print("/{s}", .{new}) catch return error.NoSpaceLeft; |
| 472 | const merged_path = remove_dot_segments(aux.getWritten()); | 472 | const merged_path = remove_dot_segments(aux.buffered()); |
| 473 | aux_buf.* = aux_buf.*[merged_path.percent_encoded.len..]; | 473 | aux_buf.* = aux_buf.*[merged_path.percent_encoded.len..]; |
| 474 | return merged_path; | 474 | return merged_path; |
| 475 | } | 475 | } |
lib/std/compress/flate/Decompress.zig+6-6| ... | @@ -8,7 +8,7 @@ const Writer = std.io.Writer; | ... | @@ -8,7 +8,7 @@ const Writer = std.io.Writer; |
| 8 | const Reader = std.io.Reader; | 8 | const Reader = std.io.Reader; |
| 9 | 9 | ||
| 10 | input: *Reader, | 10 | input: *Reader, |
| 11 | interface: Reader, | 11 | reader: Reader, |
| 12 | /// Hashes, produces checksum, of uncompressed data for gzip/zlib footer. | 12 | /// Hashes, produces checksum, of uncompressed data for gzip/zlib footer. |
| 13 | hasher: Container.Hasher, | 13 | hasher: Container.Hasher, |
| 14 | 14 | ||
| ... | @@ -51,7 +51,7 @@ pub const Error = Container.Error || error{ | ... | @@ -51,7 +51,7 @@ pub const Error = Container.Error || error{ |
| 51 | 51 | ||
| 52 | pub fn init(input: *Reader, container: Container, buffer: []u8) Decompress { | 52 | pub fn init(input: *Reader, container: Container, buffer: []u8) Decompress { |
| 53 | return .{ | 53 | return .{ |
| 54 | .interface = .{ | 54 | .reader = .{ |
| 55 | // TODO populate discard so that when an amount is discarded that | 55 | // TODO populate discard so that when an amount is discarded that |
| 56 | // includes an entire frame, skip decoding that frame. | 56 | // includes an entire frame, skip decoding that frame. |
| 57 | .vtable = &.{ .stream = stream }, | 57 | .vtable = &.{ .stream = stream }, |
| ... | @@ -130,7 +130,7 @@ fn decodeSymbol(self: *Decompress, decoder: anytype) !Symbol { | ... | @@ -130,7 +130,7 @@ fn decodeSymbol(self: *Decompress, decoder: anytype) !Symbol { |
| 130 | } | 130 | } |
| 131 | 131 | ||
| 132 | pub fn stream(r: *Reader, w: *Writer, limit: std.io.Limit) Reader.StreamError!usize { | 132 | pub fn stream(r: *Reader, w: *Writer, limit: std.io.Limit) Reader.StreamError!usize { |
| 133 | const d: *Decompress = @alignCast(@fieldParentPtr("interface", r)); | 133 | const d: *Decompress = @alignCast(@fieldParentPtr("reader", r)); |
| 134 | return readInner(d, w, limit) catch |err| switch (err) { | 134 | return readInner(d, w, limit) catch |err| switch (err) { |
| 135 | error.EndOfStream => return error.EndOfStream, | 135 | error.EndOfStream => return error.EndOfStream, |
| 136 | error.WriteFailed => return error.WriteFailed, | 136 | error.WriteFailed => return error.WriteFailed, |
| ... | @@ -247,7 +247,7 @@ fn readInner(d: *Decompress, w: *Writer, limit: std.io.Limit) (Error || Reader.S | ... | @@ -247,7 +247,7 @@ fn readInner(d: *Decompress, w: *Writer, limit: std.io.Limit) (Error || Reader.S |
| 247 | } | 247 | } |
| 248 | }, | 248 | }, |
| 249 | .stored_block => |remaining_len| { | 249 | .stored_block => |remaining_len| { |
| 250 | const out = try w.writableSliceGreedyPreserving(flate.history_len, 1); | 250 | const out = try w.writableSliceGreedyPreserve(flate.history_len, 1); |
| 251 | const limited_out = limit.min(.limited(remaining_len)).slice(out); | 251 | const limited_out = limit.min(.limited(remaining_len)).slice(out); |
| 252 | const n = try d.input.readVec(&.{limited_out}); | 252 | const n = try d.input.readVec(&.{limited_out}); |
| 253 | if (remaining_len - n == 0) { | 253 | if (remaining_len - n == 0) { |
| ... | @@ -263,7 +263,7 @@ fn readInner(d: *Decompress, w: *Writer, limit: std.io.Limit) (Error || Reader.S | ... | @@ -263,7 +263,7 @@ fn readInner(d: *Decompress, w: *Writer, limit: std.io.Limit) (Error || Reader.S |
| 263 | while (@intFromEnum(limit) > w.count - start) { | 263 | while (@intFromEnum(limit) > w.count - start) { |
| 264 | const code = try d.readFixedCode(); | 264 | const code = try d.readFixedCode(); |
| 265 | switch (code) { | 265 | switch (code) { |
| 266 | 0...255 => try w.writeBytePreserving(flate.history_len, @intCast(code)), | 266 | 0...255 => try w.writeBytePreserve(flate.history_len, @intCast(code)), |
| 267 | 256 => { | 267 | 256 => { |
| 268 | d.state = if (d.final_block) .protocol_footer else .block_header; | 268 | d.state = if (d.final_block) .protocol_footer else .block_header; |
| 269 | return w.count - start; | 269 | return w.count - start; |
| ... | @@ -289,7 +289,7 @@ fn readInner(d: *Decompress, w: *Writer, limit: std.io.Limit) (Error || Reader.S | ... | @@ -289,7 +289,7 @@ fn readInner(d: *Decompress, w: *Writer, limit: std.io.Limit) (Error || Reader.S |
| 289 | const sym = try d.decodeSymbol(&d.lit_dec); | 289 | const sym = try d.decodeSymbol(&d.lit_dec); |
| 290 | 290 | ||
| 291 | switch (sym.kind) { | 291 | switch (sym.kind) { |
| 292 | .literal => try w.writeBytePreserving(flate.history_len, sym.symbol), | 292 | .literal => try w.writeBytePreserve(flate.history_len, sym.symbol), |
| 293 | .match => { | 293 | .match => { |
| 294 | // Decode match backreference <length, distance> | 294 | // Decode match backreference <length, distance> |
| 295 | const length = try d.decodeLength(sym.symbol); | 295 | const length = try d.decodeLength(sym.symbol); |
lib/std/compress/zstd/Decompress.zig+3-3| ... | @@ -7,7 +7,7 @@ const zstd = @import("../zstd.zig"); | ... | @@ -7,7 +7,7 @@ const zstd = @import("../zstd.zig"); |
| 7 | const Writer = std.io.Writer; | 7 | const Writer = std.io.Writer; |
| 8 | 8 | ||
| 9 | input: *Reader, | 9 | input: *Reader, |
| 10 | interface: Reader, | 10 | reader: Reader, |
| 11 | state: State, | 11 | state: State, |
| 12 | verify_checksum: bool, | 12 | verify_checksum: bool, |
| 13 | err: ?Error = null, | 13 | err: ?Error = null, |
| ... | @@ -68,7 +68,7 @@ pub fn init(input: *Reader, buffer: []u8, options: Options) Decompress { | ... | @@ -68,7 +68,7 @@ pub fn init(input: *Reader, buffer: []u8, options: Options) Decompress { |
| 68 | .input = input, | 68 | .input = input, |
| 69 | .state = .new_frame, | 69 | .state = .new_frame, |
| 70 | .verify_checksum = options.verify_checksum, | 70 | .verify_checksum = options.verify_checksum, |
| 71 | .interface = .{ | 71 | .reader = .{ |
| 72 | .vtable = &.{ .stream = stream }, | 72 | .vtable = &.{ .stream = stream }, |
| 73 | .buffer = buffer, | 73 | .buffer = buffer, |
| 74 | .seek = 0, | 74 | .seek = 0, |
| ... | @@ -78,7 +78,7 @@ pub fn init(input: *Reader, buffer: []u8, options: Options) Decompress { | ... | @@ -78,7 +78,7 @@ pub fn init(input: *Reader, buffer: []u8, options: Options) Decompress { |
| 78 | } | 78 | } |
| 79 | 79 | ||
| 80 | fn stream(r: *Reader, w: *Writer, limit: Limit) Reader.StreamError!usize { | 80 | fn stream(r: *Reader, w: *Writer, limit: Limit) Reader.StreamError!usize { |
| 81 | const d: *Decompress = @alignCast(@fieldParentPtr("interface", r)); | 81 | const d: *Decompress = @alignCast(@fieldParentPtr("reader", r)); |
| 82 | const in = d.input; | 82 | const in = d.input; |
| 83 | 83 | ||
| 84 | switch (d.state) { | 84 | switch (d.state) { |
lib/std/debug/Dwarf.zig+1-1| ... | @@ -2242,7 +2242,7 @@ pub const ElfModule = struct { | ... | @@ -2242,7 +2242,7 @@ pub const ElfModule = struct { |
| 2242 | 2242 | ||
| 2243 | var zlib_stream: std.compress.flate.Decompress = .init(&section_reader, .zlib, &.{}); | 2243 | var zlib_stream: std.compress.flate.Decompress = .init(&section_reader, .zlib, &.{}); |
| 2244 | 2244 | ||
| 2245 | const decompressed_section = zlib_stream.interface.allocRemaining(gpa, .limited(ch_size)) catch | 2245 | const decompressed_section = zlib_stream.reader.allocRemaining(gpa, .limited(ch_size)) catch |
| 2246 | continue; | 2246 | continue; |
| 2247 | if (decompressed_section.len != ch_size) { | 2247 | if (decompressed_section.len != ch_size) { |
| 2248 | gpa.free(decompressed_section); | 2248 | gpa.free(decompressed_section); |
lib/std/http.zig+49-55| ... | @@ -328,6 +328,9 @@ pub const Header = struct { | ... | @@ -328,6 +328,9 @@ pub const Header = struct { |
| 328 | 328 | ||
| 329 | pub const Reader = struct { | 329 | pub const Reader = struct { |
| 330 | in: *std.io.Reader, | 330 | in: *std.io.Reader, |
| 331 | /// This is preallocated memory that might be used by `bodyReader`. That | ||
| 332 | /// function might return a pointer to this field, or a different | ||
| 333 | /// `*std.io.Reader`. Advisable to not access this field directly. | ||
| 331 | interface: std.io.Reader, | 334 | interface: std.io.Reader, |
| 332 | /// Keeps track of whether the stream is ready to accept a new request, | 335 | /// Keeps track of whether the stream is ready to accept a new request, |
| 333 | /// making invalid API usage cause assertion failures rather than HTTP | 336 | /// making invalid API usage cause assertion failures rather than HTTP |
| ... | @@ -489,31 +492,31 @@ pub const Reader = struct { | ... | @@ -489,31 +492,31 @@ pub const Reader = struct { |
| 489 | content_encoding: ContentEncoding, | 492 | content_encoding: ContentEncoding, |
| 490 | decompressor: *Decompressor, | 493 | decompressor: *Decompressor, |
| 491 | decompression_buffer: []u8, | 494 | decompression_buffer: []u8, |
| 492 | ) std.io.Reader { | 495 | ) *std.io.Reader { |
| 493 | if (transfer_encoding == .none and content_length == null) { | 496 | if (transfer_encoding == .none and content_length == null) { |
| 494 | assert(reader.state == .received_head); | 497 | assert(reader.state == .received_head); |
| 495 | reader.state = .body_none; | 498 | reader.state = .body_none; |
| 496 | switch (content_encoding) { | 499 | switch (content_encoding) { |
| 497 | .identity => { | 500 | .identity => { |
| 498 | return reader.in.reader(); | 501 | return reader.in; |
| 499 | }, | 502 | }, |
| 500 | .deflate => { | 503 | .deflate => { |
| 501 | decompressor.compression = .{ .deflate = .init(reader.in) }; | 504 | decompressor.* = .{ .flate = .init(reader.in, .raw, decompression_buffer) }; |
| 502 | return decompressor.compression.deflate.reader(); | 505 | return &decompressor.flate.reader; |
| 503 | }, | 506 | }, |
| 504 | .gzip => { | 507 | .gzip => { |
| 505 | decompressor.compression = .{ .gzip = .init(reader.in) }; | 508 | decompressor.* = .{ .flate = .init(reader.in, .gzip, decompression_buffer) }; |
| 506 | return decompressor.compression.gzip.reader(); | 509 | return &decompressor.flate.reader; |
| 507 | }, | 510 | }, |
| 508 | .zstd => { | 511 | .zstd => { |
| 509 | decompressor.compression = .{ .zstd = .init(reader.in, .{ .verify_checksum = false }) }; | 512 | decompressor.* = .{ .zstd = .init(reader.in, decompression_buffer, .{ .verify_checksum = false }) }; |
| 510 | return decompressor.compression.zstd.reader(); | 513 | return &decompressor.zstd.reader; |
| 511 | }, | 514 | }, |
| 512 | .compress => unreachable, | 515 | .compress => unreachable, |
| 513 | } | 516 | } |
| 514 | } | 517 | } |
| 515 | const transfer_reader = bodyReader(reader, transfer_encoding, content_length); | 518 | const transfer_reader = bodyReader(reader, &.{}, transfer_encoding, content_length); |
| 516 | return decompressor.reader(transfer_reader, decompression_buffer, content_encoding); | 519 | return decompressor.init(transfer_reader, decompression_buffer, content_encoding); |
| 517 | } | 520 | } |
| 518 | 521 | ||
| 519 | fn contentLengthStream( | 522 | fn contentLengthStream( |
| ... | @@ -711,42 +714,33 @@ pub const Reader = struct { | ... | @@ -711,42 +714,33 @@ pub const Reader = struct { |
| 711 | } | 714 | } |
| 712 | }; | 715 | }; |
| 713 | 716 | ||
| 714 | pub const Decompressor = struct { | 717 | pub const Decompressor = union(enum) { |
| 715 | compression: Compression, | 718 | flate: std.compress.flate.Decompress, |
| 716 | buffered_reader: std.io.Reader, | 719 | zstd: std.compress.zstd.Decompress, |
| 717 | 720 | none: *std.io.Reader, | |
| 718 | pub const Compression = union(enum) { | ||
| 719 | deflate: std.compress.flate.Decompressor, | ||
| 720 | gzip: std.compress.flate.Decompressor, | ||
| 721 | zstd: std.compress.zstd.Decompress, | ||
| 722 | none: void, | ||
| 723 | }; | ||
| 724 | 721 | ||
| 725 | pub fn reader( | 722 | pub fn init( |
| 726 | decompressor: *Decompressor, | 723 | decompressor: *Decompressor, |
| 727 | transfer_reader: std.io.Reader, | 724 | transfer_reader: *std.io.Reader, |
| 728 | buffer: []u8, | 725 | buffer: []u8, |
| 729 | content_encoding: ContentEncoding, | 726 | content_encoding: ContentEncoding, |
| 730 | ) std.io.Reader { | 727 | ) *std.io.Reader { |
| 731 | switch (content_encoding) { | 728 | switch (content_encoding) { |
| 732 | .identity => { | 729 | .identity => { |
| 733 | decompressor.compression = .none; | 730 | decompressor.* = .{ .none = transfer_reader }; |
| 734 | return transfer_reader; | 731 | return transfer_reader; |
| 735 | }, | 732 | }, |
| 736 | .deflate => { | 733 | .deflate => { |
| 737 | decompressor.buffered_reader = transfer_reader.buffered(buffer); | 734 | decompressor.* = .{ .flate = .init(transfer_reader, .raw, buffer) }; |
| 738 | decompressor.compression = .{ .deflate = .init(&decompressor.buffered_reader) }; | 735 | return &decompressor.flate.reader; |
| 739 | return decompressor.compression.deflate.reader(); | ||
| 740 | }, | 736 | }, |
| 741 | .gzip => { | 737 | .gzip => { |
| 742 | decompressor.buffered_reader = transfer_reader.buffered(buffer); | 738 | decompressor.* = .{ .flate = .init(transfer_reader, .gzip, buffer) }; |
| 743 | decompressor.compression = .{ .gzip = .init(&decompressor.buffered_reader) }; | 739 | return &decompressor.flate.reader; |
| 744 | return decompressor.compression.gzip.reader(); | ||
| 745 | }, | 740 | }, |
| 746 | .zstd => { | 741 | .zstd => { |
| 747 | decompressor.buffered_reader = transfer_reader.buffered(buffer); | 742 | decompressor.* = .{ .zstd = .init(transfer_reader, buffer, .{ .verify_checksum = false }) }; |
| 748 | decompressor.compression = .{ .zstd = .init(&decompressor.buffered_reader, .{}) }; | 743 | return &decompressor.zstd.reader; |
| 749 | return decompressor.compression.gzip.reader(); | ||
| 750 | }, | 744 | }, |
| 751 | .compress => unreachable, | 745 | .compress => unreachable, |
| 752 | } | 746 | } |
| ... | @@ -759,7 +753,7 @@ pub const BodyWriter = struct { | ... | @@ -759,7 +753,7 @@ pub const BodyWriter = struct { |
| 759 | /// state of this other than via methods of `BodyWriter`. | 753 | /// state of this other than via methods of `BodyWriter`. |
| 760 | http_protocol_output: *Writer, | 754 | http_protocol_output: *Writer, |
| 761 | state: State, | 755 | state: State, |
| 762 | interface: Writer, | 756 | writer: Writer, |
| 763 | 757 | ||
| 764 | pub const Error = Writer.Error; | 758 | pub const Error = Writer.Error; |
| 765 | 759 | ||
| ... | @@ -797,7 +791,7 @@ pub const BodyWriter = struct { | ... | @@ -797,7 +791,7 @@ pub const BodyWriter = struct { |
| 797 | }; | 791 | }; |
| 798 | 792 | ||
| 799 | pub fn isEliding(w: *const BodyWriter) bool { | 793 | pub fn isEliding(w: *const BodyWriter) bool { |
| 800 | return w.interface.vtable.drain == Writer.discardingDrain; | 794 | return w.writer.vtable.drain == Writer.discardingDrain; |
| 801 | } | 795 | } |
| 802 | 796 | ||
| 803 | /// Sends all buffered data across `BodyWriter.http_protocol_output`. | 797 | /// Sends all buffered data across `BodyWriter.http_protocol_output`. |
| ... | @@ -924,42 +918,42 @@ pub const BodyWriter = struct { | ... | @@ -924,42 +918,42 @@ pub const BodyWriter = struct { |
| 924 | w.state = .end; | 918 | w.state = .end; |
| 925 | } | 919 | } |
| 926 | 920 | ||
| 927 | fn contentLengthDrain(w: *Writer, data: []const []const u8, splat: usize) Error!usize { | 921 | pub fn contentLengthDrain(w: *Writer, data: []const []const u8, splat: usize) Error!usize { |
| 928 | const bw: *BodyWriter = @fieldParentPtr("interface", w); | 922 | const bw: *BodyWriter = @fieldParentPtr("writer", w); |
| 929 | assert(!bw.isEliding()); | 923 | assert(!bw.isEliding()); |
| 930 | const out = w.http_protocol_output; | 924 | const out = bw.http_protocol_output; |
| 931 | const n = try w.drainTo(out, data, splat); | 925 | const n = try w.drainTo(out, data, splat); |
| 932 | w.state.content_length -= n; | 926 | bw.state.content_length -= n; |
| 933 | return n; | 927 | return n; |
| 934 | } | 928 | } |
| 935 | 929 | ||
| 936 | fn noneDrain(w: *Writer, data: []const []const u8, splat: usize) Error!usize { | 930 | pub fn noneDrain(w: *Writer, data: []const []const u8, splat: usize) Error!usize { |
| 937 | const bw: *BodyWriter = @fieldParentPtr("interface", w); | 931 | const bw: *BodyWriter = @fieldParentPtr("writer", w); |
| 938 | assert(!bw.isEliding()); | 932 | assert(!bw.isEliding()); |
| 939 | const out = w.http_protocol_output; | 933 | const out = bw.http_protocol_output; |
| 940 | return try w.drainTo(out, data, splat); | 934 | return try w.drainTo(out, data, splat); |
| 941 | } | 935 | } |
| 942 | 936 | ||
| 943 | /// Returns `null` if size cannot be computed without making any syscalls. | 937 | /// Returns `null` if size cannot be computed without making any syscalls. |
| 944 | fn noneSendFile(w: *Writer, file_reader: *File.Reader, limit: std.io.Limit) Writer.FileError!usize { | 938 | pub fn noneSendFile(w: *Writer, file_reader: *File.Reader, limit: std.io.Limit) Writer.FileError!usize { |
| 945 | const bw: *BodyWriter = @fieldParentPtr("interface", w); | 939 | const bw: *BodyWriter = @fieldParentPtr("writer", w); |
| 946 | assert(!bw.isEliding()); | 940 | assert(!bw.isEliding()); |
| 947 | return w.sendFileTo(bw.http_protocol_output, file_reader, limit); | 941 | return w.sendFileTo(bw.http_protocol_output, file_reader, limit); |
| 948 | } | 942 | } |
| 949 | 943 | ||
| 950 | fn contentLengthSendFile(w: *Writer, file_reader: *File.Reader, limit: std.io.Limit) Writer.FileError!usize { | 944 | pub fn contentLengthSendFile(w: *Writer, file_reader: *File.Reader, limit: std.io.Limit) Writer.FileError!usize { |
| 951 | const bw: *BodyWriter = @fieldParentPtr("interface", w); | 945 | const bw: *BodyWriter = @fieldParentPtr("writer", w); |
| 952 | assert(!bw.isEliding()); | 946 | assert(!bw.isEliding()); |
| 953 | const n = try w.sendFileTo(bw.http_protocol_output, file_reader, limit); | 947 | const n = try w.sendFileTo(bw.http_protocol_output, file_reader, limit); |
| 954 | bw.state.content_length -= n; | 948 | bw.state.content_length -= n; |
| 955 | return n; | 949 | return n; |
| 956 | } | 950 | } |
| 957 | 951 | ||
| 958 | fn chunkedSendFile(w: *Writer, file_reader: *File.Reader, limit: std.io.Limit) Writer.FileError!usize { | 952 | pub fn chunkedSendFile(w: *Writer, file_reader: *File.Reader, limit: std.io.Limit) Writer.FileError!usize { |
| 959 | const bw: *BodyWriter = @fieldParentPtr("interface", w); | 953 | const bw: *BodyWriter = @fieldParentPtr("writer", w); |
| 960 | assert(!bw.isEliding()); | 954 | assert(!bw.isEliding()); |
| 961 | const data_len = w.countSendFileUpperBound(file_reader, limit) orelse { | 955 | const data_len = if (file_reader.getSize()) |x| w.end + x else |_| { |
| 962 | // If the file size is unknown, we cannot lower to a `writeFile` since we would | 956 | // If the file size is unknown, we cannot lower to a `sendFile` since we would |
| 963 | // have to flush the chunk header before knowing the chunk length. | 957 | // have to flush the chunk header before knowing the chunk length. |
| 964 | return error.Unimplemented; | 958 | return error.Unimplemented; |
| 965 | }; | 959 | }; |
| ... | @@ -1003,10 +997,10 @@ pub const BodyWriter = struct { | ... | @@ -1003,10 +997,10 @@ pub const BodyWriter = struct { |
| 1003 | } | 997 | } |
| 1004 | } | 998 | } |
| 1005 | 999 | ||
| 1006 | fn chunkedDrain(w: *Writer, data: []const []const u8, splat: usize) Error!usize { | 1000 | pub fn chunkedDrain(w: *Writer, data: []const []const u8, splat: usize) Error!usize { |
| 1007 | const bw: *BodyWriter = @fieldParentPtr("interface", w); | 1001 | const bw: *BodyWriter = @fieldParentPtr("writer", w); |
| 1008 | assert(!bw.isEliding()); | 1002 | assert(!bw.isEliding()); |
| 1009 | const out = w.http_protocol_output; | 1003 | const out = bw.http_protocol_output; |
| 1010 | const data_len = Writer.countSplat(w.end, data, splat); | 1004 | const data_len = Writer.countSplat(w.end, data, splat); |
| 1011 | const chunked = &bw.state.chunked; | 1005 | const chunked = &bw.state.chunked; |
| 1012 | state: switch (chunked.*) { | 1006 | state: switch (chunked.*) { |
| ... | @@ -1018,7 +1012,7 @@ pub const BodyWriter = struct { | ... | @@ -1018,7 +1012,7 @@ pub const BodyWriter = struct { |
| 1018 | const buffered_len = out.end - offset - chunk_header_template.len; | 1012 | const buffered_len = out.end - offset - chunk_header_template.len; |
| 1019 | const chunk_len = data_len + buffered_len; | 1013 | const chunk_len = data_len + buffered_len; |
| 1020 | writeHex(out.buffer[offset..][0..chunk_len_digits], chunk_len); | 1014 | writeHex(out.buffer[offset..][0..chunk_len_digits], chunk_len); |
| 1021 | const n = try w.drainTo(w, data, splat); | 1015 | const n = try w.drainTo(out, data, splat); |
| 1022 | chunked.* = .{ .chunk_len = data_len + 2 - n }; | 1016 | chunked.* = .{ .chunk_len = data_len + 2 - n }; |
| 1023 | return n; | 1017 | return n; |
| 1024 | }, | 1018 | }, |
| ... | @@ -1041,7 +1035,7 @@ pub const BodyWriter = struct { | ... | @@ -1041,7 +1035,7 @@ pub const BodyWriter = struct { |
| 1041 | continue :l 1; | 1035 | continue :l 1; |
| 1042 | }, | 1036 | }, |
| 1043 | else => { | 1037 | else => { |
| 1044 | const n = try w.drainToLimit(data, splat, .limited(chunk_len - 2)); | 1038 | const n = try w.drainToLimit(out, data, splat, .limited(chunk_len - 2)); |
| 1045 | chunked.chunk_len = chunk_len - n; | 1039 | chunked.chunk_len = chunk_len - n; |
| 1046 | return n; | 1040 | return n; |
| 1047 | }, | 1041 | }, |
lib/std/http/Client.zig+28-26| ... | @@ -408,7 +408,7 @@ pub const Connection = struct { | ... | @@ -408,7 +408,7 @@ pub const Connection = struct { |
| 408 | 408 | ||
| 409 | /// HTTP protocol from server to client. | 409 | /// HTTP protocol from server to client. |
| 410 | /// This either comes directly from `stream_reader`, or from a TLS client. | 410 | /// This either comes directly from `stream_reader`, or from a TLS client. |
| 411 | pub fn reader(c: *const Connection) *Reader { | 411 | pub fn reader(c: *Connection) *Reader { |
| 412 | return switch (c.protocol) { | 412 | return switch (c.protocol) { |
| 413 | .tls => { | 413 | .tls => { |
| 414 | if (disable_tls) unreachable; | 414 | if (disable_tls) unreachable; |
| ... | @@ -682,7 +682,7 @@ pub const Response = struct { | ... | @@ -682,7 +682,7 @@ pub const Response = struct { |
| 682 | /// | 682 | /// |
| 683 | /// See also: | 683 | /// See also: |
| 684 | /// * `readerDecompressing` | 684 | /// * `readerDecompressing` |
| 685 | pub fn reader(response: *Response, buffer: []u8) Reader { | 685 | pub fn reader(response: *Response, buffer: []u8) *Reader { |
| 686 | const req = response.request; | 686 | const req = response.request; |
| 687 | if (!req.method.responseHasBody()) return .ending; | 687 | if (!req.method.responseHasBody()) return .ending; |
| 688 | const head = &response.head; | 688 | const head = &response.head; |
| ... | @@ -702,7 +702,7 @@ pub const Response = struct { | ... | @@ -702,7 +702,7 @@ pub const Response = struct { |
| 702 | response: *Response, | 702 | response: *Response, |
| 703 | decompressor: *http.Decompressor, | 703 | decompressor: *http.Decompressor, |
| 704 | decompression_buffer: []u8, | 704 | decompression_buffer: []u8, |
| 705 | ) Reader { | 705 | ) *Reader { |
| 706 | const head = &response.head; | 706 | const head = &response.head; |
| 707 | return response.request.reader.bodyReaderDecompressing( | 707 | return response.request.reader.bodyReaderDecompressing( |
| 708 | head.transfer_encoding, | 708 | head.transfer_encoding, |
| ... | @@ -864,14 +864,14 @@ pub const Request = struct { | ... | @@ -864,14 +864,14 @@ pub const Request = struct { |
| 864 | pub fn sendBodyUnflushed(r: *Request, buffer: []u8) Writer.Error!http.BodyWriter { | 864 | pub fn sendBodyUnflushed(r: *Request, buffer: []u8) Writer.Error!http.BodyWriter { |
| 865 | assert(r.method.requestHasBody()); | 865 | assert(r.method.requestHasBody()); |
| 866 | try sendHead(r); | 866 | try sendHead(r); |
| 867 | const http_protocol_output = &r.connection.?.writer; | 867 | const http_protocol_output = r.connection.?.writer(); |
| 868 | return switch (r.transfer_encoding) { | 868 | return switch (r.transfer_encoding) { |
| 869 | .chunked => .{ | 869 | .chunked => .{ |
| 870 | .http_protocol_output = http_protocol_output, | 870 | .http_protocol_output = http_protocol_output, |
| 871 | .state = .{ .chunked = .init }, | 871 | .state = .{ .chunked = .init }, |
| 872 | .interface = .{ | 872 | .writer = .{ |
| 873 | .buffer = buffer, | 873 | .buffer = buffer, |
| 874 | .interface = &.{ | 874 | .vtable = &.{ |
| 875 | .drain = http.BodyWriter.chunkedDrain, | 875 | .drain = http.BodyWriter.chunkedDrain, |
| 876 | .sendFile = http.BodyWriter.chunkedSendFile, | 876 | .sendFile = http.BodyWriter.chunkedSendFile, |
| 877 | }, | 877 | }, |
| ... | @@ -880,9 +880,9 @@ pub const Request = struct { | ... | @@ -880,9 +880,9 @@ pub const Request = struct { |
| 880 | .content_length => |len| .{ | 880 | .content_length => |len| .{ |
| 881 | .http_protocol_output = http_protocol_output, | 881 | .http_protocol_output = http_protocol_output, |
| 882 | .state = .{ .content_length = len }, | 882 | .state = .{ .content_length = len }, |
| 883 | .interface = .{ | 883 | .writer = .{ |
| 884 | .buffer = buffer, | 884 | .buffer = buffer, |
| 885 | .interface = &.{ | 885 | .vtable = &.{ |
| 886 | .drain = http.BodyWriter.contentLengthDrain, | 886 | .drain = http.BodyWriter.contentLengthDrain, |
| 887 | .sendFile = http.BodyWriter.contentLengthSendFile, | 887 | .sendFile = http.BodyWriter.contentLengthSendFile, |
| 888 | }, | 888 | }, |
| ... | @@ -891,9 +891,9 @@ pub const Request = struct { | ... | @@ -891,9 +891,9 @@ pub const Request = struct { |
| 891 | .none => .{ | 891 | .none => .{ |
| 892 | .http_protocol_output = http_protocol_output, | 892 | .http_protocol_output = http_protocol_output, |
| 893 | .state = .none, | 893 | .state = .none, |
| 894 | .interface = .{ | 894 | .writer = .{ |
| 895 | .buffer = buffer, | 895 | .buffer = buffer, |
| 896 | .interface = &.{ | 896 | .vtable = &.{ |
| 897 | .drain = http.BodyWriter.noneDrain, | 897 | .drain = http.BodyWriter.noneDrain, |
| 898 | .sendFile = http.BodyWriter.noneSendFile, | 898 | .sendFile = http.BodyWriter.noneSendFile, |
| 899 | }, | 899 | }, |
| ... | @@ -906,7 +906,7 @@ pub const Request = struct { | ... | @@ -906,7 +906,7 @@ pub const Request = struct { |
| 906 | fn sendHead(r: *Request) Writer.Error!void { | 906 | fn sendHead(r: *Request) Writer.Error!void { |
| 907 | const uri = r.uri; | 907 | const uri = r.uri; |
| 908 | const connection = r.connection.?; | 908 | const connection = r.connection.?; |
| 909 | const w = &connection.writer; | 909 | const w = connection.writer(); |
| 910 | 910 | ||
| 911 | try r.method.write(w); | 911 | try r.method.write(w); |
| 912 | try w.writeByte(' '); | 912 | try w.writeByte(' '); |
| ... | @@ -1085,7 +1085,7 @@ pub const Request = struct { | ... | @@ -1085,7 +1085,7 @@ pub const Request = struct { |
| 1085 | if (head.status.class() == .redirect and r.redirect_behavior != .unhandled) { | 1085 | if (head.status.class() == .redirect and r.redirect_behavior != .unhandled) { |
| 1086 | if (r.redirect_behavior == .not_allowed) { | 1086 | if (r.redirect_behavior == .not_allowed) { |
| 1087 | // Connection can still be reused by skipping the body. | 1087 | // Connection can still be reused by skipping the body. |
| 1088 | var reader = r.reader.bodyReader(head.transfer_encoding, head.content_length); | 1088 | const reader = r.reader.bodyReader(&.{}, head.transfer_encoding, head.content_length); |
| 1089 | _ = reader.discardRemaining() catch |err| switch (err) { | 1089 | _ = reader.discardRemaining() catch |err| switch (err) { |
| 1090 | error.ReadFailed => connection.closing = true, | 1090 | error.ReadFailed => connection.closing = true, |
| 1091 | }; | 1091 | }; |
| ... | @@ -1117,7 +1117,7 @@ pub const Request = struct { | ... | @@ -1117,7 +1117,7 @@ pub const Request = struct { |
| 1117 | { | 1117 | { |
| 1118 | // Skip the body of the redirect response to leave the connection in | 1118 | // Skip the body of the redirect response to leave the connection in |
| 1119 | // the correct state. This causes `new_location` to be invalidated. | 1119 | // the correct state. This causes `new_location` to be invalidated. |
| 1120 | var reader = r.reader.bodyReader(head.transfer_encoding, head.content_length); | 1120 | const reader = r.reader.bodyReader(&.{}, head.transfer_encoding, head.content_length); |
| 1121 | _ = reader.discardRemaining() catch |err| switch (err) { | 1121 | _ = reader.discardRemaining() catch |err| switch (err) { |
| 1122 | error.ReadFailed => return r.reader.body_err.?, | 1122 | error.ReadFailed => return r.reader.body_err.?, |
| 1123 | }; | 1123 | }; |
| ... | @@ -1169,8 +1169,10 @@ pub const Request = struct { | ... | @@ -1169,8 +1169,10 @@ pub const Request = struct { |
| 1169 | r.uri = new_uri; | 1169 | r.uri = new_uri; |
| 1170 | r.connection = new_connection; | 1170 | r.connection = new_connection; |
| 1171 | r.reader = .{ | 1171 | r.reader = .{ |
| 1172 | .in = &new_connection.reader, | 1172 | .in = new_connection.reader(), |
| 1173 | .state = .ready, | 1173 | .state = .ready, |
| 1174 | // Populated when `http.Reader.bodyReader` is called. | ||
| 1175 | .interface = undefined, | ||
| 1174 | }; | 1176 | }; |
| 1175 | r.redirect_behavior.subtractOne(); | 1177 | r.redirect_behavior.subtractOne(); |
| 1176 | } | 1178 | } |
| ... | @@ -1292,12 +1294,12 @@ pub const basic_authorization = struct { | ... | @@ -1292,12 +1294,12 @@ pub const basic_authorization = struct { |
| 1292 | 1294 | ||
| 1293 | pub fn write(uri: Uri, out: *Writer) Writer.Error!void { | 1295 | pub fn write(uri: Uri, out: *Writer) Writer.Error!void { |
| 1294 | var buf: [max_user_len + ":".len + max_password_len]u8 = undefined; | 1296 | var buf: [max_user_len + ":".len + max_password_len]u8 = undefined; |
| 1295 | var bw: Writer = .fixed(&buf); | 1297 | var w: Writer = .fixed(&buf); |
| 1296 | bw.print("{fuser}:{fpassword}", .{ | 1298 | w.print("{fuser}:{fpassword}", .{ |
| 1297 | uri.user orelse Uri.Component.empty, | 1299 | uri.user orelse Uri.Component.empty, |
| 1298 | uri.password orelse Uri.Component.empty, | 1300 | uri.password orelse Uri.Component.empty, |
| 1299 | }) catch unreachable; | 1301 | }) catch unreachable; |
| 1300 | try out.print("Basic {b64}", .{bw.getWritten()}); | 1302 | try out.print("Basic {b64}", .{w.buffered()}); |
| 1301 | } | 1303 | } |
| 1302 | }; | 1304 | }; |
| 1303 | 1305 | ||
| ... | @@ -1622,8 +1624,10 @@ pub fn request( | ... | @@ -1622,8 +1624,10 @@ pub fn request( |
| 1622 | .client = client, | 1624 | .client = client, |
| 1623 | .connection = connection, | 1625 | .connection = connection, |
| 1624 | .reader = .{ | 1626 | .reader = .{ |
| 1625 | .in = &connection.reader, | 1627 | .in = connection.reader(), |
| 1626 | .state = .ready, | 1628 | .state = .ready, |
| 1629 | // Populated when `http.Reader.bodyReader` is called. | ||
| 1630 | .interface = undefined, | ||
| 1627 | }, | 1631 | }, |
| 1628 | .keep_alive = options.keep_alive, | 1632 | .keep_alive = options.keep_alive, |
| 1629 | .method = method, | 1633 | .method = method, |
| ... | @@ -1711,9 +1715,8 @@ pub fn fetch(client: *Client, options: FetchOptions) FetchError!FetchResult { | ... | @@ -1711,9 +1715,8 @@ pub fn fetch(client: *Client, options: FetchOptions) FetchError!FetchResult { |
| 1711 | 1715 | ||
| 1712 | if (options.payload) |payload| { | 1716 | if (options.payload) |payload| { |
| 1713 | req.transfer_encoding = .{ .content_length = payload.len }; | 1717 | req.transfer_encoding = .{ .content_length = payload.len }; |
| 1714 | var body = try req.sendBody(); | 1718 | var body = try req.sendBody(&.{}); |
| 1715 | var bw = body.writer().unbuffered(); | 1719 | try body.writer.writeAll(payload); |
| 1716 | try bw.writeAll(payload); | ||
| 1717 | try body.end(); | 1720 | try body.end(); |
| 1718 | } else { | 1721 | } else { |
| 1719 | try req.sendBodiless(); | 1722 | try req.sendBodiless(); |
| ... | @@ -1726,7 +1729,7 @@ pub fn fetch(client: *Client, options: FetchOptions) FetchError!FetchResult { | ... | @@ -1726,7 +1729,7 @@ pub fn fetch(client: *Client, options: FetchOptions) FetchError!FetchResult { |
| 1726 | var response = try req.receiveHead(redirect_buffer); | 1729 | var response = try req.receiveHead(redirect_buffer); |
| 1727 | 1730 | ||
| 1728 | const storage = options.response_storage orelse { | 1731 | const storage = options.response_storage orelse { |
| 1729 | var reader = response.reader(); | 1732 | const reader = response.reader(&.{}); |
| 1730 | _ = reader.discardRemaining() catch |err| switch (err) { | 1733 | _ = reader.discardRemaining() catch |err| switch (err) { |
| 1731 | error.ReadFailed => return response.bodyErr().?, | 1734 | error.ReadFailed => return response.bodyErr().?, |
| 1732 | }; | 1735 | }; |
| ... | @@ -1741,18 +1744,17 @@ pub fn fetch(client: *Client, options: FetchOptions) FetchError!FetchResult { | ... | @@ -1741,18 +1744,17 @@ pub fn fetch(client: *Client, options: FetchOptions) FetchError!FetchResult { |
| 1741 | defer if (options.decompress_buffer == null) client.allocator.free(decompress_buffer); | 1744 | defer if (options.decompress_buffer == null) client.allocator.free(decompress_buffer); |
| 1742 | 1745 | ||
| 1743 | var decompressor: http.Decompressor = undefined; | 1746 | var decompressor: http.Decompressor = undefined; |
| 1744 | var reader = response.readerDecompressing(&decompressor, decompress_buffer); | 1747 | const reader = response.readerDecompressing(&decompressor, decompress_buffer); |
| 1745 | const list = storage.list; | 1748 | const list = storage.list; |
| 1746 | 1749 | ||
| 1747 | if (storage.allocator) |allocator| { | 1750 | if (storage.allocator) |allocator| { |
| 1748 | reader.readRemainingArrayList(allocator, null, list, storage.append_limit, 128) catch |err| switch (err) { | 1751 | reader.appendRemaining(allocator, null, list, storage.append_limit) catch |err| switch (err) { |
| 1749 | error.ReadFailed => return response.bodyErr().?, | 1752 | error.ReadFailed => return response.bodyErr().?, |
| 1750 | else => |e| return e, | 1753 | else => |e| return e, |
| 1751 | }; | 1754 | }; |
| 1752 | } else { | 1755 | } else { |
| 1753 | var br = reader.unbuffered(); | ||
| 1754 | const buf = storage.append_limit.slice(list.unusedCapacitySlice()); | 1756 | const buf = storage.append_limit.slice(list.unusedCapacitySlice()); |
| 1755 | list.items.len += br.readSliceShort(buf) catch |err| switch (err) { | 1757 | list.items.len += reader.readSliceShort(buf) catch |err| switch (err) { |
| 1756 | error.ReadFailed => return response.bodyErr().?, | 1758 | error.ReadFailed => return response.bodyErr().?, |
| 1757 | }; | 1759 | }; |
| 1758 | } | 1760 | } |
lib/std/http/Server.zig+17-12| ... | @@ -26,6 +26,8 @@ pub fn init(in: *std.io.Reader, out: *Writer) Server { | ... | @@ -26,6 +26,8 @@ pub fn init(in: *std.io.Reader, out: *Writer) Server { |
| 26 | .reader = .{ | 26 | .reader = .{ |
| 27 | .in = in, | 27 | .in = in, |
| 28 | .state = .ready, | 28 | .state = .ready, |
| 29 | // Populated when `http.Reader.bodyReader` is called. | ||
| 30 | .interface = undefined, | ||
| 29 | }, | 31 | }, |
| 30 | .out = out, | 32 | .out = out, |
| 31 | }; | 33 | }; |
| ... | @@ -58,7 +60,7 @@ pub const Request = struct { | ... | @@ -58,7 +60,7 @@ pub const Request = struct { |
| 58 | /// Pointers in this struct are invalidated with the next call to | 60 | /// Pointers in this struct are invalidated with the next call to |
| 59 | /// `receiveHead`. | 61 | /// `receiveHead`. |
| 60 | head: Head, | 62 | head: Head, |
| 61 | respond_err: ?RespondError, | 63 | respond_err: ?RespondError = null, |
| 62 | 64 | ||
| 63 | pub const RespondError = error{ | 65 | pub const RespondError = error{ |
| 64 | /// The request contained an `expect` header with an unrecognized value. | 66 | /// The request contained an `expect` header with an unrecognized value. |
| ... | @@ -243,6 +245,7 @@ pub const Request = struct { | ... | @@ -243,6 +245,7 @@ pub const Request = struct { |
| 243 | .in = undefined, | 245 | .in = undefined, |
| 244 | .state = .received_head, | 246 | .state = .received_head, |
| 245 | .head_buffer = @constCast(request_bytes), | 247 | .head_buffer = @constCast(request_bytes), |
| 248 | .interface = undefined, | ||
| 246 | }, | 249 | }, |
| 247 | .out = undefined, | 250 | .out = undefined, |
| 248 | }; | 251 | }; |
| ... | @@ -381,8 +384,6 @@ pub const Request = struct { | ... | @@ -381,8 +384,6 @@ pub const Request = struct { |
| 381 | content_length: ?u64 = null, | 384 | content_length: ?u64 = null, |
| 382 | /// Options that are shared with the `respond` method. | 385 | /// Options that are shared with the `respond` method. |
| 383 | respond_options: RespondOptions = .{}, | 386 | respond_options: RespondOptions = .{}, |
| 384 | /// Used by `http.BodyWriter`. | ||
| 385 | buffer: []u8, | ||
| 386 | }; | 387 | }; |
| 387 | 388 | ||
| 388 | /// The header is not guaranteed to be sent until `BodyWriter.flush` or | 389 | /// The header is not guaranteed to be sent until `BodyWriter.flush` or |
| ... | @@ -400,7 +401,11 @@ pub const Request = struct { | ... | @@ -400,7 +401,11 @@ pub const Request = struct { |
| 400 | /// be done to satisfy the request. | 401 | /// be done to satisfy the request. |
| 401 | /// | 402 | /// |
| 402 | /// Asserts status is not `continue`. | 403 | /// Asserts status is not `continue`. |
| 403 | pub fn respondStreaming(request: *Request, options: RespondStreamingOptions) Writer.Error!http.BodyWriter { | 404 | pub fn respondStreaming( |
| 405 | request: *Request, | ||
| 406 | buffer: []u8, | ||
| 407 | options: RespondStreamingOptions, | ||
| 408 | ) ExpectContinueError!http.BodyWriter { | ||
| 404 | try writeExpectContinue(request); | 409 | try writeExpectContinue(request); |
| 405 | const o = options.respond_options; | 410 | const o = options.respond_options; |
| 406 | assert(o.status != .@"continue"); | 411 | assert(o.status != .@"continue"); |
| ... | @@ -448,12 +453,12 @@ pub const Request = struct { | ... | @@ -448,12 +453,12 @@ pub const Request = struct { |
| 448 | return if (elide_body) .{ | 453 | return if (elide_body) .{ |
| 449 | .http_protocol_output = request.server.out, | 454 | .http_protocol_output = request.server.out, |
| 450 | .state = state, | 455 | .state = state, |
| 451 | .interface = .discarding(options.buffer), | 456 | .writer = .discarding(buffer), |
| 452 | } else .{ | 457 | } else .{ |
| 453 | .http_protocol_output = request.server.out, | 458 | .http_protocol_output = request.server.out, |
| 454 | .state = state, | 459 | .state = state, |
| 455 | .interface = .{ | 460 | .writer = .{ |
| 456 | .buffer = options.buffer, | 461 | .buffer = buffer, |
| 457 | .vtable = switch (state) { | 462 | .vtable = switch (state) { |
| 458 | .none => &.{ | 463 | .none => &.{ |
| 459 | .drain = http.BodyWriter.noneDrain, | 464 | .drain = http.BodyWriter.noneDrain, |
| ... | @@ -559,11 +564,11 @@ pub const Request = struct { | ... | @@ -559,11 +564,11 @@ pub const Request = struct { |
| 559 | /// | 564 | /// |
| 560 | /// See `readerExpectNone` for an infallible alternative that cannot write | 565 | /// See `readerExpectNone` for an infallible alternative that cannot write |
| 561 | /// to the server output stream. | 566 | /// to the server output stream. |
| 562 | pub fn readerExpectContinue(request: *Request) ExpectContinueError!std.io.Reader { | 567 | pub fn readerExpectContinue(request: *Request, buffer: []u8) ExpectContinueError!*std.io.Reader { |
| 563 | const flush = request.head.expect != null; | 568 | const flush = request.head.expect != null; |
| 564 | try writeExpectContinue(request); | 569 | try writeExpectContinue(request); |
| 565 | if (flush) try request.server.out.flush(); | 570 | if (flush) try request.server.out.flush(); |
| 566 | return readerExpectNone(request); | 571 | return readerExpectNone(request, buffer); |
| 567 | } | 572 | } |
| 568 | 573 | ||
| 569 | /// Asserts the expect header is `null`. The caller must handle the | 574 | /// Asserts the expect header is `null`. The caller must handle the |
| ... | @@ -571,11 +576,11 @@ pub const Request = struct { | ... | @@ -571,11 +576,11 @@ pub const Request = struct { |
| 571 | /// this function. | 576 | /// this function. |
| 572 | /// | 577 | /// |
| 573 | /// Asserts that this function is only called once. | 578 | /// Asserts that this function is only called once. |
| 574 | pub fn readerExpectNone(request: *Request) std.io.Reader { | 579 | pub fn readerExpectNone(request: *Request, buffer: []u8) *std.io.Reader { |
| 575 | assert(request.server.reader.state == .received_head); | 580 | assert(request.server.reader.state == .received_head); |
| 576 | assert(request.head.expect == null); | 581 | assert(request.head.expect == null); |
| 577 | if (!request.head.method.requestHasBody()) return .ending; | 582 | if (!request.head.method.requestHasBody()) return .ending; |
| 578 | return request.server.reader.bodyReader(request.head.transfer_encoding, request.head.content_length); | 583 | return request.server.reader.bodyReader(buffer, request.head.transfer_encoding, request.head.content_length); |
| 579 | } | 584 | } |
| 580 | 585 | ||
| 581 | pub const ExpectContinueError = error{ | 586 | pub const ExpectContinueError = error{ |
| ... | @@ -611,7 +616,7 @@ pub const Request = struct { | ... | @@ -611,7 +616,7 @@ pub const Request = struct { |
| 611 | .received_head => { | 616 | .received_head => { |
| 612 | if (request.head.method.requestHasBody()) { | 617 | if (request.head.method.requestHasBody()) { |
| 613 | assert(request.head.transfer_encoding != .none or request.head.content_length != null); | 618 | assert(request.head.transfer_encoding != .none or request.head.content_length != null); |
| 614 | const reader_interface = request.reader() catch return false; | 619 | const reader_interface = request.readerExpectContinue(&.{}) catch return false; |
| 615 | _ = reader_interface.discardRemaining() catch return false; | 620 | _ = reader_interface.discardRemaining() catch return false; |
| 616 | assert(r.state == .ready); | 621 | assert(r.state == .ready); |
| 617 | } else { | 622 | } else { |
lib/std/http/test.zig+97-117| ... | @@ -21,7 +21,7 @@ test "trailers" { | ... | @@ -21,7 +21,7 @@ test "trailers" { |
| 21 | 21 | ||
| 22 | var connection_br = connection.stream.reader(&recv_buffer); | 22 | var connection_br = connection.stream.reader(&recv_buffer); |
| 23 | var connection_bw = connection.stream.writer(&send_buffer); | 23 | var connection_bw = connection.stream.writer(&send_buffer); |
| 24 | var server = http.Server.init(&connection_br, &connection_bw); | 24 | var server = http.Server.init(connection_br.interface(), &connection_bw.interface); |
| 25 | 25 | ||
| 26 | try expectEqual(.ready, server.reader.state); | 26 | try expectEqual(.ready, server.reader.state); |
| 27 | var request = try server.receiveHead(); | 27 | var request = try server.receiveHead(); |
| ... | @@ -33,11 +33,10 @@ test "trailers" { | ... | @@ -33,11 +33,10 @@ test "trailers" { |
| 33 | fn serve(request: *http.Server.Request) !void { | 33 | fn serve(request: *http.Server.Request) !void { |
| 34 | try expectEqualStrings(request.head.target, "/trailer"); | 34 | try expectEqualStrings(request.head.target, "/trailer"); |
| 35 | 35 | ||
| 36 | var response = try request.respondStreaming(.{}); | 36 | var response = try request.respondStreaming(&.{}, .{}); |
| 37 | var bw = response.writer().unbuffered(); | 37 | try response.writer.writeAll("Hello, "); |
| 38 | try bw.writeAll("Hello, "); | ||
| 39 | try response.flush(); | 38 | try response.flush(); |
| 40 | try bw.writeAll("World!\n"); | 39 | try response.writer.writeAll("World!\n"); |
| 41 | try response.flush(); | 40 | try response.flush(); |
| 42 | try response.endChunked(.{ | 41 | try response.endChunked(.{ |
| 43 | .trailers = &.{ | 42 | .trailers = &.{ |
| ... | @@ -66,7 +65,7 @@ test "trailers" { | ... | @@ -66,7 +65,7 @@ test "trailers" { |
| 66 | try req.sendBodiless(); | 65 | try req.sendBodiless(); |
| 67 | var response = try req.receiveHead(&.{}); | 66 | var response = try req.receiveHead(&.{}); |
| 68 | 67 | ||
| 69 | const body = try response.reader().readRemainingAlloc(gpa, .limited(8192)); | 68 | const body = try response.reader(&.{}).allocRemaining(gpa, .limited(8192)); |
| 70 | defer gpa.free(body); | 69 | defer gpa.free(body); |
| 71 | 70 | ||
| 72 | try expectEqualStrings("Hello, World!\n", body); | 71 | try expectEqualStrings("Hello, World!\n", body); |
| ... | @@ -104,13 +103,13 @@ test "HTTP server handles a chunked transfer coding request" { | ... | @@ -104,13 +103,13 @@ test "HTTP server handles a chunked transfer coding request" { |
| 104 | 103 | ||
| 105 | var connection_br = connection.stream.reader(&recv_buffer); | 104 | var connection_br = connection.stream.reader(&recv_buffer); |
| 106 | var connection_bw = connection.stream.writer(&send_buffer); | 105 | var connection_bw = connection.stream.writer(&send_buffer); |
| 107 | var server = http.Server.init(&connection_br, &connection_bw); | 106 | var server = http.Server.init(connection_br.interface(), &connection_bw.interface); |
| 108 | var request = try server.receiveHead(); | 107 | var request = try server.receiveHead(); |
| 109 | 108 | ||
| 110 | try expect(request.head.transfer_encoding == .chunked); | 109 | try expect(request.head.transfer_encoding == .chunked); |
| 111 | 110 | ||
| 112 | var buf: [128]u8 = undefined; | 111 | var buf: [128]u8 = undefined; |
| 113 | var br = (try request.reader()).unbuffered(); | 112 | var br = try request.readerExpectContinue(&.{}); |
| 114 | const n = try br.readSliceShort(&buf); | 113 | const n = try br.readSliceShort(&buf); |
| 115 | try expectEqualStrings("ABCD", buf[0..n]); | 114 | try expectEqualStrings("ABCD", buf[0..n]); |
| 116 | 115 | ||
| ... | @@ -141,9 +140,8 @@ test "HTTP server handles a chunked transfer coding request" { | ... | @@ -141,9 +140,8 @@ test "HTTP server handles a chunked transfer coding request" { |
| 141 | const gpa = std.testing.allocator; | 140 | const gpa = std.testing.allocator; |
| 142 | const stream = try std.net.tcpConnectToHost(gpa, "127.0.0.1", test_server.port()); | 141 | const stream = try std.net.tcpConnectToHost(gpa, "127.0.0.1", test_server.port()); |
| 143 | defer stream.close(); | 142 | defer stream.close(); |
| 144 | var stream_writer = stream.writer(); | 143 | var stream_writer = stream.writer(&.{}); |
| 145 | var writer = stream_writer.interface().unbuffered(); | 144 | try stream_writer.interface.writeAll(request_bytes); |
| 146 | try writer.writeAll(request_bytes); | ||
| 147 | 145 | ||
| 148 | const expected_response = | 146 | const expected_response = |
| 149 | "HTTP/1.1 200 OK\r\n" ++ | 147 | "HTTP/1.1 200 OK\r\n" ++ |
| ... | @@ -152,8 +150,8 @@ test "HTTP server handles a chunked transfer coding request" { | ... | @@ -152,8 +150,8 @@ test "HTTP server handles a chunked transfer coding request" { |
| 152 | "content-type: text/plain\r\n" ++ | 150 | "content-type: text/plain\r\n" ++ |
| 153 | "\r\n" ++ | 151 | "\r\n" ++ |
| 154 | "message from server!\n"; | 152 | "message from server!\n"; |
| 155 | var stream_reader = stream.reader(); | 153 | var stream_reader = stream.reader(&.{}); |
| 156 | const response = try stream_reader.interface().readRemainingAlloc(gpa, .limited(expected_response.len)); | 154 | const response = try stream_reader.interface().allocRemaining(gpa, .limited(expected_response.len)); |
| 157 | defer gpa.free(response); | 155 | defer gpa.free(response); |
| 158 | try expectEqualStrings(expected_response, response); | 156 | try expectEqualStrings(expected_response, response); |
| 159 | } | 157 | } |
| ... | @@ -169,11 +167,9 @@ test "echo content server" { | ... | @@ -169,11 +167,9 @@ test "echo content server" { |
| 169 | const connection = try net_server.accept(); | 167 | const connection = try net_server.accept(); |
| 170 | defer connection.stream.close(); | 168 | defer connection.stream.close(); |
| 171 | 169 | ||
| 172 | var stream_reader = connection.stream.reader(); | 170 | var connection_br = connection.stream.reader(&recv_buffer); |
| 173 | var stream_writer = connection.stream.writer(); | 171 | var connection_bw = connection.stream.writer(&send_buffer); |
| 174 | var connection_br = stream_reader.interface().buffered(&recv_buffer); | 172 | var http_server = http.Server.init(connection_br.interface(), &connection_bw.interface); |
| 175 | var connection_bw = stream_writer.interface().buffered(&send_buffer); | ||
| 176 | var http_server = http.Server.init(&connection_br, &connection_bw); | ||
| 177 | 173 | ||
| 178 | while (http_server.reader.state == .ready) { | 174 | while (http_server.reader.state == .ready) { |
| 179 | var request = http_server.receiveHead() catch |err| switch (err) { | 175 | var request = http_server.receiveHead() catch |err| switch (err) { |
| ... | @@ -185,7 +181,7 @@ test "echo content server" { | ... | @@ -185,7 +181,7 @@ test "echo content server" { |
| 185 | } | 181 | } |
| 186 | if (request.head.expect) |expect_header_value| { | 182 | if (request.head.expect) |expect_header_value| { |
| 187 | if (mem.eql(u8, expect_header_value, "garbage")) { | 183 | if (mem.eql(u8, expect_header_value, "garbage")) { |
| 188 | try expectError(error.HttpExpectationFailed, request.reader()); | 184 | try expectError(error.HttpExpectationFailed, request.readerExpectContinue(&.{})); |
| 189 | try request.respond("", .{ .keep_alive = false }); | 185 | try request.respond("", .{ .keep_alive = false }); |
| 190 | continue; | 186 | continue; |
| 191 | } | 187 | } |
| ... | @@ -207,14 +203,14 @@ test "echo content server" { | ... | @@ -207,14 +203,14 @@ test "echo content server" { |
| 207 | // request.head.target, | 203 | // request.head.target, |
| 208 | //}); | 204 | //}); |
| 209 | 205 | ||
| 210 | const body = try (try request.reader()).readRemainingAlloc(std.testing.allocator, .limited(8192)); | 206 | const body = try (try request.readerExpectContinue(&.{})).allocRemaining(std.testing.allocator, .limited(8192)); |
| 211 | defer std.testing.allocator.free(body); | 207 | defer std.testing.allocator.free(body); |
| 212 | 208 | ||
| 213 | try expect(mem.startsWith(u8, request.head.target, "/echo-content")); | 209 | try expect(mem.startsWith(u8, request.head.target, "/echo-content")); |
| 214 | try expectEqualStrings("Hello, World!\n", body); | 210 | try expectEqualStrings("Hello, World!\n", body); |
| 215 | try expectEqualStrings("text/plain", request.head.content_type.?); | 211 | try expectEqualStrings("text/plain", request.head.content_type.?); |
| 216 | 212 | ||
| 217 | var response = try request.respondStreaming(.{ | 213 | var response = try request.respondStreaming(&.{}, .{ |
| 218 | .content_length = switch (request.head.transfer_encoding) { | 214 | .content_length = switch (request.head.transfer_encoding) { |
| 219 | .chunked => null, | 215 | .chunked => null, |
| 220 | .none => len: { | 216 | .none => len: { |
| ... | @@ -224,9 +220,9 @@ test "echo content server" { | ... | @@ -224,9 +220,9 @@ test "echo content server" { |
| 224 | }, | 220 | }, |
| 225 | }); | 221 | }); |
| 226 | try response.flush(); // Test an early flush to send the HTTP headers before the body. | 222 | try response.flush(); // Test an early flush to send the HTTP headers before the body. |
| 227 | var bw = response.writer().unbuffered(); | 223 | const w = &response.writer; |
| 228 | try bw.writeAll("Hello, "); | 224 | try w.writeAll("Hello, "); |
| 229 | try bw.writeAll("World!\n"); | 225 | try w.writeAll("World!\n"); |
| 230 | try response.end(); | 226 | try response.end(); |
| 231 | //std.debug.print(" server finished responding\n", .{}); | 227 | //std.debug.print(" server finished responding\n", .{}); |
| 232 | } | 228 | } |
| ... | @@ -259,27 +255,25 @@ test "Server.Request.respondStreaming non-chunked, unknown content-length" { | ... | @@ -259,27 +255,25 @@ test "Server.Request.respondStreaming non-chunked, unknown content-length" { |
| 259 | const connection = try net_server.accept(); | 255 | const connection = try net_server.accept(); |
| 260 | defer connection.stream.close(); | 256 | defer connection.stream.close(); |
| 261 | 257 | ||
| 262 | var stream_reader = connection.stream.reader(); | 258 | var connection_br = connection.stream.reader(&recv_buffer); |
| 263 | var stream_writer = connection.stream.writer(); | 259 | var connection_bw = connection.stream.writer(&send_buffer); |
| 264 | var connection_br = stream_reader.interface().buffered(&recv_buffer); | 260 | var server = http.Server.init(connection_br.interface(), &connection_bw.interface); |
| 265 | var connection_bw = stream_writer.interface().buffered(&send_buffer); | ||
| 266 | var server = http.Server.init(&connection_br, &connection_bw); | ||
| 267 | 261 | ||
| 268 | try expectEqual(.ready, server.reader.state); | 262 | try expectEqual(.ready, server.reader.state); |
| 269 | var request = try server.receiveHead(); | 263 | var request = try server.receiveHead(); |
| 270 | try expectEqualStrings(request.head.target, "/foo"); | 264 | try expectEqualStrings(request.head.target, "/foo"); |
| 271 | var response = try request.respondStreaming(.{ | 265 | var buf: [30]u8 = undefined; |
| 266 | var response = try request.respondStreaming(&buf, .{ | ||
| 272 | .respond_options = .{ | 267 | .respond_options = .{ |
| 273 | .transfer_encoding = .none, | 268 | .transfer_encoding = .none, |
| 274 | }, | 269 | }, |
| 275 | }); | 270 | }); |
| 276 | var buf: [30]u8 = undefined; | 271 | const w = &response.writer; |
| 277 | var bw = response.writer().buffered(&buf); | ||
| 278 | for (0..500) |i| { | 272 | for (0..500) |i| { |
| 279 | try bw.print("{d}, ah ha ha!\n", .{i}); | 273 | try w.print("{d}, ah ha ha!\n", .{i}); |
| 280 | } | 274 | } |
| 281 | try expectEqual(7390, bw.count); | 275 | try expectEqual(7390, w.count); |
| 282 | try bw.flush(); | 276 | try w.flush(); |
| 283 | try response.end(); | 277 | try response.end(); |
| 284 | try expectEqual(.closing, server.reader.state); | 278 | try expectEqual(.closing, server.reader.state); |
| 285 | } | 279 | } |
| ... | @@ -291,12 +285,11 @@ test "Server.Request.respondStreaming non-chunked, unknown content-length" { | ... | @@ -291,12 +285,11 @@ test "Server.Request.respondStreaming non-chunked, unknown content-length" { |
| 291 | const gpa = std.testing.allocator; | 285 | const gpa = std.testing.allocator; |
| 292 | const stream = try std.net.tcpConnectToHost(gpa, "127.0.0.1", test_server.port()); | 286 | const stream = try std.net.tcpConnectToHost(gpa, "127.0.0.1", test_server.port()); |
| 293 | defer stream.close(); | 287 | defer stream.close(); |
| 294 | var stream_writer = stream.writer(); | 288 | var stream_writer = stream.writer(&.{}); |
| 295 | var writer = stream_writer.interface().unbuffered(); | 289 | try stream_writer.interface.writeAll(request_bytes); |
| 296 | try writer.writeAll(request_bytes); | ||
| 297 | 290 | ||
| 298 | var stream_reader = stream.reader(); | 291 | var stream_reader = stream.reader(&.{}); |
| 299 | const response = try stream_reader.interface().readRemainingAlloc(gpa, .limited(8192)); | 292 | const response = try stream_reader.interface().allocRemaining(gpa, .limited(8192)); |
| 300 | defer gpa.free(response); | 293 | defer gpa.free(response); |
| 301 | 294 | ||
| 302 | var expected_response = std.ArrayList(u8).init(gpa); | 295 | var expected_response = std.ArrayList(u8).init(gpa); |
| ... | @@ -329,11 +322,9 @@ test "receiving arbitrary http headers from the client" { | ... | @@ -329,11 +322,9 @@ test "receiving arbitrary http headers from the client" { |
| 329 | const connection = try net_server.accept(); | 322 | const connection = try net_server.accept(); |
| 330 | defer connection.stream.close(); | 323 | defer connection.stream.close(); |
| 331 | 324 | ||
| 332 | var stream_reader = connection.stream.reader(); | 325 | var connection_br = connection.stream.reader(&recv_buffer); |
| 333 | var stream_writer = connection.stream.writer(); | 326 | var connection_bw = connection.stream.writer(&send_buffer); |
| 334 | var connection_br = stream_reader.interface().buffered(&recv_buffer); | 327 | var server = http.Server.init(connection_br.interface(), &connection_bw.interface); |
| 335 | var connection_bw = stream_writer.interface().buffered(&send_buffer); | ||
| 336 | var server = http.Server.init(&connection_br, &connection_bw); | ||
| 337 | 328 | ||
| 338 | try expectEqual(.ready, server.reader.state); | 329 | try expectEqual(.ready, server.reader.state); |
| 339 | var request = try server.receiveHead(); | 330 | var request = try server.receiveHead(); |
| ... | @@ -364,12 +355,11 @@ test "receiving arbitrary http headers from the client" { | ... | @@ -364,12 +355,11 @@ test "receiving arbitrary http headers from the client" { |
| 364 | const gpa = std.testing.allocator; | 355 | const gpa = std.testing.allocator; |
| 365 | const stream = try std.net.tcpConnectToHost(gpa, "127.0.0.1", test_server.port()); | 356 | const stream = try std.net.tcpConnectToHost(gpa, "127.0.0.1", test_server.port()); |
| 366 | defer stream.close(); | 357 | defer stream.close(); |
| 367 | var stream_writer = stream.writer(); | 358 | var stream_writer = stream.writer(&.{}); |
| 368 | var writer = stream_writer.interface().unbuffered(); | 359 | try stream_writer.interface.writeAll(request_bytes); |
| 369 | try writer.writeAll(request_bytes); | ||
| 370 | 360 | ||
| 371 | var stream_reader = stream.reader(); | 361 | var stream_reader = stream.reader(&.{}); |
| 372 | const response = try stream_reader.interface().readRemainingAlloc(gpa, .limited(8192)); | 362 | const response = try stream_reader.interface().allocRemaining(gpa, .limited(8192)); |
| 373 | defer gpa.free(response); | 363 | defer gpa.free(response); |
| 374 | 364 | ||
| 375 | var expected_response = std.ArrayList(u8).init(gpa); | 365 | var expected_response = std.ArrayList(u8).init(gpa); |
| ... | @@ -397,11 +387,9 @@ test "general client/server API coverage" { | ... | @@ -397,11 +387,9 @@ test "general client/server API coverage" { |
| 397 | var connection = try net_server.accept(); | 387 | var connection = try net_server.accept(); |
| 398 | defer connection.stream.close(); | 388 | defer connection.stream.close(); |
| 399 | 389 | ||
| 400 | var stream_reader = connection.stream.reader(); | 390 | var connection_br = connection.stream.reader(&recv_buffer); |
| 401 | var stream_writer = connection.stream.writer(); | 391 | var connection_bw = connection.stream.writer(&send_buffer); |
| 402 | var connection_br = stream_reader.interface().buffered(&recv_buffer); | 392 | var http_server = http.Server.init(connection_br.interface(), &connection_bw.interface); |
| 403 | var connection_bw = stream_writer.interface().buffered(&send_buffer); | ||
| 404 | var http_server = http.Server.init(&connection_br, &connection_bw); | ||
| 405 | 393 | ||
| 406 | while (http_server.reader.state == .ready) { | 394 | while (http_server.reader.state == .ready) { |
| 407 | var request = http_server.receiveHead() catch |err| switch (err) { | 395 | var request = http_server.receiveHead() catch |err| switch (err) { |
| ... | @@ -424,11 +412,11 @@ test "general client/server API coverage" { | ... | @@ -424,11 +412,11 @@ test "general client/server API coverage" { |
| 424 | }); | 412 | }); |
| 425 | 413 | ||
| 426 | const gpa = std.testing.allocator; | 414 | const gpa = std.testing.allocator; |
| 427 | const body = try (try request.reader()).readRemainingAlloc(gpa, .limited(8192)); | 415 | const body = try (try request.readerExpectContinue(&.{})).allocRemaining(gpa, .limited(8192)); |
| 428 | defer gpa.free(body); | 416 | defer gpa.free(body); |
| 429 | 417 | ||
| 430 | if (mem.startsWith(u8, request.head.target, "/get")) { | 418 | if (mem.startsWith(u8, request.head.target, "/get")) { |
| 431 | var response = try request.respondStreaming(.{ | 419 | var response = try request.respondStreaming(&.{}, .{ |
| 432 | .content_length = if (mem.indexOf(u8, request.head.target, "?chunked") == null) | 420 | .content_length = if (mem.indexOf(u8, request.head.target, "?chunked") == null) |
| 433 | 14 | 421 | 14 |
| 434 | else | 422 | else |
| ... | @@ -439,35 +427,35 @@ test "general client/server API coverage" { | ... | @@ -439,35 +427,35 @@ test "general client/server API coverage" { |
| 439 | }, | 427 | }, |
| 440 | }, | 428 | }, |
| 441 | }); | 429 | }); |
| 442 | var bw = response.writer().unbuffered(); | 430 | const w = &response.writer; |
| 443 | try bw.writeAll("Hello, "); | 431 | try w.writeAll("Hello, "); |
| 444 | try bw.writeAll("World!\n"); | 432 | try w.writeAll("World!\n"); |
| 445 | try response.end(); | 433 | try response.end(); |
| 446 | // Writing again would cause an assertion failure. | 434 | // Writing again would cause an assertion failure. |
| 447 | } else if (mem.startsWith(u8, request.head.target, "/large")) { | 435 | } else if (mem.startsWith(u8, request.head.target, "/large")) { |
| 448 | var response = try request.respondStreaming(.{ | 436 | var response = try request.respondStreaming(&.{}, .{ |
| 449 | .content_length = 14 * 1024 + 14 * 10, | 437 | .content_length = 14 * 1024 + 14 * 10, |
| 450 | }); | 438 | }); |
| 451 | 439 | ||
| 452 | try response.flush(); // Test an early flush to send the HTTP headers before the body. | 440 | try response.flush(); // Test an early flush to send the HTTP headers before the body. |
| 453 | 441 | ||
| 454 | var bw = response.writer().unbuffered(); | 442 | const w = &response.writer; |
| 455 | 443 | ||
| 456 | var i: u32 = 0; | 444 | var i: u32 = 0; |
| 457 | while (i < 5) : (i += 1) { | 445 | while (i < 5) : (i += 1) { |
| 458 | try bw.writeAll("Hello, World!\n"); | 446 | try w.writeAll("Hello, World!\n"); |
| 459 | } | 447 | } |
| 460 | 448 | ||
| 461 | try bw.writeAll("Hello, World!\n" ** 1024); | 449 | try w.writeAll("Hello, World!\n" ** 1024); |
| 462 | 450 | ||
| 463 | i = 0; | 451 | i = 0; |
| 464 | while (i < 5) : (i += 1) { | 452 | while (i < 5) : (i += 1) { |
| 465 | try bw.writeAll("Hello, World!\n"); | 453 | try w.writeAll("Hello, World!\n"); |
| 466 | } | 454 | } |
| 467 | 455 | ||
| 468 | try response.end(); | 456 | try response.end(); |
| 469 | } else if (mem.eql(u8, request.head.target, "/redirect/1")) { | 457 | } else if (mem.eql(u8, request.head.target, "/redirect/1")) { |
| 470 | var response = try request.respondStreaming(.{ | 458 | var response = try request.respondStreaming(&.{}, .{ |
| 471 | .respond_options = .{ | 459 | .respond_options = .{ |
| 472 | .status = .found, | 460 | .status = .found, |
| 473 | .extra_headers = &.{ | 461 | .extra_headers = &.{ |
| ... | @@ -476,9 +464,9 @@ test "general client/server API coverage" { | ... | @@ -476,9 +464,9 @@ test "general client/server API coverage" { |
| 476 | }, | 464 | }, |
| 477 | }); | 465 | }); |
| 478 | 466 | ||
| 479 | var bw = response.writer().unbuffered(); | 467 | const w = &response.writer; |
| 480 | try bw.writeAll("Hello, "); | 468 | try w.writeAll("Hello, "); |
| 481 | try bw.writeAll("Redirected!\n"); | 469 | try w.writeAll("Redirected!\n"); |
| 482 | try response.end(); | 470 | try response.end(); |
| 483 | } else if (mem.eql(u8, request.head.target, "/redirect/2")) { | 471 | } else if (mem.eql(u8, request.head.target, "/redirect/2")) { |
| 484 | try request.respond("Hello, Redirected!\n", .{ | 472 | try request.respond("Hello, Redirected!\n", .{ |
| ... | @@ -567,7 +555,7 @@ test "general client/server API coverage" { | ... | @@ -567,7 +555,7 @@ test "general client/server API coverage" { |
| 567 | try req.sendBodiless(); | 555 | try req.sendBodiless(); |
| 568 | var response = try req.receiveHead(&redirect_buffer); | 556 | var response = try req.receiveHead(&redirect_buffer); |
| 569 | 557 | ||
| 570 | const body = try response.reader().readRemainingAlloc(gpa, .limited(8192)); | 558 | const body = try response.reader(&.{}).allocRemaining(gpa, .limited(8192)); |
| 571 | defer gpa.free(body); | 559 | defer gpa.free(body); |
| 572 | 560 | ||
| 573 | try expectEqualStrings("Hello, World!\n", body); | 561 | try expectEqualStrings("Hello, World!\n", body); |
| ... | @@ -590,7 +578,7 @@ test "general client/server API coverage" { | ... | @@ -590,7 +578,7 @@ test "general client/server API coverage" { |
| 590 | try req.sendBodiless(); | 578 | try req.sendBodiless(); |
| 591 | var response = try req.receiveHead(&redirect_buffer); | 579 | var response = try req.receiveHead(&redirect_buffer); |
| 592 | 580 | ||
| 593 | const body = try response.reader().readRemainingAlloc(gpa, .limited(8192 * 1024)); | 581 | const body = try response.reader(&.{}).allocRemaining(gpa, .limited(8192 * 1024)); |
| 594 | defer gpa.free(body); | 582 | defer gpa.free(body); |
| 595 | 583 | ||
| 596 | try expectEqual(@as(usize, 14 * 1024 + 14 * 10), body.len); | 584 | try expectEqual(@as(usize, 14 * 1024 + 14 * 10), body.len); |
| ... | @@ -612,7 +600,7 @@ test "general client/server API coverage" { | ... | @@ -612,7 +600,7 @@ test "general client/server API coverage" { |
| 612 | try req.sendBodiless(); | 600 | try req.sendBodiless(); |
| 613 | var response = try req.receiveHead(&redirect_buffer); | 601 | var response = try req.receiveHead(&redirect_buffer); |
| 614 | 602 | ||
| 615 | const body = try response.reader().readRemainingAlloc(gpa, .limited(8192)); | 603 | const body = try response.reader(&.{}).allocRemaining(gpa, .limited(8192)); |
| 616 | defer gpa.free(body); | 604 | defer gpa.free(body); |
| 617 | 605 | ||
| 618 | try expectEqualStrings("", body); | 606 | try expectEqualStrings("", body); |
| ... | @@ -636,7 +624,7 @@ test "general client/server API coverage" { | ... | @@ -636,7 +624,7 @@ test "general client/server API coverage" { |
| 636 | try req.sendBodiless(); | 624 | try req.sendBodiless(); |
| 637 | var response = try req.receiveHead(&redirect_buffer); | 625 | var response = try req.receiveHead(&redirect_buffer); |
| 638 | 626 | ||
| 639 | const body = try response.reader().readRemainingAlloc(gpa, .limited(8192)); | 627 | const body = try response.reader(&.{}).allocRemaining(gpa, .limited(8192)); |
| 640 | defer gpa.free(body); | 628 | defer gpa.free(body); |
| 641 | 629 | ||
| 642 | try expectEqualStrings("Hello, World!\n", body); | 630 | try expectEqualStrings("Hello, World!\n", body); |
| ... | @@ -659,7 +647,7 @@ test "general client/server API coverage" { | ... | @@ -659,7 +647,7 @@ test "general client/server API coverage" { |
| 659 | try req.sendBodiless(); | 647 | try req.sendBodiless(); |
| 660 | var response = try req.receiveHead(&redirect_buffer); | 648 | var response = try req.receiveHead(&redirect_buffer); |
| 661 | 649 | ||
| 662 | const body = try response.reader().readRemainingAlloc(gpa, .limited(8192)); | 650 | const body = try response.reader(&.{}).allocRemaining(gpa, .limited(8192)); |
| 663 | defer gpa.free(body); | 651 | defer gpa.free(body); |
| 664 | 652 | ||
| 665 | try expectEqualStrings("", body); | 653 | try expectEqualStrings("", body); |
| ... | @@ -685,7 +673,7 @@ test "general client/server API coverage" { | ... | @@ -685,7 +673,7 @@ test "general client/server API coverage" { |
| 685 | try req.sendBodiless(); | 673 | try req.sendBodiless(); |
| 686 | var response = try req.receiveHead(&redirect_buffer); | 674 | var response = try req.receiveHead(&redirect_buffer); |
| 687 | 675 | ||
| 688 | const body = try response.reader().readRemainingAlloc(gpa, .limited(8192)); | 676 | const body = try response.reader(&.{}).allocRemaining(gpa, .limited(8192)); |
| 689 | defer gpa.free(body); | 677 | defer gpa.free(body); |
| 690 | 678 | ||
| 691 | try expectEqualStrings("Hello, World!\n", body); | 679 | try expectEqualStrings("Hello, World!\n", body); |
| ... | @@ -714,7 +702,7 @@ test "general client/server API coverage" { | ... | @@ -714,7 +702,7 @@ test "general client/server API coverage" { |
| 714 | 702 | ||
| 715 | try std.testing.expectEqual(.ok, response.head.status); | 703 | try std.testing.expectEqual(.ok, response.head.status); |
| 716 | 704 | ||
| 717 | const body = try response.reader().readRemainingAlloc(gpa, .limited(8192)); | 705 | const body = try response.reader(&.{}).allocRemaining(gpa, .limited(8192)); |
| 718 | defer gpa.free(body); | 706 | defer gpa.free(body); |
| 719 | 707 | ||
| 720 | try expectEqualStrings("", body); | 708 | try expectEqualStrings("", body); |
| ... | @@ -751,7 +739,7 @@ test "general client/server API coverage" { | ... | @@ -751,7 +739,7 @@ test "general client/server API coverage" { |
| 751 | try req.sendBodiless(); | 739 | try req.sendBodiless(); |
| 752 | var response = try req.receiveHead(&redirect_buffer); | 740 | var response = try req.receiveHead(&redirect_buffer); |
| 753 | 741 | ||
| 754 | const body = try response.reader().readRemainingAlloc(gpa, .limited(8192)); | 742 | const body = try response.reader(&.{}).allocRemaining(gpa, .limited(8192)); |
| 755 | defer gpa.free(body); | 743 | defer gpa.free(body); |
| 756 | 744 | ||
| 757 | try expectEqualStrings("Hello, World!\n", body); | 745 | try expectEqualStrings("Hello, World!\n", body); |
| ... | @@ -773,7 +761,7 @@ test "general client/server API coverage" { | ... | @@ -773,7 +761,7 @@ test "general client/server API coverage" { |
| 773 | try req.sendBodiless(); | 761 | try req.sendBodiless(); |
| 774 | var response = try req.receiveHead(&redirect_buffer); | 762 | var response = try req.receiveHead(&redirect_buffer); |
| 775 | 763 | ||
| 776 | const body = try response.reader().readRemainingAlloc(gpa, .limited(8192)); | 764 | const body = try response.reader(&.{}).allocRemaining(gpa, .limited(8192)); |
| 777 | defer gpa.free(body); | 765 | defer gpa.free(body); |
| 778 | 766 | ||
| 779 | try expectEqualStrings("Hello, World!\n", body); | 767 | try expectEqualStrings("Hello, World!\n", body); |
| ... | @@ -795,7 +783,7 @@ test "general client/server API coverage" { | ... | @@ -795,7 +783,7 @@ test "general client/server API coverage" { |
| 795 | try req.sendBodiless(); | 783 | try req.sendBodiless(); |
| 796 | var response = try req.receiveHead(&redirect_buffer); | 784 | var response = try req.receiveHead(&redirect_buffer); |
| 797 | 785 | ||
| 798 | const body = try response.reader().readRemainingAlloc(gpa, .limited(8192)); | 786 | const body = try response.reader(&.{}).allocRemaining(gpa, .limited(8192)); |
| 799 | defer gpa.free(body); | 787 | defer gpa.free(body); |
| 800 | 788 | ||
| 801 | try expectEqualStrings("Hello, World!\n", body); | 789 | try expectEqualStrings("Hello, World!\n", body); |
| ... | @@ -836,7 +824,7 @@ test "general client/server API coverage" { | ... | @@ -836,7 +824,7 @@ test "general client/server API coverage" { |
| 836 | try req.sendBodiless(); | 824 | try req.sendBodiless(); |
| 837 | var response = try req.receiveHead(&redirect_buffer); | 825 | var response = try req.receiveHead(&redirect_buffer); |
| 838 | 826 | ||
| 839 | const body = try response.reader().readRemainingAlloc(gpa, .limited(8192)); | 827 | const body = try response.reader(&.{}).allocRemaining(gpa, .limited(8192)); |
| 840 | defer gpa.free(body); | 828 | defer gpa.free(body); |
| 841 | 829 | ||
| 842 | try expectEqualStrings("Encoded redirect successful!\n", body); | 830 | try expectEqualStrings("Encoded redirect successful!\n", body); |
| ... | @@ -878,20 +866,18 @@ test "Server streams both reading and writing" { | ... | @@ -878,20 +866,18 @@ test "Server streams both reading and writing" { |
| 878 | const connection = try net_server.accept(); | 866 | const connection = try net_server.accept(); |
| 879 | defer connection.stream.close(); | 867 | defer connection.stream.close(); |
| 880 | 868 | ||
| 881 | var stream_reader = connection.stream.reader(); | 869 | var connection_br = connection.stream.reader(&recv_buffer); |
| 882 | var stream_writer = connection.stream.writer(); | 870 | var connection_bw = connection.stream.writer(&send_buffer); |
| 883 | var connection_br = stream_reader.interface().buffered(&recv_buffer); | 871 | var server = http.Server.init(connection_br.interface(), &connection_bw.interface); |
| 884 | var connection_bw = stream_writer.interface().buffered(&send_buffer); | ||
| 885 | var server = http.Server.init(&connection_br, &connection_bw); | ||
| 886 | var request = try server.receiveHead(); | 872 | var request = try server.receiveHead(); |
| 887 | var read_buffer: [100]u8 = undefined; | 873 | var read_buffer: [100]u8 = undefined; |
| 888 | var br = (try request.reader()).buffered(&read_buffer); | 874 | var br = try request.readerExpectContinue(&read_buffer); |
| 889 | var response = try request.respondStreaming(.{ | 875 | var response = try request.respondStreaming(&.{}, .{ |
| 890 | .respond_options = .{ | 876 | .respond_options = .{ |
| 891 | .transfer_encoding = .none, // Causes keep_alive=false | 877 | .transfer_encoding = .none, // Causes keep_alive=false |
| 892 | }, | 878 | }, |
| 893 | }); | 879 | }); |
| 894 | var bw = response.writer().unbuffered(); | 880 | const w = &response.writer; |
| 895 | 881 | ||
| 896 | while (true) { | 882 | while (true) { |
| 897 | try response.flush(); | 883 | try response.flush(); |
| ... | @@ -901,7 +887,7 @@ test "Server streams both reading and writing" { | ... | @@ -901,7 +887,7 @@ test "Server streams both reading and writing" { |
| 901 | }; | 887 | }; |
| 902 | br.toss(buf.len); | 888 | br.toss(buf.len); |
| 903 | for (buf) |*b| b.* = std.ascii.toUpper(b.*); | 889 | for (buf) |*b| b.* = std.ascii.toUpper(b.*); |
| 904 | try bw.writeAll(buf); | 890 | try w.writeAll(buf); |
| 905 | } | 891 | } |
| 906 | try response.end(); | 892 | try response.end(); |
| 907 | } | 893 | } |
| ... | @@ -921,15 +907,14 @@ test "Server streams both reading and writing" { | ... | @@ -921,15 +907,14 @@ test "Server streams both reading and writing" { |
| 921 | defer req.deinit(); | 907 | defer req.deinit(); |
| 922 | 908 | ||
| 923 | req.transfer_encoding = .chunked; | 909 | req.transfer_encoding = .chunked; |
| 924 | var body_writer = try req.sendBody(); | 910 | var body_writer = try req.sendBody(&.{}); |
| 925 | var response = try req.receiveHead(&redirect_buffer); | 911 | var response = try req.receiveHead(&redirect_buffer); |
| 926 | 912 | ||
| 927 | var w = body_writer.writer().unbuffered(); | 913 | try body_writer.writer.writeAll("one "); |
| 928 | try w.writeAll("one "); | 914 | try body_writer.writer.writeAll("fish"); |
| 929 | try w.writeAll("fish"); | ||
| 930 | try body_writer.end(); | 915 | try body_writer.end(); |
| 931 | 916 | ||
| 932 | const body = try response.reader().readRemainingAlloc(std.testing.allocator, .limited(8192)); | 917 | const body = try response.reader(&.{}).allocRemaining(std.testing.allocator, .limited(8192)); |
| 933 | defer std.testing.allocator.free(body); | 918 | defer std.testing.allocator.free(body); |
| 934 | 919 | ||
| 935 | try expectEqualStrings("ONE FISH", body); | 920 | try expectEqualStrings("ONE FISH", body); |
| ... | @@ -954,15 +939,14 @@ fn echoTests(client: *http.Client, port: u16) !void { | ... | @@ -954,15 +939,14 @@ fn echoTests(client: *http.Client, port: u16) !void { |
| 954 | 939 | ||
| 955 | req.transfer_encoding = .{ .content_length = 14 }; | 940 | req.transfer_encoding = .{ .content_length = 14 }; |
| 956 | 941 | ||
| 957 | var body_writer = try req.sendBody(); | 942 | var body_writer = try req.sendBody(&.{}); |
| 958 | var w = body_writer.writer().unbuffered(); | 943 | try body_writer.writer.writeAll("Hello, "); |
| 959 | try w.writeAll("Hello, "); | 944 | try body_writer.writer.writeAll("World!\n"); |
| 960 | try w.writeAll("World!\n"); | ||
| 961 | try body_writer.end(); | 945 | try body_writer.end(); |
| 962 | 946 | ||
| 963 | var response = try req.receiveHead(&redirect_buffer); | 947 | var response = try req.receiveHead(&redirect_buffer); |
| 964 | 948 | ||
| 965 | const body = try response.reader().readRemainingAlloc(gpa, .limited(8192)); | 949 | const body = try response.reader(&.{}).allocRemaining(gpa, .limited(8192)); |
| 966 | defer gpa.free(body); | 950 | defer gpa.free(body); |
| 967 | 951 | ||
| 968 | try expectEqualStrings("Hello, World!\n", body); | 952 | try expectEqualStrings("Hello, World!\n", body); |
| ... | @@ -988,15 +972,14 @@ fn echoTests(client: *http.Client, port: u16) !void { | ... | @@ -988,15 +972,14 @@ fn echoTests(client: *http.Client, port: u16) !void { |
| 988 | 972 | ||
| 989 | req.transfer_encoding = .chunked; | 973 | req.transfer_encoding = .chunked; |
| 990 | 974 | ||
| 991 | var body_writer = try req.sendBody(); | 975 | var body_writer = try req.sendBody(&.{}); |
| 992 | var w = body_writer.writer().unbuffered(); | 976 | try body_writer.writer.writeAll("Hello, "); |
| 993 | try w.writeAll("Hello, "); | 977 | try body_writer.writer.writeAll("World!\n"); |
| 994 | try w.writeAll("World!\n"); | ||
| 995 | try body_writer.end(); | 978 | try body_writer.end(); |
| 996 | 979 | ||
| 997 | var response = try req.receiveHead(&redirect_buffer); | 980 | var response = try req.receiveHead(&redirect_buffer); |
| 998 | 981 | ||
| 999 | const body = try response.reader().readRemainingAlloc(gpa, .limited(8192)); | 982 | const body = try response.reader(&.{}).allocRemaining(gpa, .limited(8192)); |
| 1000 | defer gpa.free(body); | 983 | defer gpa.free(body); |
| 1001 | 984 | ||
| 1002 | try expectEqualStrings("Hello, World!\n", body); | 985 | try expectEqualStrings("Hello, World!\n", body); |
| ... | @@ -1042,16 +1025,15 @@ fn echoTests(client: *http.Client, port: u16) !void { | ... | @@ -1042,16 +1025,15 @@ fn echoTests(client: *http.Client, port: u16) !void { |
| 1042 | 1025 | ||
| 1043 | req.transfer_encoding = .chunked; | 1026 | req.transfer_encoding = .chunked; |
| 1044 | 1027 | ||
| 1045 | var body_writer = try req.sendBody(); | 1028 | var body_writer = try req.sendBody(&.{}); |
| 1046 | var w = body_writer.writer().unbuffered(); | 1029 | try body_writer.writer.writeAll("Hello, "); |
| 1047 | try w.writeAll("Hello, "); | 1030 | try body_writer.writer.writeAll("World!\n"); |
| 1048 | try w.writeAll("World!\n"); | ||
| 1049 | try body_writer.end(); | 1031 | try body_writer.end(); |
| 1050 | 1032 | ||
| 1051 | var response = try req.receiveHead(&redirect_buffer); | 1033 | var response = try req.receiveHead(&redirect_buffer); |
| 1052 | try expectEqual(.ok, response.head.status); | 1034 | try expectEqual(.ok, response.head.status); |
| 1053 | 1035 | ||
| 1054 | const body = try response.reader().readRemainingAlloc(gpa, .limited(8192)); | 1036 | const body = try response.reader(&.{}).allocRemaining(gpa, .limited(8192)); |
| 1055 | defer gpa.free(body); | 1037 | defer gpa.free(body); |
| 1056 | 1038 | ||
| 1057 | try expectEqualStrings("Hello, World!\n", body); | 1039 | try expectEqualStrings("Hello, World!\n", body); |
| ... | @@ -1073,11 +1055,11 @@ fn echoTests(client: *http.Client, port: u16) !void { | ... | @@ -1073,11 +1055,11 @@ fn echoTests(client: *http.Client, port: u16) !void { |
| 1073 | 1055 | ||
| 1074 | req.transfer_encoding = .chunked; | 1056 | req.transfer_encoding = .chunked; |
| 1075 | 1057 | ||
| 1076 | var body_writer = try req.sendBody(); | 1058 | var body_writer = try req.sendBody(&.{}); |
| 1077 | try body_writer.flush(); | 1059 | try body_writer.flush(); |
| 1078 | var response = try req.receiveHead(&redirect_buffer); | 1060 | var response = try req.receiveHead(&redirect_buffer); |
| 1079 | try expectEqual(.expectation_failed, response.head.status); | 1061 | try expectEqual(.expectation_failed, response.head.status); |
| 1080 | _ = try response.reader().discardRemaining(); | 1062 | _ = try response.reader(&.{}).discardRemaining(); |
| 1081 | } | 1063 | } |
| 1082 | } | 1064 | } |
| 1083 | 1065 | ||
| ... | @@ -1128,11 +1110,9 @@ test "redirect to different connection" { | ... | @@ -1128,11 +1110,9 @@ test "redirect to different connection" { |
| 1128 | const connection = try net_server.accept(); | 1110 | const connection = try net_server.accept(); |
| 1129 | defer connection.stream.close(); | 1111 | defer connection.stream.close(); |
| 1130 | 1112 | ||
| 1131 | var stream_reader = connection.stream.reader(); | 1113 | var connection_br = connection.stream.reader(&recv_buffer); |
| 1132 | var stream_writer = connection.stream.writer(); | 1114 | var connection_bw = connection.stream.writer(&send_buffer); |
| 1133 | var connection_br = stream_reader.interface().buffered(&recv_buffer); | 1115 | var server = http.Server.init(connection_br.interface(), &connection_bw.interface); |
| 1134 | var connection_bw = stream_writer.interface().buffered(&send_buffer); | ||
| 1135 | var server = http.Server.init(&connection_br, &connection_bw); | ||
| 1136 | var request = try server.receiveHead(); | 1116 | var request = try server.receiveHead(); |
| 1137 | try expectEqualStrings(request.head.target, "/ok"); | 1117 | try expectEqualStrings(request.head.target, "/ok"); |
| 1138 | try request.respond("good job, you pass", .{}); | 1118 | try request.respond("good job, you pass", .{}); |
| ... | @@ -1161,7 +1141,7 @@ test "redirect to different connection" { | ... | @@ -1161,7 +1141,7 @@ test "redirect to different connection" { |
| 1161 | 1141 | ||
| 1162 | var connection_br = connection.stream.reader(&recv_buffer); | 1142 | var connection_br = connection.stream.reader(&recv_buffer); |
| 1163 | var connection_bw = connection.stream.writer(&send_buffer); | 1143 | var connection_bw = connection.stream.writer(&send_buffer); |
| 1164 | var server = http.Server.init(&connection_br, &connection_bw); | 1144 | var server = http.Server.init(connection_br.interface(), &connection_bw.interface); |
| 1165 | var request = try server.receiveHead(); | 1145 | var request = try server.receiveHead(); |
| 1166 | try expectEqualStrings(request.head.target, "/help"); | 1146 | try expectEqualStrings(request.head.target, "/help"); |
| 1167 | try request.respond("", .{ | 1147 | try request.respond("", .{ |
lib/std/io/Reader.zig+4-1| ... | @@ -96,7 +96,10 @@ pub const failing: Reader = .{ | ... | @@ -96,7 +96,10 @@ pub const failing: Reader = .{ |
| 96 | .end = 0, | 96 | .end = 0, |
| 97 | }; | 97 | }; |
| 98 | 98 | ||
| 99 | pub const ending: Reader = .fixed(&.{}); | 99 | /// This is generally safe to `@constCast` because it has an empty buffer, so |
| 100 | /// there is not really a way to accidentally attempt mutation of these fields. | ||
| 101 | const ending_state: Reader = .fixed(&.{}); | ||
| 102 | pub const ending: *Reader = @constCast(&ending_state); | ||
| 100 | 103 | ||
| 101 | pub fn limited(r: *Reader, limit: Limit, buffer: []u8) Limited { | 104 | pub fn limited(r: *Reader, limit: Limit, buffer: []u8) Limited { |
| 102 | return Limited.init(r, limit, buffer); | 105 | return Limited.init(r, limit, buffer); |
lib/std/io/Writer.zig+79-10| ... | @@ -221,11 +221,12 @@ pub fn writeSplatLimit( | ... | @@ -221,11 +221,12 @@ pub fn writeSplatLimit( |
| 221 | /// `end`. | 221 | /// `end`. |
| 222 | pub fn flush(w: *Writer) Error!void { | 222 | pub fn flush(w: *Writer) Error!void { |
| 223 | assert(0 == try w.vtable.drain(w, &.{}, 0)); | 223 | assert(0 == try w.vtable.drain(w, &.{}, 0)); |
| 224 | if (w.end != 0) assert(w.vtable.drain == &fixedDrain); | ||
| 224 | } | 225 | } |
| 225 | 226 | ||
| 226 | /// Calls `VTable.drain` but hides the last `preserve_length` bytes from the | 227 | /// Calls `VTable.drain` but hides the last `preserve_length` bytes from the |
| 227 | /// implementation, keeping them buffered. | 228 | /// implementation, keeping them buffered. |
| 228 | pub fn drainLimited(w: *Writer, preserve_length: usize) Error!void { | 229 | pub fn drainPreserve(w: *Writer, preserve_length: usize) Error!void { |
| 229 | const temp_end = w.end -| preserve_length; | 230 | const temp_end = w.end -| preserve_length; |
| 230 | const preserved = w.buffer[temp_end..w.end]; | 231 | const preserved = w.buffer[temp_end..w.end]; |
| 231 | w.end = temp_end; | 232 | w.end = temp_end; |
| ... | @@ -235,6 +236,67 @@ pub fn drainLimited(w: *Writer, preserve_length: usize) Error!void { | ... | @@ -235,6 +236,67 @@ pub fn drainLimited(w: *Writer, preserve_length: usize) Error!void { |
| 235 | @memmove(w.buffer[w.end..][0..preserved.len], preserved); | 236 | @memmove(w.buffer[w.end..][0..preserved.len], preserved); |
| 236 | } | 237 | } |
| 237 | 238 | ||
| 239 | /// Forwards a `drain` to a second `Writer` instance. `w` is only used for its | ||
| 240 | /// buffer, but it has its `end` and `count` adjusted accordingly depending on | ||
| 241 | /// how much was consumed. | ||
| 242 | /// | ||
| 243 | /// Returns how many bytes from `data` were consumed. | ||
| 244 | pub fn drainTo(noalias w: *Writer, noalias other: *Writer, data: []const []const u8, splat: usize) Error!usize { | ||
| 245 | assert(w != other); | ||
| 246 | const header = w.buffered(); | ||
| 247 | const new_end = other.end + header.len; | ||
| 248 | if (new_end <= other.buffer.len) { | ||
| 249 | @memcpy(other.buffer[other.end..][0..header.len], header); | ||
| 250 | other.end = new_end; | ||
| 251 | other.count += header.len; | ||
| 252 | w.end = 0; | ||
| 253 | const n = try other.vtable.drain(other, data, splat); | ||
| 254 | other.count += n; | ||
| 255 | return n; | ||
| 256 | } | ||
| 257 | if (other.vtable == &VectorWrapper.vtable) { | ||
| 258 | const wrapper: *VectorWrapper = @fieldParentPtr("writer", w); | ||
| 259 | while (wrapper.it.next()) |dest| { | ||
| 260 | _ = dest; | ||
| 261 | @panic("TODO"); | ||
| 262 | } | ||
| 263 | } | ||
| 264 | var vecs: [8][]const u8 = undefined; // Arbitrarily chosen size. | ||
| 265 | var i: usize = 1; | ||
| 266 | vecs[0] = header; | ||
| 267 | for (data) |buf| { | ||
| 268 | if (buf.len == 0) continue; | ||
| 269 | vecs[i] = buf; | ||
| 270 | i += 1; | ||
| 271 | if (vecs.len - i == 0) break; | ||
| 272 | } | ||
| 273 | const new_splat = if (vecs[i - 1].ptr == data[data.len - 1].ptr) splat else 1; | ||
| 274 | const n = try other.vtable.drain(other, vecs[0..i], new_splat); | ||
| 275 | other.count += n; | ||
| 276 | if (n < header.len) { | ||
| 277 | const remaining = w.buffer[n..w.end]; | ||
| 278 | @memmove(w.buffer[0..remaining.len], remaining); | ||
| 279 | w.end = remaining.len; | ||
| 280 | return 0; | ||
| 281 | } | ||
| 282 | defer w.end = 0; | ||
| 283 | return n - header.len; | ||
| 284 | } | ||
| 285 | |||
| 286 | pub fn drainToLimit( | ||
| 287 | noalias w: *Writer, | ||
| 288 | noalias other: *Writer, | ||
| 289 | data: []const []const u8, | ||
| 290 | splat: usize, | ||
| 291 | limit: Limit, | ||
| 292 | ) Error!usize { | ||
| 293 | assert(w != other); | ||
| 294 | _ = data; | ||
| 295 | _ = splat; | ||
| 296 | _ = limit; | ||
| 297 | @panic("TODO"); | ||
| 298 | } | ||
| 299 | |||
| 238 | pub fn unusedCapacitySlice(w: *const Writer) []u8 { | 300 | pub fn unusedCapacitySlice(w: *const Writer) []u8 { |
| 239 | return w.buffer[w.end..]; | 301 | return w.buffer[w.end..]; |
| 240 | } | 302 | } |
| ... | @@ -285,10 +347,10 @@ pub fn writableSliceGreedy(w: *Writer, minimum_length: usize) Error![]u8 { | ... | @@ -285,10 +347,10 @@ pub fn writableSliceGreedy(w: *Writer, minimum_length: usize) Error![]u8 { |
| 285 | /// remain buffered. | 347 | /// remain buffered. |
| 286 | /// | 348 | /// |
| 287 | /// If `preserve_length` is zero, this is equivalent to `writableSliceGreedy`. | 349 | /// If `preserve_length` is zero, this is equivalent to `writableSliceGreedy`. |
| 288 | pub fn writableSliceGreedyPreserving(w: *Writer, preserve_length: usize, minimum_length: usize) Error![]u8 { | 350 | pub fn writableSliceGreedyPreserve(w: *Writer, preserve_length: usize, minimum_length: usize) Error![]u8 { |
| 289 | assert(w.buffer.len >= preserve_length + minimum_length); | 351 | assert(w.buffer.len >= preserve_length + minimum_length); |
| 290 | while (w.buffer.len - w.end < minimum_length) { | 352 | while (w.buffer.len - w.end < minimum_length) { |
| 291 | try drainLimited(w, preserve_length); | 353 | try drainPreserve(w, preserve_length); |
| 292 | } else { | 354 | } else { |
| 293 | @branchHint(.likely); | 355 | @branchHint(.likely); |
| 294 | return w.buffer[w.end..]; | 356 | return w.buffer[w.end..]; |
| ... | @@ -444,7 +506,7 @@ pub fn write(w: *Writer, bytes: []const u8) Error!usize { | ... | @@ -444,7 +506,7 @@ pub fn write(w: *Writer, bytes: []const u8) Error!usize { |
| 444 | } | 506 | } |
| 445 | 507 | ||
| 446 | /// Asserts `buffer` capacity exceeds `preserve_length`. | 508 | /// Asserts `buffer` capacity exceeds `preserve_length`. |
| 447 | pub fn writePreserving(w: *Writer, preserve_length: usize, bytes: []const u8) Error!usize { | 509 | pub fn writePreserve(w: *Writer, preserve_length: usize, bytes: []const u8) Error!usize { |
| 448 | assert(preserve_length <= w.buffer.len); | 510 | assert(preserve_length <= w.buffer.len); |
| 449 | if (w.end + bytes.len <= w.buffer.len) { | 511 | if (w.end + bytes.len <= w.buffer.len) { |
| 450 | @branchHint(.likely); | 512 | @branchHint(.likely); |
| ... | @@ -478,9 +540,9 @@ pub fn writeAll(w: *Writer, bytes: []const u8) Error!void { | ... | @@ -478,9 +540,9 @@ pub fn writeAll(w: *Writer, bytes: []const u8) Error!void { |
| 478 | /// remain buffered. | 540 | /// remain buffered. |
| 479 | /// | 541 | /// |
| 480 | /// Asserts `buffer` capacity exceeds `preserve_length`. | 542 | /// Asserts `buffer` capacity exceeds `preserve_length`. |
| 481 | pub fn writeAllPreserving(w: *Writer, preserve_length: usize, bytes: []const u8) Error!void { | 543 | pub fn writeAllPreserve(w: *Writer, preserve_length: usize, bytes: []const u8) Error!void { |
| 482 | var index: usize = 0; | 544 | var index: usize = 0; |
| 483 | while (index < bytes.len) index += try w.writePreserving(preserve_length, bytes[index..]); | 545 | while (index < bytes.len) index += try w.writePreserve(preserve_length, bytes[index..]); |
| 484 | } | 546 | } |
| 485 | 547 | ||
| 486 | pub fn print(w: *Writer, comptime format: []const u8, args: anytype) Error!void { | 548 | pub fn print(w: *Writer, comptime format: []const u8, args: anytype) Error!void { |
| ... | @@ -505,9 +567,9 @@ pub fn writeByte(w: *Writer, byte: u8) Error!void { | ... | @@ -505,9 +567,9 @@ pub fn writeByte(w: *Writer, byte: u8) Error!void { |
| 505 | 567 | ||
| 506 | /// When draining the buffer, ensures that at least `preserve_length` bytes | 568 | /// When draining the buffer, ensures that at least `preserve_length` bytes |
| 507 | /// remain buffered. | 569 | /// remain buffered. |
| 508 | pub fn writeBytePreserving(w: *Writer, preserve_length: usize, byte: u8) Error!void { | 570 | pub fn writeBytePreserve(w: *Writer, preserve_length: usize, byte: u8) Error!void { |
| 509 | while (w.buffer.len - w.end == 0) { | 571 | while (w.buffer.len - w.end == 0) { |
| 510 | try drainLimited(w, preserve_length); | 572 | try drainPreserve(w, preserve_length); |
| 511 | } else { | 573 | } else { |
| 512 | @branchHint(.likely); | 574 | @branchHint(.likely); |
| 513 | w.buffer[w.end] = byte; | 575 | w.buffer[w.end] = byte; |
| ... | @@ -615,19 +677,26 @@ pub fn sendFile(w: *Writer, file_reader: *File.Reader, limit: Limit) FileError!u | ... | @@ -615,19 +677,26 @@ pub fn sendFile(w: *Writer, file_reader: *File.Reader, limit: Limit) FileError!u |
| 615 | /// on how much was consumed. | 677 | /// on how much was consumed. |
| 616 | /// | 678 | /// |
| 617 | /// Returns how many bytes from `file_reader` were consumed. | 679 | /// Returns how many bytes from `file_reader` were consumed. |
| 618 | pub fn sendFileTo(w: *Writer, other: *Writer, file_reader: *File.Reader, limit: Limit) FileError!usize { | 680 | pub fn sendFileTo( |
| 681 | noalias w: *Writer, | ||
| 682 | noalias other: *Writer, | ||
| 683 | file_reader: *File.Reader, | ||
| 684 | limit: Limit, | ||
| 685 | ) FileError!usize { | ||
| 686 | assert(w != other); | ||
| 619 | const header = w.buffered(); | 687 | const header = w.buffered(); |
| 620 | const new_end = other.end + header.len; | 688 | const new_end = other.end + header.len; |
| 621 | if (new_end <= other.buffer.len) { | 689 | if (new_end <= other.buffer.len) { |
| 622 | @memcpy(other.buffer[other.end..][0..header.len], header); | 690 | @memcpy(other.buffer[other.end..][0..header.len], header); |
| 623 | other.end = new_end; | 691 | other.end = new_end; |
| 692 | other.count += header.len; | ||
| 624 | w.end = 0; | 693 | w.end = 0; |
| 625 | return other.vtable.sendFile(other, file_reader, limit); | 694 | return other.vtable.sendFile(other, file_reader, limit); |
| 626 | } | 695 | } |
| 627 | assert(header.len > 0); | 696 | assert(header.len > 0); |
| 628 | var vec_buf: [2][]const u8 = .{ header, undefined }; | 697 | var vec_buf: [2][]const u8 = .{ header, undefined }; |
| 629 | var vec_i: usize = 1; | 698 | var vec_i: usize = 1; |
| 630 | const buffered_contents = limit.slice(file_reader.buffered()); | 699 | const buffered_contents = limit.slice(file_reader.interface.buffered()); |
| 631 | if (buffered_contents.len > 0) { | 700 | if (buffered_contents.len > 0) { |
| 632 | vec_buf[vec_i] = buffered_contents; | 701 | vec_buf[vec_i] = buffered_contents; |
| 633 | vec_i += 1; | 702 | vec_i += 1; |