An HTTP/1.1 server for zig
0

Configure Feed

Select the types of activity you want to include in your feed.

Remove dependency on deprecated std.net.Stream read and write methods (#159)

author
philipmv
committer
GitHub
date (Sep 29, 2025, 8:16 AM +0800) commit 86e86b2e parent c32853a4
+409 -147
+300 -96
src/httpz.zig
··· 699 699 700 700 var reply_buf: [512]u8 = undefined; 701 701 const reply = try websocket.Handshake.createReply(key, null, compression, &reply_buf); 702 - try http_conn.stream.writeAll(reply); 702 + var writer = http_conn.stream.writer(&.{}); 703 + const w = &writer.interface; 704 + try w.writeAll(reply); 705 + try w.flush(); 706 + 703 707 if (comptime std.meta.hasFn(H, "afterInit")) { 704 708 const params = @typeInfo(@TypeOf(H.afterInit)).@"fn".params; 705 709 try if (comptime params.len == 1) hc.handler.?.afterInit() else hc.handler.?.afterInit(ctx); ··· 945 949 test "httpz: invalid request" { 946 950 const stream = testStream(5992); 947 951 defer stream.close(); 948 - try stream.writeAll("TEA HTTP/1.1\r\n\r\n"); 952 + var writer = stream.writer(&.{}); 953 + const w = &writer.interface; 954 + try w.writeAll("TEA HTTP/1.1\r\n\r\n"); 955 + try w.flush(); 949 956 950 957 var buf: [100]u8 = undefined; 951 958 try t.expectString("HTTP/1.1 400 \r\nConnection: Close\r\nContent-Length: 15\r\n\r\nInvalid Request", testReadAll(stream, &buf)); ··· 954 961 test "httpz: invalid request path" { 955 962 const stream = testStream(5992); 956 963 defer stream.close(); 957 - try stream.writeAll("TEA /hello\rn\nWorld:test HTTP/1.1\r\n\r\n"); 964 + var writer = stream.writer(&.{}); 965 + const w = &writer.interface; 966 + try w.writeAll("TEA /hello\rn\nWorld:test HTTP/1.1\r\n\r\n"); 967 + try w.flush(); 958 968 959 969 var buf: [100]u8 = undefined; 960 970 try t.expectString("HTTP/1.1 400 \r\nConnection: Close\r\nContent-Length: 15\r\n\r\nInvalid Request", testReadAll(stream, &buf)); ··· 963 973 test "httpz: invalid header name" { 964 974 const stream = testStream(5992); 965 975 defer stream.close(); 966 - try stream.writeAll("GET / HTTP/1.1\r\nOver: 9000\r\nHel\tlo:World\r\n\r\n"); 976 + var writer = stream.writer(&.{}); 977 + const w = &writer.interface; 978 + try w.writeAll("GET / HTTP/1.1\r\nOver: 9000\r\nHel\tlo:World\r\n\r\n"); 979 + try w.flush(); 967 980 968 981 var buf: [100]u8 = undefined; 969 982 try t.expectString("HTTP/1.1 400 \r\nConnection: Close\r\nContent-Length: 15\r\n\r\nInvalid Request", testReadAll(stream, &buf)); ··· 972 985 test "httpz: invalid content length value (1)" { 973 986 const stream = testStream(5992); 974 987 defer stream.close(); 975 - try stream.writeAll("GET / HTTP/1.1\r\nContent-Length: HaHA\r\n\r\n"); 988 + var writer = stream.writer(&.{}); 989 + const w = &writer.interface; 990 + try w.writeAll("GET / HTTP/1.1\r\nContent-Length: HaHA\r\n\r\n"); 991 + try w.flush(); 976 992 977 993 var buf: [100]u8 = undefined; 978 994 try t.expectString("HTTP/1.1 400 \r\nConnection: Close\r\nContent-Length: 15\r\n\r\nInvalid Request", testReadAll(stream, &buf)); ··· 981 997 test "httpz: invalid content length value (2)" { 982 998 const stream = testStream(5992); 983 999 defer stream.close(); 984 - try stream.writeAll("GET / HTTP/1.1\r\nContent-Length: 1.0\r\n\r\n"); 1000 + var writer = stream.writer(&.{}); 1001 + const w = &writer.interface; 1002 + try w.writeAll("GET / HTTP/1.1\r\nContent-Length: 1.0\r\n\r\n"); 1003 + try w.flush(); 985 1004 986 1005 var buf: [100]u8 = undefined; 987 1006 try t.expectString("HTTP/1.1 400 \r\nConnection: Close\r\nContent-Length: 15\r\n\r\nInvalid Request", testReadAll(stream, &buf)); ··· 990 1009 test "httpz: body too big" { 991 1010 const stream = testStream(5993); 992 1011 defer stream.close(); 993 - try stream.writeAll("POST / HTTP/1.1\r\nContent-Length: 999999999999999999\r\n\r\n"); 1012 + var writer = stream.writer(&.{}); 1013 + const w = &writer.interface; 1014 + try w.writeAll("POST / HTTP/1.1\r\nContent-Length: 999999999999999999\r\n\r\n"); 1015 + try w.flush(); 994 1016 995 1017 var buf: [100]u8 = undefined; 996 1018 try t.expectString("HTTP/1.1 413 \r\nConnection: Close\r\nContent-Length: 23\r\n\r\nRequest body is too big", testReadAll(stream, &buf)); ··· 999 1021 test "httpz: overflow content length" { 1000 1022 const stream = testStream(5992); 1001 1023 defer stream.close(); 1002 - try stream.writeAll("GET / HTTP/1.1\r\nContent-Length: 999999999999999999999999999\r\n\r\n"); 1024 + var writer = stream.writer(&.{}); 1025 + const w = &writer.interface; 1026 + try w.writeAll("GET / HTTP/1.1\r\nContent-Length: 999999999999999999999999999\r\n\r\n"); 1027 + try w.flush(); 1003 1028 1004 1029 var buf: [100]u8 = undefined; 1005 1030 try t.expectString("HTTP/1.1 400 \r\nConnection: Close\r\nContent-Length: 15\r\n\r\nInvalid Request", testReadAll(stream, &buf)); ··· 1008 1033 test "httpz: no route" { 1009 1034 const stream = testStream(5992); 1010 1035 defer stream.close(); 1011 - try stream.writeAll("GET / HTTP/1.1\r\n\r\n"); 1036 + var writer = stream.writer(&.{}); 1037 + const w = &writer.interface; 1038 + try w.writeAll("GET / HTTP/1.1\r\n\r\n"); 1039 + try w.flush(); 1012 1040 1013 1041 var buf: [100]u8 = undefined; 1014 1042 try t.expectString("HTTP/1.1 404 \r\nContent-Length: 9\r\n\r\nNot Found", testReadAll(stream, &buf)); ··· 1017 1045 test "httpz: no route with custom notFound handler" { 1018 1046 const stream = testStream(5993); 1019 1047 defer stream.close(); 1020 - try stream.writeAll("GET /not_found HTTP/1.1\r\n\r\n"); 1048 + var writer = stream.writer(&.{}); 1049 + const w = &writer.interface; 1050 + try w.writeAll("GET /not_found HTTP/1.1\r\n\r\n"); 1051 + try w.flush(); 1021 1052 1022 1053 var buf: [100]u8 = undefined; 1023 1054 try t.expectString("HTTP/1.1 404 \r\nstate: 3\r\nContent-Length: 10\r\n\r\nwhere lah?", testReadAll(stream, &buf)); ··· 1029 1060 1030 1061 const stream = testStream(5992); 1031 1062 defer stream.close(); 1032 - try stream.writeAll("GET /fail HTTP/1.1\r\n\r\n"); 1063 + var writer = stream.writer(&.{}); 1064 + const w = &writer.interface; 1065 + try w.writeAll("GET /fail HTTP/1.1\r\n\r\n"); 1066 + try w.flush(); 1033 1067 1034 1068 var buf: [150]u8 = undefined; 1035 1069 try t.expectString("HTTP/1.1 500 \r\nContent-Length: 21\r\n\r\nInternal Server Error", testReadAll(stream, &buf)); ··· 1041 1075 1042 1076 const stream = testStream(5993); 1043 1077 defer stream.close(); 1044 - try stream.writeAll("GET /fail HTTP/1.1\r\n\r\n"); 1078 + var writer = stream.writer(&.{}); 1079 + const w = &writer.interface; 1080 + try w.writeAll("GET /fail HTTP/1.1\r\n\r\n"); 1081 + try w.flush(); 1045 1082 1046 1083 var buf: [150]u8 = undefined; 1047 1084 try t.expectString("HTTP/1.1 500 \r\nstate: 3\r\nerr: TestUnhandledError\r\nContent-Length: 29\r\n\r\n#/why/arent/tags/hierarchical", testReadAll(stream, &buf)); ··· 1052 1089 defer stream.close(); 1053 1090 1054 1091 { 1055 - try stream.writeAll("GET /test/method HTTP/1.1\r\n\r\n"); 1092 + var writer = stream.writer(&.{}); 1093 + const w = &writer.interface; 1094 + try w.writeAll("GET /test/method HTTP/1.1\r\n\r\n"); 1095 + try w.flush(); 1056 1096 var res = testReadParsed(stream); 1057 1097 defer res.deinit(); 1058 1098 try res.expectJson(.{ .method = "GET", .string = "" }); 1059 1099 } 1060 1100 1061 1101 { 1062 - try stream.writeAll("PUT /test/method HTTP/1.1\r\n\r\n"); 1102 + var writer = stream.writer(&.{}); 1103 + const w = &writer.interface; 1104 + try w.writeAll("PUT /test/method HTTP/1.1\r\n\r\n"); 1105 + try w.flush(); 1063 1106 var res = testReadParsed(stream); 1064 1107 defer res.deinit(); 1065 1108 try res.expectJson(.{ .method = "PUT", .string = "" }); 1066 1109 } 1067 1110 1068 1111 { 1069 - try stream.writeAll("TEA /test/method HTTP/1.1\r\n\r\n"); 1112 + var writer = stream.writer(&.{}); 1113 + const w = &writer.interface; 1114 + try w.writeAll("TEA /test/method HTTP/1.1\r\n\r\n"); 1115 + try w.flush(); 1070 1116 var res = testReadParsed(stream); 1071 1117 defer res.deinit(); 1072 1118 try res.expectJson(.{ .method = "OTHER", .string = "TEA" }); 1073 1119 } 1074 1120 1075 1121 { 1076 - try stream.writeAll("PING /test/method HTTP/1.1\r\n\r\n"); 1122 + var writer = stream.writer(&.{}); 1123 + const w = &writer.interface; 1124 + try w.writeAll("PING /test/method HTTP/1.1\r\n\r\n"); 1125 + try w.flush(); 1077 1126 var res = testReadParsed(stream); 1078 1127 defer res.deinit(); 1079 1128 try res.expectJson(.{ .method = "OTHER", .string = "PING" }); 1080 1129 } 1081 1130 1082 1131 { 1083 - try stream.writeAll("TEA /test/other HTTP/1.1\r\n\r\n"); 1132 + var writer = stream.writer(&.{}); 1133 + const w = &writer.interface; 1134 + try w.writeAll("TEA /test/other HTTP/1.1\r\n\r\n"); 1135 + try w.flush(); 1084 1136 var buf: [100]u8 = undefined; 1085 1137 try t.expectString("HTTP/1.1 404 \r\nContent-Length: 9\r\n\r\nNot Found", testReadAll(stream, &buf)); 1086 1138 } ··· 1089 1141 test "httpz: route params" { 1090 1142 const stream = testStream(5993); 1091 1143 defer stream.close(); 1092 - try stream.writeAll("GET /api/v2/users/9001 HTTP/1.1\r\n\r\n"); 1144 + var writer = stream.writer(&.{}); 1145 + const w = &writer.interface; 1146 + try w.writeAll("GET /api/v2/users/9001 HTTP/1.1\r\n\r\n"); 1147 + try w.flush(); 1093 1148 1094 1149 var buf: [100]u8 = undefined; 1095 1150 try t.expectString("HTTP/1.1 200 \r\nContent-Length: 20\r\n\r\nversion=v2,user=9001", testReadAll(stream, &buf)); ··· 1098 1153 test "httpz: request and response headers" { 1099 1154 const stream = testStream(5993); 1100 1155 defer stream.close(); 1101 - try stream.writeAll("GET /test/headers HTTP/1.1\r\nHeader-Name: Header-Value\r\n\r\n"); 1156 + var writer = stream.writer(&.{}); 1157 + const w = &writer.interface; 1158 + try w.writeAll("GET /test/headers HTTP/1.1\r\nHeader-Name: Header-Value\r\n\r\n"); 1159 + try w.flush(); 1102 1160 1103 1161 var buf: [100]u8 = undefined; 1104 1162 try t.expectString("HTTP/1.1 200 \r\nstate: 3\r\nEcho: Header-Value\r\nother: test-value\r\nContent-Length: 0\r\n\r\n", testReadAll(stream, &buf)); ··· 1107 1165 test "httpz: content-length body" { 1108 1166 const stream = testStream(5993); 1109 1167 defer stream.close(); 1110 - try stream.writeAll("GET /test/body/cl HTTP/1.1\r\nHeader-Name: Header-Value\r\nContent-Length: 4\r\n\r\nabcz"); 1168 + var writer = stream.writer(&.{}); 1169 + const w = &writer.interface; 1170 + try w.writeAll("GET /test/body/cl HTTP/1.1\r\nHeader-Name: Header-Value\r\nContent-Length: 4\r\n\r\nabcz"); 1171 + try w.flush(); 1111 1172 1112 1173 var buf: [100]u8 = undefined; 1113 1174 try t.expectString("HTTP/1.1 200 \r\nEcho-Body: abcz\r\nContent-Length: 0\r\n\r\n", testReadAll(stream, &buf)); ··· 1116 1177 test "httpz: json response" { 1117 1178 const stream = testStream(5992); 1118 1179 defer stream.close(); 1119 - try stream.writeAll("GET /test/json HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1180 + var writer = stream.writer(&.{}); 1181 + const w = &writer.interface; 1182 + try w.writeAll("GET /test/json HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1183 + try w.flush(); 1120 1184 1121 1185 var buf: [200]u8 = undefined; 1122 1186 try t.expectString("HTTP/1.1 201 \r\nContent-Type: application/json; charset=UTF-8\r\nContent-Length: 26\r\n\r\n{\"over\":9000,\"teg\":\"soup\"}", testReadAll(stream, &buf)); ··· 1125 1189 test "httpz: query" { 1126 1190 const stream = testStream(5992); 1127 1191 defer stream.close(); 1128 - try stream.writeAll("GET /test/query?fav=keemun%20te%61%21 HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1192 + var writer = stream.writer(&.{}); 1193 + const w = &writer.interface; 1194 + try w.writeAll("GET /test/query?fav=keemun%20te%61%21 HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1195 + try w.flush(); 1129 1196 1130 1197 var buf: [200]u8 = undefined; 1131 1198 try t.expectString("HTTP/1.1 200 \r\nContent-Length: 11\r\n\r\nkeemun tea!", testReadAll(stream, &buf)); ··· 1134 1201 test "httpz: chunked" { 1135 1202 const stream = testStream(5992); 1136 1203 defer stream.close(); 1137 - try stream.writeAll("GET /test/chunked HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1204 + var writer = stream.writer(&.{}); 1205 + const w = &writer.interface; 1206 + try w.writeAll("GET /test/chunked HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1207 + try w.flush(); 1138 1208 1139 1209 var buf: [1000]u8 = undefined; 1140 1210 try t.expectString("HTTP/1.1 200 \r\nOver: 9000!\r\nTransfer-Encoding: chunked\r\n\r\n7\r\nChunk 1\r\n11\r\nand another chunk\r\n0\r\n\r\n", testReadAll(stream, &buf)); ··· 1143 1213 test "httpz: route-specific dispatcher" { 1144 1214 const stream = testStream(5992); 1145 1215 defer stream.close(); 1146 - try stream.writeAll("HEAD /test/dispatcher HTTP/1.1\r\n\r\n"); 1216 + var writer = stream.writer(&.{}); 1217 + const w = &writer.interface; 1218 + try w.writeAll("HEAD /test/dispatcher HTTP/1.1\r\n\r\n"); 1219 + try w.flush(); 1147 1220 1148 1221 var buf: [200]u8 = undefined; 1149 1222 try t.expectString("HTTP/1.1 200 \r\ndispatcher: test-dispatcher-1\r\nContent-Length: 6\r\n\r\naction", testReadAll(stream, &buf)); ··· 1152 1225 test "httpz: middlewares" { 1153 1226 const stream = testStream(5992); 1154 1227 defer stream.close(); 1228 + var writer = stream.writer(&.{}); 1229 + const w = &writer.interface; 1155 1230 1156 1231 { 1157 - try stream.writeAll("GET /test/middlewares HTTP/1.1\r\n\r\n"); 1232 + try w.writeAll("GET /test/middlewares HTTP/1.1\r\n\r\n"); 1233 + try w.flush(); 1158 1234 var res = testReadParsed(stream); 1159 1235 defer res.deinit(); 1160 1236 ··· 1168 1244 defer stream.close(); 1169 1245 1170 1246 { 1171 - try stream.writeAll("GET /echo HTTP/1.1\r\n\r\n"); 1247 + var writer = stream.writer(&.{}); 1248 + const w = &writer.interface; 1249 + try w.writeAll("GET /echo HTTP/1.1\r\n\r\n"); 1250 + try w.flush(); 1172 1251 var res = testReadParsed(stream); 1173 1252 defer res.deinit(); 1174 1253 try t.expectEqual(null, res.headers.get("Access-Control-Max-Age")); ··· 1179 1258 1180 1259 { 1181 1260 // cors endpoint but not cors options 1182 - try stream.writeAll("OPTIONS /test/cors HTTP/1.1\r\nSec-Fetch-Mode: navigate\r\n\r\n"); 1261 + var writer = stream.writer(&.{}); 1262 + const w = &writer.interface; 1263 + try w.writeAll("OPTIONS /test/cors HTTP/1.1\r\nSec-Fetch-Mode: navigate\r\n\r\n"); 1264 + try w.flush(); 1183 1265 var res = testReadParsed(stream); 1184 1266 defer res.deinit(); 1185 1267 ··· 1191 1273 1192 1274 { 1193 1275 // cors request 1194 - try stream.writeAll("OPTIONS /test/cors HTTP/1.1\r\nSec-Fetch-Mode: cors\r\n\r\n"); 1276 + var writer = stream.writer(&.{}); 1277 + const w = &writer.interface; 1278 + try w.writeAll("OPTIONS /test/cors HTTP/1.1\r\nSec-Fetch-Mode: cors\r\n\r\n"); 1279 + try w.flush(); 1195 1280 var res = testReadParsed(stream); 1196 1281 defer res.deinit(); 1197 1282 ··· 1203 1288 1204 1289 { 1205 1290 // cors request, non-options 1206 - try stream.writeAll("GET /test/cors HTTP/1.1\r\nSec-Fetch-Mode: cors\r\n\r\n"); 1291 + var writer = stream.writer(&.{}); 1292 + const w = &writer.interface; 1293 + try w.writeAll("GET /test/cors HTTP/1.1\r\nSec-Fetch-Mode: cors\r\n\r\n"); 1294 + try w.flush(); 1207 1295 var res = testReadParsed(stream); 1208 1296 defer res.deinit(); 1209 1297 ··· 1219 1307 defer stream.close(); 1220 1308 1221 1309 { 1222 - try stream.writeAll("GET / HTTP/1.1\r\n\r\n"); 1310 + var writer = stream.writer(&.{}); 1311 + const w = &writer.interface; 1312 + try w.writeAll("GET / HTTP/1.1\r\n\r\n"); 1313 + try w.flush(); 1223 1314 var res = testReadParsed(stream); 1224 1315 defer res.deinit(); 1225 1316 ··· 1228 1319 } 1229 1320 1230 1321 { 1231 - try stream.writeAll("GET /admin/users HTTP/1.1\r\n\r\n"); 1322 + var writer = stream.writer(&.{}); 1323 + const w = &writer.interface; 1324 + try w.writeAll("GET /admin/users HTTP/1.1\r\n\r\n"); 1325 + try w.flush(); 1232 1326 var res = testReadParsed(stream); 1233 1327 defer res.deinit(); 1234 1328 ··· 1237 1331 } 1238 1332 1239 1333 { 1240 - try stream.writeAll("PUT /admin/users/:id HTTP/1.1\r\n\r\n"); 1334 + var writer = stream.writer(&.{}); 1335 + const w = &writer.interface; 1336 + try w.writeAll("PUT /admin/users/:id HTTP/1.1\r\n\r\n"); 1337 + try w.flush(); 1241 1338 var res = testReadParsed(stream); 1242 1339 defer res.deinit(); 1243 1340 ··· 1246 1343 } 1247 1344 1248 1345 { 1249 - try stream.writeAll("HEAD /debug/ping HTTP/1.1\r\n\r\n"); 1346 + var writer = stream.writer(&.{}); 1347 + const w = &writer.interface; 1348 + try w.writeAll("HEAD /debug/ping HTTP/1.1\r\n\r\n"); 1349 + try w.flush(); 1250 1350 var res = testReadParsed(stream); 1251 1351 defer res.deinit(); 1252 1352 ··· 1255 1355 } 1256 1356 1257 1357 { 1258 - try stream.writeAll("OPTIONS /debug/stats HTTP/1.1\r\n\r\n"); 1358 + var writer = stream.writer(&.{}); 1359 + const w = &writer.interface; 1360 + try w.writeAll("OPTIONS /debug/stats HTTP/1.1\r\n\r\n"); 1361 + try w.flush(); 1259 1362 var res = testReadParsed(stream); 1260 1363 defer res.deinit(); 1261 1364 ··· 1264 1367 } 1265 1368 1266 1369 { 1267 - try stream.writeAll("POST /login HTTP/1.1\r\n\r\n"); 1370 + var writer = stream.writer(&.{}); 1371 + const w = &writer.interface; 1372 + try w.writeAll("POST /login HTTP/1.1\r\n\r\n"); 1373 + try w.flush(); 1268 1374 var res = testReadParsed(stream); 1269 1375 defer res.deinit(); 1270 1376 ··· 1276 1382 test "httpz: event stream" { 1277 1383 const stream = testStream(5992); 1278 1384 defer stream.close(); 1279 - try stream.writeAll("GET /test/stream HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1385 + var writer = stream.writer(&.{}); 1386 + const w = &writer.interface; 1387 + try w.writeAll("GET /test/stream HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1388 + try w.flush(); 1280 1389 1281 1390 var res = testReadParsed(stream); 1282 1391 defer res.deinit(); ··· 1292 1401 test "httpz: event stream sync" { 1293 1402 const stream = testStream(5992); 1294 1403 defer stream.close(); 1295 - try stream.writeAll("GET /test/streamsync HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1404 + var writer = stream.writer(&.{}); 1405 + const w = &writer.interface; 1406 + try w.writeAll("GET /test/streamsync HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1407 + try w.flush(); 1296 1408 1297 1409 var res = testReadParsed(stream); 1298 1410 defer res.deinit(); ··· 1308 1420 test "httpz: keepalive" { 1309 1421 const stream = testStream(5993); 1310 1422 defer stream.close(); 1311 - try stream.writeAll("GET /api/v2/users/9001 HTTP/1.1\r\n\r\n"); 1423 + var writer = stream.writer(&.{}); 1424 + const w = &writer.interface; 1425 + try w.writeAll("GET /api/v2/users/9001 HTTP/1.1\r\n\r\n"); 1426 + try w.flush(); 1312 1427 1313 1428 var buf: [100]u8 = undefined; 1314 1429 try t.expectString("HTTP/1.1 200 \r\nContent-Length: 20\r\n\r\nversion=v2,user=9001", testReadAll(stream, &buf)); 1315 1430 1316 - try stream.writeAll("GET /api/v2/users/123 HTTP/1.1\r\n\r\n"); 1431 + try w.writeAll("GET /api/v2/users/123 HTTP/1.1\r\n\r\n"); 1432 + try w.flush(); 1317 1433 try t.expectString("HTTP/1.1 200 \r\nContent-Length: 19\r\n\r\nversion=v2,user=123", testReadAll(stream, &buf)); 1318 1434 } 1319 1435 1320 1436 test "httpz: route data" { 1321 1437 const stream = testStream(5992); 1322 1438 defer stream.close(); 1323 - try stream.writeAll("GET /test/route_data HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1439 + var writer = stream.writer(&.{}); 1440 + const w = &writer.interface; 1441 + try w.writeAll("GET /test/route_data HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1442 + try w.flush(); 1324 1443 1325 1444 var res = testReadParsed(stream); 1326 1445 defer res.deinit(); ··· 1330 1449 test "httpz: keepalive with explicit write" { 1331 1450 const stream = testStream(5993); 1332 1451 defer stream.close(); 1333 - try stream.writeAll("GET /write/9001 HTTP/1.1\r\n\r\n"); 1452 + var writer = stream.writer(&.{}); 1453 + const w = &writer.interface; 1454 + try w.writeAll("GET /write/9001 HTTP/1.1\r\n\r\n"); 1455 + try w.flush(); 1334 1456 1335 1457 var buf: [1000]u8 = undefined; 1336 1458 try t.expectString("HTTP/1.1 200 \r\nContent-Length: 47\r\n\r\n{\"state\":3,\"method\":\"GET\",\"path\":\"/write/9001\"}", testReadAll(stream, &buf)); 1337 1459 1338 - try stream.writeAll("GET /write/123 HTTP/1.1\r\n\r\n"); 1460 + try w.writeAll("GET /write/123 HTTP/1.1\r\n\r\n"); 1461 + try w.flush(); 1339 1462 try t.expectString("HTTP/1.1 200 \r\nContent-Length: 46\r\n\r\n{\"state\":3,\"method\":\"GET\",\"path\":\"/write/123\"}", testReadAll(stream, &buf)); 1340 1463 } 1341 1464 1342 1465 test "httpz: request in chunks" { 1343 1466 const stream = testStream(5993); 1344 1467 defer stream.close(); 1345 - try stream.writeAll("GET /api/v2/use"); 1468 + var writer = stream.writer(&.{}); 1469 + const w = &writer.interface; 1470 + try w.writeAll("GET /api/v2/use"); 1471 + try w.flush(); 1346 1472 std.Thread.sleep(std.time.ns_per_ms * 10); 1347 - try stream.writeAll("rs/11 HTTP/1.1\r\n\r\n"); 1473 + try w.writeAll("rs/11 HTTP/1.1\r\n\r\n"); 1474 + try w.flush(); 1348 1475 1349 1476 var buf: [100]u8 = undefined; 1350 1477 try t.expectString("HTTP/1.1 200 \r\nContent-Length: 18\r\n\r\nversion=v2,user=11", testReadAll(stream, &buf)); ··· 1355 1482 1356 1483 const stream = testStream(5996); 1357 1484 defer stream.close(); 1485 + var writer = stream.writer(&.{}); 1486 + const w = &writer.interface; 1358 1487 1359 1488 var expected: [10]TestUser = undefined; 1360 1489 ··· 1364 1493 .id = try std.fmt.allocPrint(t.arena.allocator(), "id-{d}", .{i}), 1365 1494 .power = i, 1366 1495 }; 1367 - try stream.writeAll(try std.fmt.bufPrint(&buf, "GET /test/writer?count={d} HTTP/1.1\r\nContent-Length: 0\r\n\r\n", .{i + 1})); 1496 + try w.writeAll(try std.fmt.bufPrint(&buf, "GET /test/writer?count={d} HTTP/1.1\r\nContent-Length: 0\r\n\r\n", .{i + 1})); 1497 + try w.flush(); 1368 1498 1369 1499 var res = testReadParsed(stream); 1370 1500 defer res.deinit(); ··· 1376 1506 test "httpz: custom dispatch without action context" { 1377 1507 const stream = testStream(5994); 1378 1508 defer stream.close(); 1379 - try stream.writeAll("GET / HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1509 + var writer = stream.writer(&.{}); 1510 + const w = &writer.interface; 1511 + try w.writeAll("GET / HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1512 + try w.flush(); 1380 1513 1381 1514 var buf: [200]u8 = undefined; 1382 1515 try t.expectString("HTTP/1.1 200 \r\nContent-Type: application/json; charset=UTF-8\r\ndstate: 10\r\ndispatch: TestHandlerDispatch\r\nContent-Length: 12\r\n\r\n{\"state\":10}", testReadAll(stream, &buf)); ··· 1385 1518 test "httpz: custom dispatch with action context" { 1386 1519 const stream = testStream(5995); 1387 1520 defer stream.close(); 1388 - try stream.writeAll("GET /?name=teg HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1521 + var writer = stream.writer(&.{}); 1522 + const w = &writer.interface; 1523 + try w.writeAll("GET /?name=teg HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1524 + try w.flush(); 1389 1525 1390 1526 var buf: [200]u8 = undefined; 1391 1527 try t.expectString("HTTP/1.1 200 \r\nContent-Type: application/json; charset=UTF-8\r\ndstate: 20\r\ndispatch: TestHandlerDispatchContext\r\nContent-Length: 12\r\n\r\n{\"other\":30}", testReadAll(stream, &buf)); ··· 1394 1530 test "httpz: custom handle" { 1395 1531 const stream = testStream(5997); 1396 1532 defer stream.close(); 1397 - try stream.writeAll("GET /whatever?name=teg HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1533 + var writer = stream.writer(&.{}); 1534 + const w = &writer.interface; 1535 + try w.writeAll("GET /whatever?name=teg HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1536 + try w.flush(); 1398 1537 1399 1538 var buf: [100]u8 = undefined; 1400 1539 try t.expectString("HTTP/1.1 200 \r\nContent-Length: 9\r\n\r\nhello teg", testReadAll(stream, &buf)); ··· 1405 1544 // no body 1406 1545 const stream = testStream(5992); 1407 1546 defer stream.close(); 1408 - try stream.writeAll("GET /test/req_reader HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1547 + 1548 + var writer = stream.writer(&.{}); 1549 + const w = &writer.interface; 1550 + 1551 + try w.writeAll("GET /test/req_reader HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1552 + try w.flush(); 1409 1553 1410 1554 var res = testReadParsed(stream); 1411 1555 defer res.deinit(); ··· 1416 1560 // small body 1417 1561 const stream = testStream(5992); 1418 1562 defer stream.close(); 1419 - try stream.writeAll("GET /test/req_reader HTTP/1.1\r\nContent-Length: 4\r\n\r\n123z"); 1563 + 1564 + var writer = stream.writer(&.{}); 1565 + const w = &writer.interface; 1566 + 1567 + try w.writeAll("GET /test/req_reader HTTP/1.1\r\nContent-Length: 4\r\n\r\n123z"); 1568 + try w.flush(); 1420 1569 1421 1570 var res = testReadParsed(stream); 1422 1571 defer res.deinit(); ··· 1428 1577 // medium body 1429 1578 const stream = testStream(5992); 1430 1579 defer stream.close(); 1431 - try stream.writeAll(std.fmt.comptimePrint("GET /test/req_reader HTTP/1.1\r\nContent-Length: {d}\r\n\r\n" ++ ("a" ** length), .{length})); 1580 + 1581 + var writer = stream.writer(&.{}); 1582 + const w = &writer.interface; 1583 + 1584 + try w.writeAll(std.fmt.comptimePrint("GET /test/req_reader HTTP/1.1\r\nContent-Length: {d}\r\n\r\n" ++ ("a" ** length), .{length})); 1585 + try w.flush(); 1432 1586 1433 1587 var res = testReadParsed(stream); 1434 1588 defer res.deinit(); ··· 1443 1597 for (0..10) |_| { 1444 1598 const stream = testStream(5992); 1445 1599 defer stream.close(); 1600 + 1601 + var buf: [1024]u8 = undefined; 1602 + var writer = stream.writer(&buf); 1603 + const w = &writer.interface; 1604 + 1446 1605 var req: []const u8 = std.fmt.comptimePrint("GET /test/req_reader HTTP/1.1\r\nContent-Length: {d}\r\n\r\n" ++ ("a" ** length), .{length}); 1447 1606 while (req.len > 0) { 1448 1607 const len = random.uintAtMost(usize, req.len - 1) + 1; 1449 - const n = stream.write(req[0..len]) catch |err| switch (err) { 1450 - error.WouldBlock => 0, 1451 - else => return err, 1452 - }; 1608 + try w.writeAll(req[0..len]); 1453 1609 std.Thread.sleep(std.time.ns_per_ms * 2); 1454 - req = req[n..]; 1610 + req = req[len..]; 1455 1611 } 1612 + 1613 + try w.flush(); 1456 1614 1457 1615 var res = testReadParsed(stream); 1458 1616 defer res.deinit(); ··· 1463 1621 test "websocket: invalid request" { 1464 1622 const stream = testStream(5998); 1465 1623 defer stream.close(); 1466 - try stream.writeAll("GET /ws HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1624 + var writer = stream.writer(&.{}); 1625 + const w = &writer.interface; 1626 + try w.writeAll("GET /ws HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1627 + try w.flush(); 1467 1628 1468 1629 var res = testReadParsed(stream); 1469 1630 defer res.deinit(); ··· 1473 1634 test "websocket: upgrade" { 1474 1635 const stream = testStream(5998); 1475 1636 defer stream.close(); 1476 - try stream.writeAll("GET /ws HTTP/1.1\r\nContent-Length: 0\r\n"); 1477 - try stream.writeAll("upgrade: WEBsocket\r\n"); 1478 - try stream.writeAll("Sec-Websocket-verSIon: 13\r\n"); 1479 - try stream.writeAll("ConnectioN: abc,upgrade,123\r\n"); 1480 - try stream.writeAll("SEC-WEBSOCKET-KeY: a-secret-key\r\n\r\n"); 1637 + var writer = stream.writer(&.{}); 1638 + const w = &writer.interface; 1639 + try w.writeAll("GET /ws HTTP/1.1\r\nContent-Length: 0\r\n"); 1640 + try w.writeAll("upgrade: WEBsocket\r\n"); 1641 + try w.writeAll("Sec-Websocket-verSIon: 13\r\n"); 1642 + try w.writeAll("ConnectioN: abc,upgrade,123\r\n"); 1643 + try w.writeAll("SEC-WEBSOCKET-KeY: a-secret-key\r\n\r\n"); 1644 + try w.flush(); 1481 1645 1482 1646 var res = testReadHeader(stream); 1483 1647 defer res.deinit(); ··· 1486 1650 try t.expectString("upgrade", res.headers.get("Connection").?); 1487 1651 try t.expectString("55eM2SNGu+68v5XXrr982mhPFkU=", res.headers.get("Sec-Websocket-Accept").?); 1488 1652 1489 - try stream.writeAll(&websocket.frameText("over 9000!")); 1490 - try stream.writeAll(&websocket.frameText("close")); 1653 + try w.writeAll(&websocket.frameText("over 9000!")); 1654 + try w.writeAll(&websocket.frameText("close")); 1655 + try w.flush(); 1491 1656 1492 1657 var pos: usize = 0; 1493 1658 var buf: [100]u8 = undefined; 1494 1659 var wait_count: usize = 0; 1660 + var reader = stream.reader(&.{}); 1661 + const r = reader.interface(); 1495 1662 while (pos < 16) { 1496 - const n = stream.read(buf[pos..]) catch |err| switch (err) { 1497 - error.WouldBlock => { 1498 - if (wait_count == 100) { 1499 - break; 1500 - } 1501 - wait_count += 1; 1502 - std.Thread.sleep(std.time.ns_per_ms); 1503 - continue; 1504 - }, 1505 - else => return err, 1506 - }; 1663 + const n = r.readSliceShort(buf[pos..]) catch |err| 1664 + switch (err) { 1665 + error.ReadFailed => { 1666 + if (reader.getError()) |e| { 1667 + switch (e) { 1668 + error.WouldBlock => { 1669 + if (wait_count == 100) { 1670 + break; 1671 + } 1672 + wait_count += 1; 1673 + std.Thread.sleep(std.time.ns_per_ms); 1674 + continue; 1675 + }, 1676 + else => {}, 1677 + } 1678 + } 1679 + return err; 1680 + }, 1681 + }; 1682 + 1507 1683 if (n == 0) { 1508 1684 break; 1509 1685 } ··· 1553 1729 fn testReadAll(stream: std.net.Stream, buf: []u8) []u8 { 1554 1730 var pos: usize = 0; 1555 1731 var blocked = false; 1732 + var reader = stream.reader(&.{}); 1733 + const r = reader.interface(); 1556 1734 while (true) { 1557 1735 std.debug.assert(pos < buf.len); 1558 - const n = stream.read(buf[pos..]) catch |err| switch (err) { 1559 - error.WouldBlock => { 1560 - if (blocked) return buf[0..pos]; 1561 - blocked = true; 1562 - std.Thread.sleep(std.time.ns_per_ms); 1563 - continue; 1564 - }, 1565 - error.ConnectionResetByPeer => return buf[0..pos], 1566 - else => @panic(@errorName(err)), 1567 - }; 1736 + var vecs: [1][]u8 = .{buf[pos..]}; 1737 + const n = r.readVec(&vecs) catch |err| 1738 + switch (err) { 1739 + error.ReadFailed => { 1740 + if (reader.getError()) |e| { 1741 + switch (e) { 1742 + error.WouldBlock => { 1743 + if (blocked) return buf[0..pos]; 1744 + blocked = true; 1745 + std.Thread.sleep(std.time.ns_per_ms); 1746 + continue; 1747 + }, 1748 + error.ConnectionResetByPeer => return buf[0..pos], 1749 + else => @panic(@errorName(e)), 1750 + } 1751 + } 1752 + @panic(@errorName(err)); 1753 + }, 1754 + error.EndOfStream => 0, 1755 + }; 1756 + 1568 1757 if (n == 0) { 1569 1758 return buf[0..pos]; 1570 1759 } ··· 1584 1773 var pos: usize = 0; 1585 1774 var blocked = false; 1586 1775 var buf: [1024]u8 = undefined; 1776 + var reader = stream.reader(&.{}); 1777 + const r = reader.interface(); 1587 1778 while (true) { 1588 1779 std.debug.assert(pos < buf.len); 1589 - const n = stream.read(buf[pos..]) catch |err| switch (err) { 1590 - error.WouldBlock => { 1591 - if (blocked) unreachable; 1592 - blocked = true; 1593 - std.Thread.sleep(std.time.ns_per_ms); 1594 - continue; 1595 - }, 1596 - else => @panic(@errorName(err)), 1597 - }; 1780 + var vecs: [1][]u8 = .{buf[pos..]}; 1781 + const n = r.readVec(&vecs) catch |err| 1782 + switch (err) { 1783 + error.ReadFailed => { 1784 + if (reader.getError()) |e| { 1785 + switch (e) { 1786 + error.WouldBlock => { 1787 + if (blocked) unreachable; 1788 + blocked = true; 1789 + std.Thread.sleep(std.time.ns_per_ms); 1790 + continue; 1791 + }, 1792 + else => @panic(@errorName(e)), 1793 + } 1794 + } 1795 + @panic(@errorName(err)); 1796 + }, 1797 + error.EndOfStream => 0, 1798 + }; 1598 1799 1599 1800 if (n == 0) unreachable; 1600 1801 ··· 1685 1886 data: []const u8, 1686 1887 1687 1888 fn handle(self: StreamContext, stream: std.net.Stream) void { 1688 - stream.writeAll(self.data) catch unreachable; 1689 - stream.writeAll("a message") catch unreachable; 1889 + var writer = stream.writer(&.{}); 1890 + const w = &writer.interface; 1891 + w.writeAll(self.data) catch unreachable; 1892 + w.writeAll("a message") catch unreachable; 1893 + w.flush() catch unreachable; 1690 1894 } 1691 1895 }; 1692 1896
+15 -12
src/request.zig
··· 705 705 } 706 706 707 707 // returns true if the header has been fully parsed 708 - pub fn parse(self: *State, req_arena: Allocator, stream: anytype) !bool { 708 + pub fn parse(self: *State, req_arena: Allocator, stream: *std.Io.Reader) !bool { 709 709 if (self.body != null) { 710 710 // if we have a body, then we've read the header. We want to read into 711 711 // self.body, not self.buf. ··· 714 714 715 715 var len = self.len; 716 716 const buf = self.buf; 717 - const n = try stream.read(buf[len..]); 717 + var vecs: [1][]u8 = .{buf[len..]}; 718 + const n = try stream.readVec(&vecs); 718 719 if (n == 0) { 719 - return error.ConnectionClosed; 720 + return false; 720 721 } 721 722 len = len + n; 722 723 self.len = len; ··· 1043 1044 return false; 1044 1045 } 1045 1046 1046 - fn readBody(self: *State, stream: anytype) !bool { 1047 + fn readBody(self: *State, stream: *std.Io.Reader) !bool { 1047 1048 const buf = self.body.?.data; 1048 1049 1049 - const n = try stream.read(buf[self.body_pos..]); 1050 - if (n == 0) { 1051 - return error.ConnectionClosed; 1052 - } 1053 - self.body_pos += n; 1050 + var vecs: [1][]u8 = .{buf[self.body_pos..]}; 1051 + self.body_pos += try stream.readVec(&vecs); 1054 1052 return (self.body_pos == self.body_len); 1055 1053 } 1056 1054 }; ··· 1706 1704 1707 1705 var conn = ctx.conn; 1708 1706 var fake_reader = ctx.fakeReader(); 1707 + const fr = &fake_reader.interface; 1709 1708 while (true) { 1710 - const done = try conn.req_state.parse(conn.req_arena.allocator(), &fake_reader); 1709 + const done = try conn.req_state.parse(conn.req_arena.allocator(), fr); 1711 1710 if (done) break; 1712 1711 } 1713 1712 ··· 1782 1781 fn testParse(input: []const u8, config: Config) !Request { 1783 1782 var ctx = t.Context.allocInit(t.arena.allocator(), .{ .request = config }); 1784 1783 ctx.write(input); 1784 + var reader = ctx.stream.reader(&.{}); 1785 + const r = reader.interface(); 1785 1786 while (true) { 1786 - const done = try ctx.conn.req_state.parse(ctx.conn.req_arena.allocator(), ctx.stream); 1787 + const done = try ctx.conn.req_state.parse(ctx.conn.req_arena.allocator(), r); 1787 1788 if (done) break; 1788 1789 } 1789 1790 return Request.init(ctx.conn.req_arena.allocator(), ctx.conn); ··· 1793 1794 var ctx = t.Context.init(.{ .request = config }); 1794 1795 defer ctx.deinit(); 1795 1796 1797 + var reader = ctx.stream.reader(&.{}); 1798 + const r = reader.interface(); 1796 1799 ctx.write(input); 1797 - try t.expectError(expected, ctx.conn.req_state.parse(ctx.conn.req_arena.allocator(), ctx.stream)); 1800 + try t.expectError(expected, ctx.conn.req_state.parse(ctx.conn.req_arena.allocator(), r)); 1798 1801 } 1799 1802 1800 1803 fn randomMethod(random: std.Random) []const u8 {
+56 -11
src/t.zig
··· 180 180 if (self.fake) { 181 181 self.to_read.appendSlice(self.arena.allocator(), data) catch unreachable; 182 182 } else { 183 - self.client.writeAll(data) catch unreachable; 183 + var buf: [1024]u8 = undefined; 184 + var writer = self.client.writer(&buf); 185 + const w = &writer.interface; 186 + w.writeAll(data) catch unreachable; 187 + w.flush() catch unreachable; 184 188 } 185 189 } 186 190 ··· 188 192 var buf: [1024]u8 = undefined; 189 193 var arr: std.ArrayList(u8) = .empty; 190 194 195 + var reader = self.client.reader(&.{}); 196 + const r = reader.interface(); 191 197 while (true) { 192 - const n = self.client.read(&buf) catch |err| switch (err) { 193 - error.WouldBlock => return arr, 194 - else => return err, 195 - }; 198 + const n = r.readSliceShort(&buf) catch |err| 199 + switch (err) { 200 + error.ReadFailed => { 201 + if (reader.getError()) |e| { 202 + switch (e) { 203 + error.WouldBlock => return arr, 204 + else => return e, 205 + } 206 + } 207 + return err; 208 + }, 209 + }; 196 210 if (n == 0) return arr; 197 211 try arr.appendSlice(a, buf[0..n]); 198 212 } ··· 203 217 var pos: usize = 0; 204 218 var buf = try allocator.alloc(u8, expected.len); 205 219 defer allocator.free(buf); 220 + var reader = self.client.reader(&.{}); 221 + const r = reader.interface(); 206 222 while (pos < buf.len) { 207 - const n = try self.client.read(buf[pos..]); 223 + const n = try r.readSliceShort(buf[pos..]); 208 224 if (n == 0) break; 209 225 pos += n; 210 226 } ··· 218 234 .usec = 1_000, 219 235 })) catch unreachable; 220 236 221 - const n: usize = self.client.read(buf[0..]) catch |err| blk: { 222 - switch (err) { 223 - error.WouldBlock => break :blk 0, 224 - else => @panic(@errorName(err)), 225 - } 237 + const n = r.readSliceShort(buf[0..]) catch |err| blk: switch (err) { 238 + error.ReadFailed => { 239 + if (reader.getError()) |e| { 240 + switch (e) { 241 + error.WouldBlock => break :blk 0, 242 + else => @panic(@errorName(e)), 243 + } 244 + } 245 + @panic(@errorName(err)); 246 + }, 226 247 }; 227 248 try expectEqual(0, n); 228 249 ··· 277 298 pos: usize, 278 299 buf: []const u8, 279 300 random: std.Random, 301 + interface: std.Io.Reader = 302 + .{ 303 + .vtable = &.{ 304 + .stream = FakeReader.stream, 305 + }, 306 + .buffer = &.{}, 307 + .seek = 0, 308 + .end = 0, 309 + }, 280 310 281 311 pub fn read( 282 312 self: *FakeReader, ··· 292 322 const to_read = self.random.intRangeAtMost(usize, 1, @min(data.len, buf.len)); 293 323 @memcpy(buf[0..to_read], data[0..to_read]); 294 324 self.pos += to_read; 325 + return to_read; 326 + } 327 + 328 + fn stream(io_reader: *std.Io.Reader, w: *std.Io.Writer, limit: std.Io.Limit) std.Io.Reader.StreamError!usize { 329 + const r: *FakeReader = @alignCast(@fieldParentPtr("interface", io_reader)); 330 + const data = r.buf[r.pos..]; 331 + 332 + if (data.len == 0 or limit == .nothing) { 333 + return 0; 334 + } 335 + 336 + // randomly fragment the data 337 + const to_read = r.random.intRangeAtMost(usize, 1, @min(data.len, limit.toInt() orelse data.len)); 338 + try w.writeAll(data[0..to_read]); 339 + r.pos += to_read; 295 340 return to_read; 296 341 } 297 342 };
+1 -1
src/testing.zig
··· 16 16 // Parse a basic request. This will put our conn.req_state into a valid state 17 17 // for creating a request. Application code can modify the request directly 18 18 // thereafter to change whatever properties they want. 19 - var base_request = std.io.fixedBufferStream("GET / HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 19 + var base_request: std.Io.Reader = .fixed("GET / HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 20 20 while (true) { 21 21 const done = conn.req_state.parse(conn.req_arena.allocator(), &base_request) catch unreachable; 22 22 if (done) {
+37 -27
src/worker.zig
··· 261 261 } 262 262 263 263 var is_first = true; 264 + var reader = stream.reader(&.{}); // Request.State does its own buffering 264 265 while (true) { 265 - const done = conn.req_state.parse(conn.req_arena.allocator(), stream) catch |err| switch (err) { 266 - error.WouldBlock => { 267 - if (is_keepalive and is_first) { 268 - metrics.timeoutKeepalive(1); 269 - } else { 270 - metrics.timeoutRequest(1); 271 - } 272 - return .close; 273 - }, 274 - error.NotOpenForReading => { 275 - // This can only happen when we're shutting down and our 276 - // listener has called posix.close(socket) to unblock 277 - // this thread. Using `.disown` is a bit of a hack, but 278 - // disown is handled in handleConnection the way we want 279 - // WE DO NOT WANT to return .close, else that would result 280 - // in posix.close(socket) being called on an already-closed 281 - // socket, which would panic. 282 - return .disown; 283 - }, 284 - else => { 285 - requestError(conn, err) catch {}; 286 - posix.close(stream.handle); 287 - return .disown; 288 - }, 266 + const done = conn.req_state.parse(conn.req_arena.allocator(), reader.interface()) catch |err| { 267 + switch (err) { 268 + error.ReadFailed => { 269 + if (reader.getError()) |e| { 270 + switch (e) { 271 + error.WouldBlock => { 272 + if (is_keepalive and is_first) { 273 + metrics.timeoutKeepalive(1); 274 + } else { 275 + metrics.timeoutRequest(1); 276 + } 277 + return .close; 278 + }, 279 + error.NotOpenForReading => { 280 + // This can only happen when we're shutting down and our 281 + // listener has called posix.close(socket) to unblock 282 + // this thread. Using `.disown` is a bit of a hack, but 283 + // disown is handled in handleConnection the way we want 284 + // WE DO NOT WANT to return .close, else that would result 285 + // in posix.close(socket) being called on an already-closed 286 + // socket, which would panic. 287 + return .disown; 288 + }, 289 + else => {}, 290 + } 291 + } 292 + }, 293 + else => {}, 294 + } 295 + requestError(conn, err) catch {}; 296 + posix.close(stream.handle); 297 + return .disown; 289 298 }; 290 299 291 300 if (done) { ··· 583 592 // can access _state directly. 584 593 585 594 const stream = http_conn.stream; 586 - const done = http_conn.req_state.parse(http_conn.req_arena.allocator(), stream) catch |err| { 595 + var reader = stream.reader(&.{}); // Request.State does its own buffering 596 + const done = http_conn.req_state.parse(http_conn.req_arena.allocator(), reader.interface()) catch |err| { 587 597 // maybe a write fail or something, doesn't matter, we're closing the connection 588 - requestError(http_conn, err) catch {}; 598 + requestError(http_conn, reader.getError() orelse err) catch {}; 589 599 590 600 // impossible to fail when false is passed 591 601 http_conn.requestDone(self.retain_allocated_bytes, false) catch unreachable; ··· 1775 1785 metrics.bodyTooBig(); 1776 1786 return writeError(handle, 413, "Request body is too big"); 1777 1787 }, 1778 - error.BrokenPipe, error.ConnectionClosed, error.ConnectionResetByPeer => return, 1788 + error.BrokenPipe, error.ConnectionClosed, error.ConnectionResetByPeer, error.EndOfStream => return, 1779 1789 else => { 1780 1790 log.err("server error: {}", .{err}); 1781 1791 metrics.internalError();