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 }