Skip to content

Commit 093719a

Browse files
committed
Enhance HTTP handling and connection management
- Added timeout checks before reading HTTP headers and body to prevent dangling async operations during shutdown. - Improved error response handling to ensure proper framing in multi-process mode. - Introduced a new method to check if the peer has closed the connection in the network layer. - Refactored HTTP response reading logic to handle peer closure and ensure compliance with Content-Length. - Updated tests to verify that original status codes are preserved and that responses are correctly handled in various scenarios.
1 parent d0448e0 commit 093719a

11 files changed

Lines changed: 179 additions & 123 deletions

File tree

.github/workflows/ci.yml

Lines changed: 28 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -206,27 +206,47 @@ jobs:
206206
cspkg list
207207
208208
# ------------------------------------------------------------------ #
209-
# 10b. Reproducible build check (.csp only — date line is volatile) #
209+
# 10b. Reproducible build check (Date header line is volatile) #
210210
# ------------------------------------------------------------------ #
211211
- name: Reproducible build check
212212
shell: bash
213213
run: |
214214
set -e
215-
for pkg in argparse netutils; do
215+
# Check every source package in the repo, not a hardcoded list.
216+
# nullglob avoids iterating the literal glob when nothing matches.
217+
shopt -s nullglob
218+
for src in *.ecs; do
219+
pkg="${src%.ecs}"
216220
echo "=== Checking ${pkg}.csp reproducibility ==="
217221
cp "${pkg}.csp" "/tmp/${pkg}.csp.bak"
222+
if [ -f "${pkg}.csym" ]; then
223+
cp "${pkg}.csym" "/tmp/${pkg}.csym.bak"
224+
fi
218225
cspkg build "${pkg}.ecs" --compile
219-
# The .csp header contains a Date line (line 3) that always
220-
# differs. Strip it before comparing; everything else must
226+
# The .csp header contains a Date line that always differs.
227+
# Strip it by pattern before comparing; everything else must
221228
# be byte-identical when the source has not changed.
222-
sed '3d' "/tmp/${pkg}.csp.bak" > "/tmp/${pkg}.csp.old"
223-
sed '3d' "${pkg}.csp" > "/tmp/${pkg}.csp.new"
229+
sed '/^# Date:/d' "/tmp/${pkg}.csp.bak" > "/tmp/${pkg}.csp.old"
230+
sed '/^# Date:/d' "${pkg}.csp" > "/tmp/${pkg}.csp.new"
224231
if ! diff -q "/tmp/${pkg}.csp.old" "/tmp/${pkg}.csp.new"; then
225232
echo "ERROR: ${pkg}.csp changed beyond the date line!"
226233
echo " Did you forget to recompile after editing ${pkg}.ecs?"
227234
diff "/tmp/${pkg}.csp.old" "/tmp/${pkg}.csp.new" || true
228235
exit 1
229236
fi
237+
# The .csym has no volatile header — it must be byte-identical.
238+
if [ -f "/tmp/${pkg}.csym.bak" ]; then
239+
if [ ! -f "${pkg}.csym" ]; then
240+
echo "ERROR: ${pkg}.csym is missing after build!"
241+
echo " Did you forget to regenerate ${pkg}.csym after editing ${pkg}.ecs?"
242+
exit 1
243+
fi
244+
if ! cmp -s "/tmp/${pkg}.csym.bak" "${pkg}.csym"; then
245+
echo "ERROR: ${pkg}.csym is stale!"
246+
echo " Did you forget to regenerate ${pkg}.csym after editing ${pkg}.ecs?"
247+
exit 1
248+
fi
249+
fi
230250
echo " ${pkg} reproducible — ok"
231251
done
232252
@@ -294,7 +314,7 @@ jobs:
294314
cs tests/test_tls_errors.csc
295315
296316
# ------------------------------------------------------------------ #
297-
# 16c. Run master-slave integration tests #
317+
# 16b. Run master-slave integration tests #
298318
# ------------------------------------------------------------------ #
299319
- name: Run master-slave integration tests
300320
shell: bash
@@ -304,7 +324,7 @@ jobs:
304324
cs tests/test_http_compliance.csc
305325
306326
# ------------------------------------------------------------------ #
307-
# 16b. Run TLS custom trust mode test (self-signed cert) #
327+
# 16c. Run TLS custom trust mode test (self-signed cert) #
308328
# ------------------------------------------------------------------ #
309329
- name: Run TLS custom trust mode test
310330
shell: bash

