Skip to content

Commit f96897e

Browse files
committed
feat(http): let another thread cancel a request in flight
There was no way to abandon a request that had been sent. A caller that stopped caring (a user pressing Esc, say) could only leave the call running on a thread and wait out readTimeoutMs, with the connection held open and the server still producing the response. send and send_stream take an optional std::stop_token. While a token is attached, every wait on the socket is made in slices of at most 50 ms with a check for a stop request between them, and bio_recv waits for the socket before reading, so a handshake or the second half of a TLS record is covered too. When the token is stopped the call returns within about one slice, and the connection is dropped rather than pooled, so the server sees it end. A close_notify goes out only when the stop comes after the handshake has finished. HttpResponse::cancelled says why the call ended, without matching text. Before a status line has arrived the response is statusCode 0 with statusText and bodyError "Cancelled"; after it, statusCode is the server's and bodyError is "cancelled". A cancelled request is not resent on a fresh connection by the stale-connection retry, and a redirect is not followed. proxy_tunnel and proxy_connect take the token as a trailing default argument and hand it to the socket or the TLS session they open, so the exchange with an HTTP, https or SOCKS5 proxy is covered as well. Without a token (stop_possible() is false) every wait is the single poll it was, and bio_recv does not poll first. Not covered: getaddrinfo, a write blocked because the peer is not reading, and connect on openkal, where connect completes synchronously. download_to_file is unchanged. Tests in test_cancel run against the in-process TLS server and a listener that accepts and says nothing: stopped while waiting for the response head, while a stream waits for its next chunk, and during the handshake, each with a 30 s read timeout so only the token can have ended it; a token already stopped sends nothing; a token never stopped leaves a 300 ms timeout alone; a cancelled connection is not reused; a cancelled request on a pooled connection is not retried; and 20 cancelled requests leave the descriptor count where it was. Not run on macOS or Windows; the change adds no platform-specific code beyond the poll the sockets already use.
1 parent f461940 commit f96897e

8 files changed

Lines changed: 570 additions & 30 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,23 @@
11
# Changelog
22

3+
## Unreleased
4+
5+
`send` and `send_stream` take an optional `std::stop_token`, so another thread
6+
can abandon a request that is in flight.
7+
8+
* When the token is stopped the call returns within about 50 ms with
9+
`HttpResponse::cancelled` set, and the connection is closed instead of being
10+
returned to the pool. It covers connect, the proxy CONNECT exchange, the TLS
11+
handshake, waiting for the response head, reading the body and the gaps
12+
between streaming callbacks.
13+
* Before the status line arrives the response is `statusCode` 0 with
14+
`statusText` and `bodyError` `Cancelled`. After it, `statusCode` is the
15+
server's and `bodyError` is `cancelled`. A cancelled request is not retried on
16+
a new connection and its redirect is not followed.
17+
* Not covered: name resolution, and a write that is blocked because the server
18+
is not reading. Without a token nothing changes.
19+
* `proxy_connect` takes the token as a trailing default argument.
20+
321
## 0.3.3
422

