lib/http/src/client/stream/model.zig

daab053ee43316e1809a84551d573ddd1e5bf3d2

  1 const std = @import("std");
  2 const alloc_phase = @import("alloc_phase");
  3 const http = @import("../../root.zig");
  4 const response = @import("../root.zig").response;
  5 const chunks = http.chunk;
  6 
  7 pub const default_header_count: usize = response.default_header_count;
  8 pub const default_head_bytes: usize = response.default_head_bytes;
  9 
 10 pub const Header = response.Header;
 11 
 12 pub const Limits = struct {
 13     stream_count: usize,
 14     header_count_per_stream: usize,
 15     head_bytes_per_stream: usize,
 16 };
 17 
 18 pub const Capacity = struct {
 19     stream_count: usize,
 20     header_count_per_stream: usize,
 21     head_bytes_per_stream: usize,
 22     header_count: usize,
 23     header_bytes: usize,
 24     head_bytes: usize,
 25     storage_bytes: usize,
 26 
 27     pub fn derive(limits: Limits) error{CapacityOverflow}!Capacity {
 28         const header_count = try alloc_phase.capacity.mul(
 29             usize,
 30             limits.stream_count,
 31             limits.header_count_per_stream,
 32         );
 33         const header_bytes = try alloc_phase.capacity.mul(
 34             usize,
 35             header_count,
 36             @sizeOf(Header),
 37         );
 38         const head_bytes = try alloc_phase.capacity.mul(
 39             usize,
 40             limits.stream_count,
 41             limits.head_bytes_per_stream,
 42         );
 43         const storage_bytes = try alloc_phase.capacity.add(
 44             usize,
 45             header_bytes,
 46             head_bytes,
 47         );
 48         return .{
 49             .stream_count = limits.stream_count,
 50             .header_count_per_stream = limits.header_count_per_stream,
 51             .head_bytes_per_stream = limits.head_bytes_per_stream,
 52             .header_count = header_count,
 53             .header_bytes = header_bytes,
 54             .head_bytes = head_bytes,
 55             .storage_bytes = storage_bytes,
 56         };
 57     }
 58 };
 59 
 60 pub const Scratch = struct {
 61     headers: []Header,
 62     head: []u8,
 63 };
 64 
 65 pub const Error = error{
 66     MalformedResponse,
 67     ReadFailed,
 68     StreamHeaderCapacityExceeded,
 69     StreamHeadCapacityExceeded,
 70 };
 71 
 72 pub const HandlerVTable = struct {
 73     onHead: *const fn (ctx: *anyopaque, status: u16, headers: []const Header) anyerror!void,
 74     onChunk: *const fn (ctx: *anyopaque, chunk: []const u8) anyerror!void,
 75 };
 76 
 77 pub const Handler = struct {
 78     ctx: *anyopaque,
 79     vtable: *const HandlerVTable,
 80 
 81     pub fn head(self: Handler, status: u16, headers: []const Header) !void {
 82         return self.vtable.onHead(self.ctx, status, headers);
 83     }
 84 
 85     pub fn emit(self: Handler, chunk: []const u8) !void {
 86         return self.vtable.onChunk(self.ctx, chunk);
 87     }
 88 };
 89 
 90 pub const ClientStream = struct {
 91     pub fn read(scratch: Scratch, reader: *std.Io.Reader, handler: Handler) anyerror!u16 {
 92         return readWithMethod(scratch, reader, null, handler);
 93     }
 94 
 95     pub fn readForMethod(
 96         scratch: Scratch,
 97         reader: *std.Io.Reader,
 98         method: []const u8,
 99         handler: Handler,
100     ) anyerror!u16 {
101         return readWithMethod(scratch, reader, method, handler);
102     }
103 };
104 
105 fn readWithMethod(
106     scratch: Scratch,
107     reader: *std.Io.Reader,
108     method: ?[]const u8,
109     handler: Handler,
110 ) anyerror!u16 {
111     var head_length: usize = 0;
112     var head_survey: ?response.HeadSurvey = null;
113     var fixed_remaining: usize = 0;
114     var decoder = chunks.Decoder.init();
115 
116     while (true) {
117         const read_buffer = readAvailable(reader) catch {
118             return error.ReadFailed;
119         };
120         const read_length = read_buffer.len;
121         if (read_length == 0) {
122             const head = head_survey orelse {
123                 if (head_length == 0) return error.ReadFailed;
124                 return error.MalformedResponse;
125             };
126             return switch (head.framing) {
127                 .none, .tunnel => head.status,
128                 .fixed => if (fixed_remaining == 0) head.status else error.MalformedResponse,
129                 .chunked => if (decoder.done) head.status else error.MalformedResponse,
130                 .close => head.status,
131             };
132         }
133         var toss_length = read_length;
134         defer reader.toss(toss_length);
135 
136         var offset: usize = 0;
137         if (head_survey == null) {
138             const previous_head_length = head_length;
139             const take = @min(scratch.head.len - head_length, read_length);
140             @memcpy(scratch.head[head_length..][0..take], read_buffer[0..take]);
141             const combined_length = head_length + take;
142             const search_start = head_length - @min(head_length, 3);
143             const relative_end = std.mem.indexOf(
144                 u8,
145                 scratch.head[search_start..combined_length],
146                 "\r\n\r\n",
147             );
148             if (relative_end == null) {
149                 head_length = combined_length;
150                 if (take != read_length) return error.StreamHeadCapacityExceeded;
151                 continue;
152             }
153             head_length = search_start + relative_end.? + 4;
154             offset = head_length - previous_head_length;
155             const head = response.surveyHead(scratch.head[0..head_length], method) catch {
156                 return error.MalformedResponse;
157             };
158             if (head.header_count > scratch.headers.len) {
159                 return error.StreamHeaderCapacityExceeded;
160             }
161             const headers = response.fillHeaders(
162                 scratch.headers,
163                 scratch.head[0..head_length],
164                 head.header_count,
165             );
166             switch (head.framing) {
167                 .none, .tunnel => toss_length = offset,
168                 .fixed => |expected| if (expected == 0) {
169                     toss_length = offset;
170                 },
171                 .chunked, .close => {},
172             }
173             try handler.head(head.status, headers);
174             head_survey = head;
175             switch (head.framing) {
176                 .none, .tunnel => {
177                     toss_length = offset;
178                     return head.status;
179                 },
180                 .fixed => |expected| {
181                     fixed_remaining = expected;
182                     if (expected == 0) {
183                         toss_length = offset;
184                         return head.status;
185                     }
186                 },
187                 .chunked, .close => {},
188             }
189         }
190 
191         const head = head_survey orelse continue;
192         const input = read_buffer[offset..read_length];
193         switch (head.framing) {
194             .none, .tunnel => return head.status,
195             .fixed => {
196                 const take = @min(fixed_remaining, input.len);
197                 if (take != 0) try handler.emit(input[0..take]);
198                 fixed_remaining -= take;
199                 if (fixed_remaining == 0) {
200                     toss_length = offset + take;
201                     return head.status;
202                 }
203             },
204             .chunked => {
205                 const consumed = decoder.feed(input, handler) catch |err| switch (err) {
206                     error.MalformedBody => return error.MalformedResponse,
207                     else => return err,
208                 };
209                 if (decoder.done) {
210                     toss_length = offset + consumed;
211                     return head.status;
212                 }
213                 std.debug.assert(consumed == input.len);
214             },
215             .close => if (input.len != 0) try handler.emit(input),
216         }
217     }
218 }
219 
220 fn readAvailable(reader: *std.Io.Reader) ![]const u8 {
221     return reader.peekGreedy(1) catch |err| switch (err) {
222         error.EndOfStream => return &.{},
223         else => return err,
224     };
225 }
226 
227 fn independentCapacity(limits: Limits) error{CapacityOverflow}!Capacity {
228     const header_count = @as(u128, limits.stream_count) * limits.header_count_per_stream;
229     const header_bytes = header_count * @sizeOf(Header);
230     const head_bytes = @as(u128, limits.stream_count) * limits.head_bytes_per_stream;
231     const storage_bytes = header_bytes + head_bytes;
232     if (header_count > std.math.maxInt(usize) or
233         header_bytes > std.math.maxInt(usize) or
234         head_bytes > std.math.maxInt(usize) or
235         storage_bytes > std.math.maxInt(usize))
236     {
237         return error.CapacityOverflow;
238     }
239     return .{
240         .stream_count = limits.stream_count,
241         .header_count_per_stream = limits.header_count_per_stream,
242         .head_bytes_per_stream = limits.head_bytes_per_stream,
243         .header_count = @intCast(header_count),
244         .header_bytes = @intCast(header_bytes),
245         .head_bytes = @intCast(head_bytes),
246         .storage_bytes = @intCast(storage_bytes),
247     };
248 }
249 
250 const TestHandler = struct {
251     status: u16 = 0,
252     head_count: usize = 0,
253     header_count: usize = 0,
254     body: [128]u8 = undefined,
255     body_length: usize = 0,
256 
257     fn handler(self: *TestHandler) Handler {
258         return .{ .ctx = self, .vtable = &vtable };
259     }
260 
261     fn onHead(context: *anyopaque, status: u16, headers: []const Header) anyerror!void {
262         const self: *TestHandler = @ptrCast(@alignCast(context));
263         self.status = status;
264         self.head_count += 1;
265         self.header_count = headers.len;
266     }
267 
268     fn onChunk(context: *anyopaque, chunk: []const u8) anyerror!void {
269         const self: *TestHandler = @ptrCast(@alignCast(context));
270         if (chunk.len > self.body.len - self.body_length) return error.TestBodyCapacityExceeded;
271         @memcpy(self.body[self.body_length..][0..chunk.len], chunk);
272         self.body_length += chunk.len;
273     }
274 
275     const vtable = HandlerVTable{
276         .onHead = &onHead,
277         .onChunk = &onChunk,
278     };
279 };
280 
281 const RejectingHead = struct {
282     fn head(_: *anyopaque, _: u16, _: []const Header) anyerror!void {
283         return error.HeadRejected;
284     }
285 
286     fn chunk(_: *anyopaque, _: []const u8) anyerror!void {
287         return error.BodyObservedAfterHeadFailure;
288     }
289 
290     const vtable = HandlerVTable{ .onHead = &head, .onChunk = &chunk };
291 };
292 
293 const RejectingBody = struct {
294     fn head(_: *anyopaque, _: u16, _: []const Header) anyerror!void {}
295 
296     fn chunk(_: *anyopaque, _: []const u8) anyerror!void {
297         return error.BodyRejected;
298     }
299 
300     const vtable = HandlerVTable{ .onHead = &head, .onChunk = &chunk };
301 };
302 
303 const FragmentedReader = struct {
304     data: []const u8,
305     offset: usize = 0,
306     fragment_length: usize,
307     storage: [8]u8 = undefined,
308     reader: std.Io.Reader = .{
309         .vtable = &vtable,
310         .buffer = undefined,
311         .seek = 0,
312         .end = 0,
313     },
314 
315     fn attach(self: *FragmentedReader) void {
316         self.reader.buffer = &self.storage;
317     }
318 
319     const vtable: std.Io.Reader.VTable = .{
320         .stream = stream,
321         .discard = discard,
322         .readVec = readVec,
323     };
324 
325     fn stream(
326         reader: *std.Io.Reader,
327         writer: *std.Io.Writer,
328         limit: std.Io.Limit,
329     ) std.Io.Reader.StreamError!usize {
330         const self: *FragmentedReader = @alignCast(@fieldParentPtr("reader", reader));
331         const buffer = limit.slice(reader.buffer);
332         if (buffer.len == 0) return 0;
333         const length = self.readInto(buffer);
334         if (length == 0) return error.EndOfStream;
335         writer.writeAll(buffer[0..length]) catch return error.WriteFailed;
336         return length;
337     }
338 
339     fn discard(reader: *std.Io.Reader, limit: std.Io.Limit) std.Io.Reader.Error!usize {
340         const self: *FragmentedReader = @alignCast(@fieldParentPtr("reader", reader));
341         const buffer = limit.slice(reader.buffer);
342         if (buffer.len == 0) return 0;
343         const length = self.readInto(buffer);
344         if (length == 0) return error.EndOfStream;
345         return length;
346     }
347 
348     fn readVec(reader: *std.Io.Reader, data: [][]u8) std.Io.Reader.Error!usize {
349         const self: *FragmentedReader = @alignCast(@fieldParentPtr("reader", reader));
350         for (data) |buffer| {
351             if (buffer.len == 0) continue;
352             const length = self.readInto(buffer);
353             if (length == 0) return error.EndOfStream;
354             return length;
355         }
356         if (reader.buffer.len == 0) return 0;
357         const destination = reader.buffer[reader.end..];
358         if (destination.len == 0) return 0;
359         const length = self.readInto(destination);
360         if (length == 0) return error.EndOfStream;
361         reader.end += length;
362         return 0;
363     }
364 
365     fn readInto(self: *FragmentedReader, destination: []u8) usize {
366         const length = @min(
367             self.fragment_length,
368             @min(destination.len, self.data.len - self.offset),
369         );
370         @memcpy(destination[0..length], self.data[self.offset..][0..length]);
371         self.offset += length;
372         return length;
373     }
374 };
375 
376 test "Client stream capacity matches independent arithmetic" {
377     comptime {
378         @stardustClaim(
379             @import("alloc_phase").capacity.witness(@import("./root.zig").ClientStreamStorage, "http_client_stream_capacity"),
380             null,
381             null,
382             null,
383             null,
384             null,
385             null,
386         );
387     }
388 
389     for (0..9) |stream_count| {
390         for (0..9) |header_count_per_stream| {
391             for (0..9) |head_bytes_per_stream| {
392                 const limits = Limits{
393                     .stream_count = stream_count,
394                     .header_count_per_stream = header_count_per_stream,
395                     .head_bytes_per_stream = head_bytes_per_stream,
396                 };
397                 try std.testing.expectEqual(
398                     try independentCapacity(limits),
399                     try Capacity.derive(limits),
400                 );
401             }
402         }
403     }
404     const maximum = std.math.maxInt(usize);
405     try std.testing.expectError(error.CapacityOverflow, Capacity.derive(.{
406         .stream_count = 2,
407         .header_count_per_stream = maximum,
408         .head_bytes_per_stream = 0,
409     }));
410     try std.testing.expectError(error.CapacityOverflow, Capacity.derive(.{
411         .stream_count = 1,
412         .header_count_per_stream = maximum,
413         .head_bytes_per_stream = 0,
414     }));
415     try std.testing.expectError(error.CapacityOverflow, Capacity.derive(.{
416         .stream_count = maximum,
417         .header_count_per_stream = 0,
418         .head_bytes_per_stream = 2,
419     }));
420     try std.testing.expectError(error.CapacityOverflow, Capacity.derive(.{
421         .stream_count = 1,
422         .header_count_per_stream = 1,
423         .head_bytes_per_stream = maximum,
424     }));
425 }
426 
427 test "Client stream accepts exact capacities" {
428     comptime {
429         @stardustClaim(
430             @import("alloc_phase").capacity.witness(@import("./root.zig").ClientStreamStorage, "http_client_stream_boundary"),
431             null,
432             null,
433             null,
434             null,
435             null,
436             null,
437         );
438     }
439 
440     const raw_head = "HTTP/1.1 200 OK\r\nContent-Type: text/plain\r\nContent-Length: 5\r\n\r\n";
441     const raw = raw_head ++ "Hello";
442     var headers: [2]Header = undefined;
443     var head: [raw_head.len]u8 = undefined;
444     var reader = std.Io.Reader.fixed(raw);
445     var handler = TestHandler{};
446     const status = try ClientStream.read(
447         .{ .headers = &headers, .head = &head },
448         &reader,
449         handler.handler(),
450     );
451     try std.testing.expectEqual(@as(u16, 200), status);
452     try std.testing.expectEqual(@as(u16, 200), handler.status);
453     try std.testing.expectEqual(@as(usize, 1), handler.head_count);
454     try std.testing.expectEqual(@as(usize, 2), handler.header_count);
455     try std.testing.expectEqualStrings("Hello", handler.body[0..handler.body_length]);
456 }
457 
458 test "Client stream capacity failures publish no head" {
459     comptime {
460         @stardustClaim(
461             @import("alloc_phase").capacity.witness(@import("./root.zig").ClientStreamStorage, "http_client_stream_atomic"),
462             null,
463             null,
464             null,
465             null,
466             null,
467             null,
468         );
469     }
470 
471     const raw_head = "HTTP/1.1 200 OK\r\nContent-Type: text/plain\r\nContent-Length: 5\r\n\r\n";
472     const raw = raw_head ++ "Hello";
473     var headers: [2]Header = undefined;
474     var short_head: [raw_head.len - 1]u8 = undefined;
475     var head_reader = std.Io.Reader.fixed(raw);
476     var head_handler = TestHandler{};
477     try std.testing.expectError(
478         error.StreamHeadCapacityExceeded,
479         ClientStream.read(
480             .{ .headers = &headers, .head = &short_head },
481             &head_reader,
482             head_handler.handler(),
483         ),
484     );
485     try std.testing.expectEqual(@as(usize, 0), head_handler.head_count);
486 
487     var no_headers: [0]Header = .{};
488     var exact_head: [raw_head.len]u8 = undefined;
489     var header_reader = std.Io.Reader.fixed(raw);
490     var header_handler = TestHandler{};
491     try std.testing.expectError(
492         error.StreamHeaderCapacityExceeded,
493         ClientStream.read(
494             .{ .headers = &no_headers, .head = &exact_head },
495             &header_reader,
496             header_handler.handler(),
497         ),
498     );
499     try std.testing.expectEqual(@as(usize, 0), header_handler.head_count);
500 }
501 
502 test "Client stream preserves callback error identity" {
503     comptime {
504         @stardustClaim(
505             @import("alloc_phase").capacity.witness(@import("./root.zig").ClientStreamStorage, "http_client_stream_callback"),
506             null,
507             null,
508             null,
509             null,
510             null,
511             null,
512         );
513     }
514 
515     const raw = "HTTP/1.1 200 OK\r\nContent-Length: 5\r\n\r\nHello";
516     var headers: [1]Header = undefined;
517     var head: [64]u8 = undefined;
518     var context: u8 = 0;
519     var head_reader = std.Io.Reader.fixed(raw);
520     try std.testing.expectError(
521         error.HeadRejected,
522         ClientStream.read(
523             .{ .headers = &headers, .head = &head },
524             &head_reader,
525             .{ .ctx = &context, .vtable = &RejectingHead.vtable },
526         ),
527     );
528     var body_reader = std.Io.Reader.fixed(raw);
529     try std.testing.expectError(
530         error.BodyRejected,
531         ClientStream.read(
532             .{ .headers = &headers, .head = &head },
533             &body_reader,
534             .{ .ctx = &context, .vtable = &RejectingBody.vtable },
535         ),
536     );
537 
538     var tunnel_reader = std.Io.Reader.fixed(
539         "HTTP/1.1 200 Connected\r\n\r\nraw",
540     );
541     try std.testing.expectError(
542         error.HeadRejected,
543         ClientStream.readForMethod(
544             .{ .headers = &headers, .head = &head },
545             &tunnel_reader,
546             "CONNECT",
547             .{ .ctx = &context, .vtable = &RejectingHead.vtable },
548         ),
549     );
550     try std.testing.expectEqualStrings(
551         "raw",
552         try tunnel_reader.peekGreedy(1),
553     );
554 }
555 
556 test "Client stream omits HEAD response body" {
557     const raw = "HTTP/1.1 200 OK\r\nContent-Length: 5\r\n\r\nHello";
558     var headers: [1]Header = undefined;
559     var head: [64]u8 = undefined;
560     var reader = std.Io.Reader.fixed(raw);
561     var handler = TestHandler{};
562     const status = try ClientStream.readForMethod(
563         .{ .headers = &headers, .head = &head },
564         &reader,
565         "HEAD",
566         handler.handler(),
567     );
568     try std.testing.expectEqual(@as(u16, 200), status);
569     try std.testing.expectEqual(@as(usize, 1), handler.head_count);
570     try std.testing.expectEqual(@as(usize, 0), handler.body_length);
571 }
572 
573 test "Client stream rejects malformed framing before publishing a head" {
574     const malformed = [_][]const u8{
575         "HTTP/1.1 200 OK\r\nBadHeader\r\n\r\n",
576         "HTTP/1.1 200 OK\r\nContent-Length: 3\r\nContent-Length: 3\r\n\r\nabc",
577         "HTTP/1.1 200 OK\r\nTransfer-Encoding: gzip\r\nTransfer-Encoding: chunked\r\n\r\n0\r\n\r\n",
578         "HTTP/1.1 200 OK\r\nTransfer-Encoding: gzip\r\nContent-Length: 3\r\n\r\nabc",
579     };
580     var headers: [4]Header = undefined;
581     var head: [256]u8 = undefined;
582     for (malformed) |raw| {
583         var reader = std.Io.Reader.fixed(raw);
584         var handler = TestHandler{};
585         try std.testing.expectError(
586             error.MalformedResponse,
587             ClientStream.readForMethod(
588                 .{ .headers = &headers, .head = &head },
589                 &reader,
590                 "GET",
591                 handler.handler(),
592             ),
593         );
594         try std.testing.expectEqual(@as(usize, 0), handler.head_count);
595         try std.testing.expectEqual(@as(usize, 0), handler.body_length);
596     }
597 }
598 
599 test "Client stream honors nonfinal transfer and tunnel framing" {
600     var headers: [2]Header = undefined;
601     var head: [128]u8 = undefined;
602 
603     var close_reader = std.Io.Reader.fixed(
604         "HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked, gzip\r\n\r\nraw",
605     );
606     var close_handler = TestHandler{};
607     const close_status = try ClientStream.readForMethod(
608         .{ .headers = &headers, .head = &head },
609         &close_reader,
610         "GET",
611         close_handler.handler(),
612     );
613     try std.testing.expectEqual(@as(u16, 200), close_status);
614     try std.testing.expectEqualStrings(
615         "raw",
616         close_handler.body[0..close_handler.body_length],
617     );
618 
619     var tunnel_reader = std.Io.Reader.fixed(
620         "HTTP/1.1 200 Connected\r\nContent-Length: 3\r\n\r\nraw",
621     );
622     var tunnel_handler = TestHandler{};
623     const tunnel_status = try ClientStream.readForMethod(
624         .{ .headers = &headers, .head = &head },
625         &tunnel_reader,
626         "CONNECT",
627         tunnel_handler.handler(),
628     );
629     try std.testing.expectEqual(@as(u16, 200), tunnel_status);
630     try std.testing.expectEqual(@as(usize, 1), tunnel_handler.head_count);
631     try std.testing.expectEqual(@as(usize, 0), tunnel_handler.body_length);
632     try std.testing.expectEqualStrings(
633         "raw",
634         try tunnel_reader.peekGreedy(1),
635     );
636 }
637 
638 test "Client stream preserves bytes after fixed and chunked bodies" {
639     var headers: [2]Header = undefined;
640     var head: [128]u8 = undefined;
641 
642     var fixed_reader = std.Io.Reader.fixed(
643         "HTTP/1.1 200 OK\r\nContent-Length: 3\r\n\r\nabcTAIL",
644     );
645     var fixed_handler = TestHandler{};
646     const fixed_status = try ClientStream.readForMethod(
647         .{ .headers = &headers, .head = &head },
648         &fixed_reader,
649         "GET",
650         fixed_handler.handler(),
651     );
652     try std.testing.expectEqual(@as(u16, 200), fixed_status);
653     try std.testing.expectEqualStrings(
654         "abc",
655         fixed_handler.body[0..fixed_handler.body_length],
656     );
657     try std.testing.expectEqualStrings("TAIL", try fixed_reader.peekGreedy(1));
658 
659     var chunked_reader = std.Io.Reader.fixed(
660         "HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n3\r\nabc\r\n0\r\n\r\nTAIL",
661     );
662     var chunked_handler = TestHandler{};
663     const chunked_status = try ClientStream.readForMethod(
664         .{ .headers = &headers, .head = &head },
665         &chunked_reader,
666         "GET",
667         chunked_handler.handler(),
668     );
669     try std.testing.expectEqual(@as(u16, 200), chunked_status);
670     try std.testing.expectEqualStrings(
671         "abc",
672         chunked_handler.body[0..chunked_handler.body_length],
673     );
674     try std.testing.expectEqualStrings("TAIL", try chunked_reader.peekGreedy(1));
675 }
676 
677 test "Client stream preserves framed tails across fragmentation" {
678     const wires = [_][]const u8{
679         "HTTP/1.1 200 OK\r\nContent-Length: 3\r\n\r\nabcTAIL",
680         "HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n3\r\nabc\r\n0\r\n\r\nTAIL",
681     };
682     var headers: [2]Header = undefined;
683     var head: [128]u8 = undefined;
684 
685     for (wires) |wire| {
686         for (1..9) |fragment_length| {
687             var read_buffer: [32]u8 = undefined;
688             var source = std.testing.Reader.init(
689                 &read_buffer,
690                 &.{.{ .buffer = wire }},
691             );
692             source.artificial_limit = .limited(fragment_length);
693             var handler = TestHandler{};
694             const status = try ClientStream.readForMethod(
695                 .{ .headers = &headers, .head = &head },
696                 &source.interface,
697                 "GET",
698                 handler.handler(),
699             );
700             try std.testing.expectEqual(@as(u16, 200), status);
701             try std.testing.expectEqualStrings(
702                 "abc",
703                 handler.body[0..handler.body_length],
704             );
705             try std.testing.expectEqualStrings(
706                 "TAIL",
707                 try source.interface.peek(4),
708             );
709         }
710     }
711 }
712 
713 test "Client stream decodes every fragmented response byte" {
714     const raw_head = "HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\nContent-Type: text/plain\r\n\r\n";
715     const raw = raw_head ++ "4;kind=test\r\nWiki\r\n5\r\npedia\r\n0\r\nTrace: value\r\n\r\n";
716     var headers: [2]Header = undefined;
717     var head: [raw_head.len]u8 = undefined;
718     var reader = FragmentedReader{
719         .data = raw,
720         .fragment_length = 1,
721     };
722     reader.attach();
723     var handler = TestHandler{};
724     const status = try ClientStream.read(
725         .{ .headers = &headers, .head = &head },
726         &reader.reader,
727         handler.handler(),
728     );
729     try std.testing.expectEqual(@as(u16, 200), status);
730     try std.testing.expectEqual(@as(usize, 1), handler.head_count);
731     try std.testing.expectEqualStrings("Wikipedia", handler.body[0..handler.body_length]);
732 }
733 
734 test "Client stream preserves framing semantics across one-byte fragments" {
735     var headers: [2]Header = undefined;
736     var head: [128]u8 = undefined;
737 
738     var close_reader = FragmentedReader{
739         .data = "HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked, tiny-coding; name=\"a,b\"\r\n\r\nraw",
740         .fragment_length = 1,
741     };
742     close_reader.attach();
743     var close_handler = TestHandler{};
744     const status = try ClientStream.readForMethod(
745         .{ .headers = &headers, .head = &head },
746         &close_reader.reader,
747         "GET",
748         close_handler.handler(),
749     );
750     try std.testing.expectEqual(@as(u16, 200), status);
751     try std.testing.expectEqualStrings(
752         "raw",
753         close_handler.body[0..close_handler.body_length],
754     );
755 
756     var reject_reader = FragmentedReader{
757         .data = "HTTP/1.1 200 OK\r\nTransfer-Encoding: gzip\r\nContent-Length: 3\r\n\r\nraw",
758         .fragment_length = 1,
759     };
760     reject_reader.attach();
761     var reject_handler = TestHandler{};
762     try std.testing.expectError(
763         error.MalformedResponse,
764         ClientStream.readForMethod(
765             .{ .headers = &headers, .head = &head },
766             &reject_reader.reader,
767             "GET",
768             reject_handler.handler(),
769         ),
770     );
771     try std.testing.expectEqual(@as(usize, 0), reject_handler.head_count);
772     try std.testing.expectEqual(@as(usize, 0), reject_handler.body_length);
773 }
774 
775 test "Client stream rejects truncated framed bodies" {
776     var headers: [2]Header = undefined;
777     var head: [128]u8 = undefined;
778 
779     var fixed_reader = std.Io.Reader.fixed(
780         "HTTP/1.1 200 OK\r\nContent-Length: 5\r\n\r\nHell",
781     );
782     var fixed_handler = TestHandler{};
783     try std.testing.expectError(
784         error.MalformedResponse,
785         ClientStream.read(
786             .{ .headers = &headers, .head = &head },
787             &fixed_reader,
788             fixed_handler.handler(),
789         ),
790     );
791 
792     var chunk_reader = std.Io.Reader.fixed(
793         "HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n5\r\nHello\r\n",
794     );
795     var chunk_handler = TestHandler{};
796     try std.testing.expectError(
797         error.MalformedResponse,
798         ClientStream.read(
799             .{ .headers = &headers, .head = &head },
800             &chunk_reader,
801             chunk_handler.handler(),
802         ),
803     );
804 
805     var terminal_reader = std.Io.Reader.fixed(
806         "HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n0\r\n",
807     );
808     var terminal_handler = TestHandler{};
809     try std.testing.expectError(
810         error.MalformedResponse,
811         ClientStream.read(
812             .{ .headers = &headers, .head = &head },
813             &terminal_reader,
814             terminal_handler.handler(),
815         ),
816     );
817 }