CNI_API.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -105,6 +105,7 @@ sock.connect_ssl("localhost", {"trust_mode": "insecure"}.to_hash_map())
105105
| `set_opt_no_delay` | `(value: boolean)` | 设置 `TCP_NODELAY` 选项(禁用 Nagle 算法) |
106106
| `set_opt_keep_alive` | `(value: boolean)` | 设置 `SO_KEEPALIVE` 选项 |
107107
| `available` | `() → int` | 可读取的字节数(非阻塞) |
108+
| `peer_closed` | `() → boolean` | 对端是否已关闭连接(非阻塞、非破坏性,使用 1 字节 `MSG_PEEK` 探测)。空闲但存活返回 `false`;收到 FIN 或连接已失效返回 `true`。有异步读挂起时返回 `false`(该读操作自身会暴露 EOF)。TLS 套接字上探测的是底层传输而非解密流 |
108109
| `receive` | `(max: int) → string` | 读取最多 `max` 字节。阻塞直到至少 1 字节可读 |
109110
| `read` | `(size: int) → string` | 读取恰好 `size` 字节。阻塞直到全部读完 |
110111
| `send` | `(data: string) → int` | 发送数据(单次部分写入,返回实际写入字节数)。需完整发送时使用 `write` |

NETUTILS.md

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -181,7 +181,7 @@ Master 分配 `rank`:先尝试使用 `deprecated_rank`(回收的编号),
181181
控制逻辑:
182182