523
**Behaviour change, and the reason for this release:** `verifySsl = true` (the

‎README.md‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,24 @@ else { /* resp.body is all of it */ }
4646
`ok()` deliberately does not consult `bodyComplete`, so existing `if (res.ok())`
4747
means exactly what it did before.
4848
49+
### Cancelling a request
50+
51+
Pass a `std::stop_token` to `send` or `send_stream` and stop its source from
52+
another thread. The call returns within about 50 ms with `cancelled` set, and
53+
the connection is closed rather than pooled.
54+
55+
```cpp
56+
std::stop_source source;
57+
auto worker = std::jthread([&] { response = client.send(request, source.get_token()); });
58+
// elsewhere:
59+
source.request_stop();
60+
```
61+
62+
If the status line had not arrived, `statusCode` is 0 and `statusText` is
63+
`Cancelled`; otherwise they are the server's and `bodyError` is `cancelled`.
64+
Name resolution and a write blocked on a server that is not reading cannot be
65+
interrupted.
66+
4967
### Configuration
5068

5169
| field | default | what it decides |

‎src/http.cppm‎

Lines changed: 67 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,11 @@ export struct HttpResponse {
5858
// "Invalid chunk size: zz" where the server had said "OK".
5959
std::string bodyError;
6060

61+
// True when the caller's stop token ended the request. Then `bodyComplete`
62+
// is false, and the status line is the server's if it had arrived and
63+
// 0 / "Cancelled" if not.
64+
bool cancelled { false };
65+
6166
bool ok() const { return statusCode >= 200 && statusCode < 300; }
6267
};
6368

@@ -880,14 +885,21 @@ public:
880885
HttpClient(HttpClient&&) = default;
881886
HttpClient& operator=(HttpClient&&) = default;
882887

883-
HttpResponse send(const HttpRequest& request) {
884-
return send_impl(request, 0);
888+
// `stop` can be stopped from another thread to abandon the request: the
889+
// call returns within about 50 ms with `cancelled` set, and the connection
890+
// is closed rather than pooled. It cannot interrupt name resolution or a
891+
// write blocked because the server is not reading.
892+
HttpResponse send(const HttpRequest& request, std::stop_token stop = {}) {
893+
return send_impl(request, 0, stop);
885894
}
886895

887896
// Streaming SSE request — reads the response body incrementally, feeding
888897
// chunks through SseParser to the caller's callback. The callback receives
889898
// each SseEvent and returns true to continue or false to stop.
890-
HttpResponse send_stream(const HttpRequest& request, SseCallbackFn callback) {
899+
// `stop` is as for send().
900+
HttpResponse send_stream(const HttpRequest& request, SseCallbackFn callback,
901+
std::stop_token stop = {}) {
902+
if (stop.stop_requested()) return cancelled_response();
891903
HttpResponse response;
892904

893905
auto parsed = parse_url(request.url);
@@ -901,12 +913,13 @@ public:
901913
PooledConnection guard(pool_, poolKey);
902914

903915
auto exchange = perform_exchange(parsed, poolKey,
904-
build_request(request, parsed), guard);
916+
build_request(request, parsed), guard, stop);
905917
if (!exchange.error.empty()) {
906918
response.statusCode = 0;
907919
response.statusText = exchange.error;
908920
response.bodyComplete = false;
909921
response.bodyError = exchange.error;
922+
response.cancelled = exchange.cancelled;
910923
return response;
911924
}
912925

@@ -921,6 +934,7 @@ public:
921934
SseParser parser;
922935

923936
auto sink = [&](std::string_view data) -> bool {
937+
if (stop.stop_requested()) return false;
924938
if (captureBody) {
925939
append_within_limit(response.body, data, stream_error_body_limit);
926940
}
@@ -932,7 +946,7 @@ public:
932946

933947
finish_body(response, exchange, request.method,
934948
/*maxBytes=*/std::numeric_limits<std::int64_t>::max(),
935-
sink, guard);
949+
sink, guard, stop);
936950
return response;
937951
}
938952

@@ -959,8 +973,18 @@ private:
959973
TlsSocket* sock { nullptr };
960974
ResponseHead head;
961975
std::string error;
976+
bool cancelled { false };
962977
};
963978

979+
static HttpResponse cancelled_response() {
980+
HttpResponse response;
981+
response.statusText = "Cancelled";
982+
response.bodyComplete = false;
983+
response.bodyError = "Cancelled";
984+
response.cancelled = true;
985+
return response;
986+
}
987+
964988
std::string build_request(const HttpRequest& request, const ParsedUrl& parsed) const {
965989
std::string reqStr;
966990
reqStr += method_to_string(request.method);
@@ -999,11 +1023,12 @@ private:
9991023

10001024
// False on failure, with `error` set when there is more to say than that the
10011025
// connection failed: a proxy that refused the tunnel says so in its own words.
1002-
bool open_connection(TlsSocket& sock, const ParsedUrl& parsed, std::string& error) {
1026+
bool open_connection(TlsSocket& sock, const ParsedUrl& parsed, std::string& error,
1027+
std::stop_token stop) {
10031028
if (config_.proxy.has_value()) {
10041029
auto tunnel = proxy_tunnel(parse_proxy_url(config_.proxy.value()),
10051030
parsed.host, parsed.port,
1006-
config_.connectTimeoutMs, config_.verifySsl);
1031+
config_.connectTimeoutMs, config_.verifySsl, stop);
10071032
if (!tunnel.ok()) {
10081033
error = std::move(tunnel.error);
10091034
return false;
@@ -1028,8 +1053,17 @@ private:
10281053
// request, so sending it again on a fresh connection is safe — once, and
10291054
// only when no response byte has arrived. See `retryOnStaleConnection`.
10301055
Exchange perform_exchange(const ParsedUrl& parsed, const std::string& poolKey,
1031-
const std::string& reqStr, PooledConnection& guard) {
1056+
const std::string& reqStr, PooledConnection& guard,
1057+
std::stop_token stop = {}) {
1058+
// Drops the connection; a stop request outranks whatever went wrong.
1059+
auto fail = [&](std::string message) -> Exchange {
1060+
guard.drop();
1061+
if (stop.stop_requested()) return { nullptr, {}, "Cancelled", true };
1062+
return { nullptr, {}, std::move(message) };
1063+
};
1064+
10321065
for (int attempt = 0; attempt < 2; ++attempt) {
1066+
if (stop.stop_requested()) return fail("");
10331067
bool reused = false;
10341068
TlsSocket* sock = nullptr;
10351069

@@ -1041,15 +1075,17 @@ private:
10411075
if (it != pool_.end()) pool_.erase(it);
10421076
auto [inserted, ok] = pool_.emplace(poolKey, TlsSocket{});
10431077
sock = &inserted->second;
1078+
}
1079+
sock->set_stop(stop);
1080+
if (!reused) {
10441081
std::string openError;
1045-
if (!open_connection(*sock, parsed, openError)) {
1082+
if (!open_connection(*sock, parsed, openError, stop)) {
10461083
// The proxy's refusal, else the TLS session's reason,
10471084
// else the TCP connection failed and there is no more to say.
10481085
std::string why = std::move(openError);
10491086
if (why.empty()) why = sock->error();
10501087
if (why.empty()) why = "Connection failed";
1051-
guard.drop();
1052-
return { nullptr, {}, std::move(why) };
1088+
return fail(std::move(why));
10531089
}
10541090
}
10551091

@@ -1065,8 +1101,7 @@ private:
10651101
// connection therefore yields exactly one execution.
10661102
if (!write_all(*sock, reqStr, config_.readTimeoutMs)) {
10671103
if (mayRetry) { guard.reset(); continue; }
1068-
guard.drop();
1069-
return { nullptr, {}, "Write failed" };
1104+
return fail("Write failed");
10701105
}
10711106

10721107
auto head = read_response_head(*sock, config_.readTimeoutMs);
@@ -1076,18 +1111,16 @@ private:
10761111
// arrived the server has seen the request, and repeating it could
10771112
// repeat its effect.
10781113
if (mayRetry && !head.error().sawBytes) { guard.reset(); continue; }
1079-
guard.drop();
1080-
return { nullptr, {}, head.error().message };
1114+
return fail(head.error().message);
10811115
}
1082-
guard.drop();
1083-
return { nullptr, {}, "No response" };
1116+
return fail("No response");
10841117
}
10851118

10861119
// Reads the body into `sink` and settles the pool guard by what the read
10871120
// found. The single place that decides whether a connection is reusable.
10881121
void finish_body(HttpResponse& response, Exchange& exchange, Method method,
10891122
std::int64_t maxBytes, const BodySink& sink,
1090-
PooledConnection& guard) {
1123+
PooledConnection& guard, std::stop_token stop = {}) {
10911124
if (!response_has_body(method, exchange.head.statusCode)) {
10921125
if (config_.keepAlive && !exchange.head.connectionClose) guard.keep();
10931126
return;
@@ -1103,18 +1136,28 @@ private:
11031136
// Not an error, but bytes are still owed on the socket.
11041137
response.bodyComplete = false;
11051138
response.bodyError = "stopped by callback";
1139+
if (stop.stop_requested()) {
1140+
response.bodyError = "cancelled";
1141+
response.cancelled = true;
1142+
}
11061143
break;
11071144
case BodyEnd::ClosedByPeer:
11081145
// The body ended where it said it would; the connection did too.
11091146
break;
11101147
case BodyEnd::Truncated:
11111148
response.bodyComplete = false;
11121149
response.bodyError = outcome.error;
1150+
if (stop.stop_requested()) {
1151+
response.bodyError = "cancelled";
1152+
response.cancelled = true;
1153+
}
11131154
break;
11141155
}
11151156
}
11161157

1117-
HttpResponse send_impl(const HttpRequest& request, int redirectCount) {
1158+
HttpResponse send_impl(const HttpRequest& request, int redirectCount,
1159+
std::stop_token stop) {
1160+
if (stop.stop_requested()) return cancelled_response();
11181161
HttpResponse response;
11191162

11201163
auto parsed = parse_url(request.url);
@@ -1128,12 +1171,13 @@ private:
11281171
PooledConnection guard(pool_, poolKey);
11291172

11301173
auto exchange = perform_exchange(parsed, poolKey,
1131-
build_request(request, parsed), guard);
1174+
build_request(request, parsed), guard, stop);
11321175
if (!exchange.error.empty()) {
11331176
response.statusCode = 0;
11341177
response.statusText = exchange.error;
11351178
response.bodyComplete = false;
11361179
response.bodyError = exchange.error;
1180+
response.cancelled = exchange.cancelled;
11371181
return response;
11381182
}
11391183

@@ -1157,10 +1201,10 @@ private:
11571201
response.body.append(data);
11581202
return true;
11591203
},
1160-
guard);
1204+
guard, stop);
11611205

11621206
// Follow 3xx redirects if configured.
1163-
if (config_.maxRedirects > 0 &&
1207+
if (config_.maxRedirects > 0 && !response.cancelled &&
11641208
response.statusCode >= 300 && response.statusCode < 400 &&
11651209
redirectCount < config_.maxRedirects) {
11661210
std::string location = find_header(response.headers, "Location");
@@ -1181,7 +1225,7 @@ private:
11811225
// connection into the pool under this same key, and a guard
11821226
// still armed would delete it when this scope ends.
11831227
guard.drop();
1184-
return send_impl(redirectReq, redirectCount + 1);
1228+
return send_impl(redirectReq, redirectCount + 1, stop);
11851229
}
11861230
}
11871231

‎src/proxy.cppm‎

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -351,8 +351,10 @@ export struct ProxyTunnel {
351351
// setting that governs the connection to the target.
352352
export ProxyTunnel proxy_tunnel(const ProxyConfig& proxy,
353353
std::string_view targetHost, int targetPort,
354-
int timeoutMs, bool verifySsl = true) {
354+
int timeoutMs, bool verifySsl = true,
355+
std::stop_token stop = {}) {
355356
ProxyTunnel tunnel;
357+
tunnel.socket.set_stop(stop);
356358
const std::string where = proxy.host + ":" + std::to_string(proxy.port);
357359

358360
if (proxy.scheme == "http" || proxy.scheme == "socks5" || proxy.scheme == "socks5h") {
@@ -366,6 +368,7 @@ export ProxyTunnel proxy_tunnel(const ProxyConfig& proxy,
366368
if (!tunnel.error.empty()) tunnel.socket.close();
367369
} else if (proxy.scheme == "https") {
368370
auto tls = std::make_unique<TlsSocket>();
371+
tls->set_stop(stop);
369372
if (!tls->connect(proxy.host.c_str(), proxy.port, timeoutMs, verifySsl)) {
370373
// The session says why when the TCP connection was up, which is
371374
// where a proxy certificate that does not verify is refused.
@@ -385,11 +388,11 @@ export ProxyTunnel proxy_tunnel(const ProxyConfig& proxy,
385388
// say why it failed; `proxy_tunnel` can.
386389
export Socket proxy_connect(std::string_view proxyHost, int proxyPort,
387390
std::string_view targetHost, int targetPort,
388-
int timeoutMs) {
391+
int timeoutMs, std::stop_token stop = {}) {
389392
ProxyConfig proxy;
390393
proxy.host = std::string(proxyHost);
391394
proxy.port = proxyPort;
392-
auto tunnel = proxy_tunnel(proxy, targetHost, targetPort, timeoutMs);
395+
auto tunnel = proxy_tunnel(proxy, targetHost, targetPort, timeoutMs, true, stop);
393396
return std::move(tunnel.socket);
394397
}
395398

0 commit comments

Comments
 (0)