183183
* 若某连接空闲时间超过 `keep_alive_timeout`,Master 会在关闭前尝试回写 `408 Request Timeout`(或直接关闭)。
184-
* 若单条连接处理的请求数超过 `max_keep_alive`Master 会主动向客户端发送 `408` 并关闭连接
184+
* 若单条连接处理的请求数达到 `max_keep_alive`服务器会正常处理最后一个请求,并在该响应中携带 `Connection: close`,随后关闭连接(自 2.1 起不再发送 `408`
185185
* Slave 的请求处理若超时(`receive_content_s` 超时),Master 将用错误码构造响应并关闭该 Slave 连接。
186186

187187
## 7. 静态文件服务与安全(wwwroot、path_normalize)
@@ -443,8 +443,8 @@ sequenceDiagram
443443
M->>C: 发回 HTTP 回复
444444
end
445445
446-
opt 超出 Keep-alive 限制
447-
M->>C: 发回 HTTP 回复(408)
446+
opt 达到 Keep-alive 请求数上限
447+
M->>C: 最后一个 HTTP 回复携带 Connection: close<br/>随后关闭连接
448448
end
449449
```
450450

README.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@ A high-performance network extension for the [Covariant Script](http://covscript
88

99
| Package | Type | Version | Description |
1010
|---|---|---|---|
11-
| `network` | C++ Extension | `1.38.0_v6.4` | TCP/UDP sockets, TLS/SSL, async I/O, event loop |
11+
| `network` | C++ Extension | `1.38.0_v6.5` | TCP/UDP sockets, TLS/SSL, async I/O, event loop |
1212
| `netutils` | CovScript | `2.1` | HTTP server/client framework with single-process, distributed master/slave, and OpenAI API modes |
1313
| `argparse` | CovScript | `1.1` | Lightweight command-line argument parser |
1414

csbuild/network.json

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@
33
"Name": "network",
44
"Info": "Socket Extension",
55
"Author": "CovScript Organization",
6-
"Version": "1.38.0_v6.4",
6+
"Version": "1.38.0_v6.5",
77
"Target": "build/imports/network.cse",
88
"Dependencies": []
99
}

include/network.hpp

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -717,6 +717,38 @@ namespace cs_impl {
717717
return sock.available();
718718
}
719719

720+
// For TLS sockets the probe reflects the raw transport, not
721+
// the decrypted stream — MSG_PEEK sees encrypted bytes, so
722+
// peer_closed() detects transport-level FIN but not TLS
723+
// close_notify. Callers that need TLS-level closure should
724+
// read the stream until eof.
725+
bool peer_closed() noexcept
726+
{
727+
if (!try_begin_io_job(io_direction::read))
728+
return false;
729+
bool closed = false;
730+
if (!sock.is_open()) {
731+
closed = true;
732+
}
733+
else {
734+
asio::error_code ec;
735+
sock.non_blocking(true, ec);
736+
if (!ec) {
737+
char probe = 0;
738+
std::size_t peeked = sock.receive(
739+
asio::buffer(&probe, 1),
740+
asio::socket_base::message_peek, ec);
741+
asio::error_code restore_ec;
742+
sock.non_blocking(false, restore_ec);
743+
closed = peeked == 0 &&
744+
ec != asio::error::would_block &&
745+
ec != asio::error::try_again;
746+
}
747+
}
748+
finish_io_job(io_direction::read);
749+
return closed;
750+
}
751+
720752
std::string receive(std::size_t maximum)
721753
{
722754
scoped_io_job job(*this, io_direction::read);

netutils.csp

Lines changed: 13 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
# Generated by Extended CovScript Compiler
22
# DO NOT MODIFY
3-
# Date: Wed Jul 15 20:58:29 2026
3+
# Date: Wed Jul 15 23:03:06 2026
44
@charset: utf8
55
import ecs as netutils_ecs
66
struct __netutils_ecs_lambda_impl_1__
@@ -479,12 +479,12 @@ function read_http_header(sock, state, keep_alive_timeout, max_body_size)
479479
var error_code = null
480480
var header_size = 0
481481
loop
482-
async.read_until(sock, state, "\r\n")
483482
var timeout = keep_alive_timeout - runtime.time()
484483
if timeout <= 0
485484
error_code = state_codes.code_408
486485
break
487486
end
487+
async.read_until(sock, state, "\r\n")
488488
if !state.wait_for(timeout)
489489
if state.has_done()
490490
if state.eof()
@@ -542,12 +542,12 @@ function read_http_header(sock, state, keep_alive_timeout, max_body_size)
542542
session.post_data = state.get_buffer(session.content_length)
543543
var remaining = session.content_length - session.post_data.size
544544
while remaining > 0
545-
state = async.read(sock, remaining)
546545
var timeout = keep_alive_timeout - runtime.time()
547546
if timeout <= 0
548547
error_code = state_codes.code_408
549548
break
550549
end
550+
state = async.read(sock, remaining)
551551
if !state.wait_for(timeout)
552552
if state.has_done()
553553
if state.eof()
@@ -642,7 +642,8 @@ function call_http_handler(session, server)
642642
if server->url_map.exist(error_code)
643643
server->url_map[error_code](*server, session)
644644
else
645-
send_error_response(session.sock, error_code)
645+
session.connection = "close"
646+
session.send_response(error_code, "", "text/html")
646647
end
647648
return false
648649
else
@@ -680,19 +681,22 @@ function simple_worker(self)
680681
if session == null
681682
break
682683
end
684+
var force_close = false
683685
if ++request_count >= self->server->max_keep_alive
684686
session.connection = "close"
687+
force_close = true
685688
end
686689
session.sock = sock
687-
if !call_http_handler(session, self->server)
688-
break
689-
end
690+
var handler_ok = call_http_handler(session, self->server)
690691
if session.response_state != null && !session.response_state.wait()
691692
log("Write response error: " + session.response_state.get_error())
692693
break
693694
end
695+
if !handler_ok
696+
break
697+
end
694698
last_request_time = runtime.time()
695-
if !sock.is_open() || session.connection == "close"
699+
if !sock.is_open() || force_close || session.connection == "close"
696700
break
697701
end
698702
end
@@ -705,7 +709,6 @@ struct http_conn
705709
var sock = null
706710
var read_state = null
707711
var state = 0
708-
var close = false
709712
var keep_alive = true
710713
var request_count = 0
711714
var last_request_time = 0
@@ -830,7 +833,6 @@ function master_request_worker(self)
830833
conn->last_request_time = runtime.time()
831834
if session.connection == "close"
832835
conn->keep_alive = false
833-
conn->close = true
834836
end
835837
conn->request_queue.push_back(move(session))
836838
conn->state = 0
@@ -864,21 +866,18 @@ function master_response_worker(self)
864866
if !response_state.wait()
865867
log("Write response error: " + response_state.get_error())
866868
conn->keep_alive = false
867-
conn->close = true
868869
conn->request_queue = new array
869870
break
870871
end
871872
conn->request_queue.pop_front()
872873
--conn->request_idx
874+
conn->last_request_time = runtime.time()
873875
else
874876
fiber.yield()
875877
break
876878
end
877879
end
878880
if conn->request_queue.empty() && !conn->keep_alive
879-
if !conn->close
880-
send_error_response(conn->sock, state_codes.code_408)
881-
end
882881
if !conn->sock.safe_shutdown()
883882
log("safe_shutdown returned false — async jobs may still be pending")
884883
end

netutils.csym

Lines changed: 26 additions & 15 deletions
Large diffs are not rendered by default.

netutils.ecs

Lines changed: 25 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -458,12 +458,14 @@ function read_http_header(sock, state, keep_alive_timeout, max_body_size)
458458
var error_code = null
459459
var header_size = 0
460460
loop
461-
async.read_until(sock, state, "\r\n")
461+
# Check the deadline before starting the read: breaking with a
462+
# just-started async op pending would leave it dangling past shutdown.
462463
var timeout = keep_alive_timeout - runtime.time()
463464
if timeout <= 0
464465
error_code = state_codes.code_408
465466
break
466467
end
468+
async.read_until(sock, state, "\r\n")
467469
if !state.wait_for(timeout)
468470
if state.has_done()
469471
if state.eof()
@@ -524,12 +526,13 @@ function read_http_header(sock, state, keep_alive_timeout, max_body_size)
524526
session.post_data = state.get_buffer(session.content_length)
525527
var remaining = session.content_length - session.post_data.size
526528
while remaining > 0
527-
state = async.read(sock, remaining)
529+
# Same as above: check the deadline before starting the read.
528530
var timeout = keep_alive_timeout - runtime.time()
529531
if timeout <= 0
530532
error_code = state_codes.code_408
531533
break
532534
end
535+
state = async.read(sock, remaining)
533536
if !state.wait_for(timeout)
534537
if state.has_done()
535538
if state.eof()
@@ -635,7 +638,11 @@ function call_http_handler(session, server)
635638
if server->url_map.exist(error_code)
636639
server->url_map[error_code](*server, session)
637640
else
638-
send_error_response(session.sock, error_code)
641+
# Route through the session so multi-process mode sends a framed
642+
# response to the master instead of raw HTTP bytes on the
643+
# slave<->master socket (which would desynchronize the framing).
644+
session.connection = "close"
645+
session.send_response(error_code, "", "text/html")
639646
end
640647
return false
641648
else
@@ -684,21 +691,27 @@ function simple_worker(self)
684691
end
685692
# Close cleanly once max_keep_alive requests are served: advertise
686693
# Connection: close on the final response
694+
var force_close = false
687695
if ++request_count >= self->server->max_keep_alive
688696
session.connection = "close"
697+
force_close = true
689698
end
690699
# Call handler
691700
session.sock = sock
692-
if !call_http_handler(session, self->server)
693-
break
694-
end
701+
var handler_ok = call_http_handler(session, self->server)
702+
# Wait for the response write (including error responses sent on
703+
# handler failure) before deciding the connection's fate.
695704
if session.response_state != null && !session.response_state.wait()
696705
log("Write response error: " + session.response_state.get_error())
697706
break
698707
end
708+
if !handler_ok
709+
break
710+
end
699711
last_request_time = runtime.time()
700-
# Keep-Alive check
701-
if !sock.is_open() || session.connection == "close"
712+
# Keep-Alive check — force_close guards against handlers that
713+
# overwrite session.connection after the limit was reached
714+
if !sock.is_open() || force_close || session.connection == "close"
702715
break
703716
end
704717
end
@@ -714,7 +727,6 @@ struct http_conn
714727
var read_state = null
715728
# -1 = close, 0 = established, 1 = busy
716729
var state = 0
717-
var close = false
718730
var keep_alive = true
719731
var request_count = 0
720732
var last_request_time = 0
@@ -853,7 +865,6 @@ function master_request_worker(self)
853865
# Check connection type
854866
if session.connection == "close"
855867
conn->keep_alive = false
856-
conn->close = true
857868
end
858869
conn->request_queue.push_back(move(session))
859870
conn->state = 0
@@ -890,21 +901,21 @@ function master_response_worker(self)
890901
if !response_state.wait()
891902
log("Write response error: " + response_state.get_error())
892903
conn->keep_alive = false
893-
conn->close = true
894904
conn->request_queue = new array
895905
break
896906
end
897907
conn->request_queue.pop_front()
898908
--conn->request_idx
909+
# Restart the keep-alive idle window only once the response is
910+
# written, matching simple_worker: slave processing time must
911+
# not eat into the client's idle allowance.
912+
conn->last_request_time = runtime.time()
899913
else
900914
fiber.yield()
901915
break
902916
end
903917
end
904918
if conn->request_queue.empty() && !conn->keep_alive
905-
if !conn->close
906-
send_error_response(conn->sock, state_codes.code_408)
907-
end
908919
if !conn->sock.safe_shutdown()
909920
log("safe_shutdown returned false — async jobs may still be pending")
910921
end

0 commit comments

Comments
 (0)