Skip to content

Commit a778ca1

Browse files
committed
fix(cancel): a token belongs to its call; downloads take one; openkal Windows has entropy
The token is set on the socket at the bottom of a connection and nowhere else. TlsSocket::set_stop forwards to the session beneath it, which is where every wait happens inside an https:// proxy's tunnel; setting only its own socket left a reused tunnel on the token of the call that opened it. Once that token was stopped, which a std::jthread does in its destructor, every wait on the tunnel failed at once and the stale-connection retry sent a POST that had already arrived a second time; a later call's own token never reached the tunnel at all. PooledConnection::keep() takes the token off the connection it returns to the pool, so nothing in the pool carries one. download_to_file takes a std::stop_token as its last argument, and DownloadToFileResult::cancelled is set by it or by isCancelled, as HttpResponse::cancelled is. proxy_connect, the wrapper nothing here calls, does not take one. Where the C library is musl, the TLS random generator is seeded from musl's getentropy rather than mbedTLS's reading of /dev/urandom, which above openkal on x86_64-windows-musl does not exist: every handshake there failed with "CTR_DRBG - The entropy source failed". Tests: a stop after the call returned does not reach the next call, and the next call's token does, each directly, through an http:// proxy and through an https:// proxy; a SOCKS5 proxy that never answers; a download cancelled before the head, during the body, with a token already stopped, and by isCancelled. test_cancel imports std before its includes under libc++, whose 20 and 22 fail to link stop_source::request_stop() the other way round. examples/openkal follows mcpp-index's openkal pins (runtime 0.15.2) and cancels a handshake against a local listener, which needs no network; it runs on x86_64-linux-gnu, x86_64-linux-musl, aarch64-linux-musl and x86_64-windows-musl.
1 parent 8996984 commit a778ca1

6 files changed

Lines changed: 463 additions & 12 deletions

File tree

‎examples/openkal/mcpp.toml‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,7 @@ version = "0.1.0"
1111
# sources either way. That is the claim this directory exists to check.
1212
[dependencies]
1313
tinyhttps = { path = "../.." }
14-
openkal-llvm-runtime = "0.9.1"
14+
openkal-llvm-runtime = "0.15.2"
1515

1616
[targets.smoke]
1717
kind = "bin"
@@ -22,3 +22,11 @@ main = "src/main.cpp"
2222
# dependency above.
2323
[toolchain]
2424
default = "llvm@22.1.8"
25+
26+
# The programs of the other targets this example is run for, as mcpp-index's
27+
# openkal measurement runs them (tests/openkal/pins.toml there).
28+
[target.x86_64-windows-musl]
29+
runner = ["wine"]
30+
31+
[target.aarch64-linux-musl]
32+
runner = ["qemu-aarch64"]

‎examples/openkal/src/main.cpp‎

Lines changed: 73 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,11 +10,82 @@
1010
// is none it says so and exits 0, because a CI job without egress should report
1111
// "not run" rather than "broken". A certificate the client refused is not an
1212
// absence of network, and exits 1.
13+
// The listener below is made with the C library's sockets, which openkal-musl
14+
// provides on every target this example runs on. C headers only: they declare
15+
// nothing `import std` also declares.
16+
#include <arpa/inet.h>
17+
#include <netinet/in.h>
18+
#include <sys/socket.h>
19+
#include <unistd.h>
20+
1321
import mcpplibs.tinyhttps;
1422
import std;
1523

1624
namespace https = mcpplibs::tinyhttps;
1725

26+
namespace {
27+
28+
// CANCELLATION ABOVE openkal, WITH NO NETWORK.
29+
//
30+
// A listener that never accepts: the connection completes in the backlog, the
31+
// client sends its ClientHello and waits for an answer that is not coming. That
32+
// wait is a `poll` in slices of 50 ms, and above openkal each slice is a bounded
33+
// read (`kal_timeout_read`) rather than a readiness query, which is the part of
34+
// the stack this checks. The read timeout is 30 s, so returning promptly can
35+
// only be the token's doing.
36+
bool cancellation_works() {
37+
const int listener = ::socket(AF_INET, SOCK_STREAM, 0);
38+
if (listener < 0) { std::println("cancellation: no socket"); return false; }
39+
sockaddr_in addr {};
40+
addr.sin_family = AF_INET;
41+
addr.sin_addr.s_addr = htonl(INADDR_LOOPBACK);
42+
addr.sin_port = 0;
43+
socklen_t len = sizeof addr;
44+
if (::bind(listener, reinterpret_cast<sockaddr*>(&addr), sizeof addr) != 0
45+
|| ::listen(listener, 4) != 0
46+
|| ::getsockname(listener, reinterpret_cast<sockaddr*>(&addr), &len) != 0) {
47+
std::println("cancellation: no listener");
48+
::close(listener);
49+
return false;
50+
}
51+
const std::string url = "https://127.0.0.1:" + std::to_string(ntohs(addr.sin_port)) + "/";
52+
53+
https::HttpClientConfig config;
54+
config.verifySsl = false;
55+
config.connectTimeoutMs = 5000;
56+
config.readTimeoutMs = 30000;
57+
https::HttpClient client(config);
58+
https::HttpRequest request;
59+
request.method = https::Method::GET;
60+
request.url = url;
61+
62+
// A token stopped before the call: nothing is sent.
63+
std::stop_source early;
64+
early.request_stop();
65+
auto before = client.send(request, early.get_token());
66+
67+
// A token stopped while the handshake waits.
68+
std::stop_source source;
69+
std::jthread stopper([&] {
70+
std::this_thread::sleep_for(std::chrono::milliseconds(200));
71+
source.request_stop();
72+
});
73+
const auto start = std::chrono::steady_clock::now();
74+
auto during = client.send(request, source.get_token());
75+
const auto ms = std::chrono::duration_cast<std::chrono::milliseconds>(
76+
std::chrono::steady_clock::now() - start).count();
77+
::close(listener);
78+
79+
const bool ok = before.cancelled && before.statusCode == 0
80+
&& during.cancelled && during.statusCode == 0
81+
&& ms >= 150 && ms < 2000;
82+
std::println("cancellation: {} (stopped at 200 ms, returned after {} ms)",
83+
ok ? "ok" : "WRONG", ms);
84+
return ok;
85+
}
86+
87+
} // namespace
88+
1889
int main() {
1990
https::Socket::platform_init();
2091

@@ -30,6 +101,8 @@ int main() {
30101
}
31102
std::println("framing parsers: ok");
32103

104+
if (!cancellation_works()) return 1;
105+
33106
https::HttpClientConfig config;
34107
config.connectTimeoutMs = 15000;
35108
config.readTimeoutMs = 20000;

‎src/http.cppm‎

Lines changed: 45 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -149,6 +149,12 @@ export struct DownloadToFileResult {
149149
// failure downstream of a truncated file otherwise blames the wrong party.
150150
bool writeFailed { false };
151151

152+
// True when the caller abandoned the transfer: `isCancelled` returned true
153+
// or the stop token was stopped. `error` is then "Cancelled" if no status
154+
// line had arrived and "cancelled" if one had, as for `HttpResponse`. A
155+
// destination that failed first is reported as that instead.
156+
bool cancelled { false };
157+
152158
bool ok() const { return statusCode >= 200 && statusCode < 300 && error.empty(); }
153159
};
154160

@@ -842,7 +848,16 @@ public:
842848

843849
// Call only where the body was read to the end its framing declared and the
844850
// response did not ask for the connection to close.
845-
void keep() noexcept { armed_ = false; }
851+
//
852+
// A connection in the pool carries no stop token. The call that used it
853+
// set its own (`perform_exchange`), and the caller may stop that token
854+
// after this call has returned — a `std::jthread` does so in its
855+
// destructor — which must not reach whatever takes the connection next.
856+
void keep() noexcept {
857+
if (!armed_) return;
858+
armed_ = false;
859+
if (auto it = pool_.find(key_); it != pool_.end()) it->second.set_stop({});
860+
}
846861

847862
// Drop now rather than at the end of the scope. Idempotent, and a no-op
848863
// after `keep()`.
@@ -953,14 +968,16 @@ public:
953968
// Download URL to file with streaming progress.
954969
// Follows redirects. Calls onProgress periodically during body read.
955970
// isCancelled is checked after each block — return true to abort.
971+
// `stop` is as for send(), and also ends the wait for the response head.
956972
DownloadToFileResult download_to_file(
957973
const std::string& url,
958974
const std::filesystem::path& destFile,
959975
DownloadProgressFn onProgress = nullptr,
960-
std::function<bool()> isCancelled = nullptr)
976+
std::function<bool()> isCancelled = nullptr,
977+
std::stop_token stop = {})
961978
{
962979
return download_to_file_impl(url, destFile, std::move(onProgress),
963-
std::move(isCancelled), 0);
980+
std::move(isCancelled), 0, stop);
964981
}
965982

966983
HttpClientConfig& config() { return config_; }
@@ -1237,9 +1254,15 @@ private:
12371254
const std::filesystem::path& destFile,
12381255
DownloadProgressFn onProgress,
12391256
std::function<bool()> isCancelled,
1240-
int redirectCount)
1257+
int redirectCount,
1258+
std::stop_token stop)
12411259
{
12421260
DownloadToFileResult result;
1261+
if (stop.stop_requested()) {
1262+
result.error = "Cancelled";
1263+
result.cancelled = true;
1264+
return result;
1265+
}
12431266

12441267
auto parsed = parse_url(url);
12451268
if (parsed.scheme != "https") {
@@ -1256,9 +1279,10 @@ private:
12561279
request.headers = { {"User-Agent", "tinyhttps/1.0"}, {"Accept", "*/*"} };
12571280

12581281
auto exchange = perform_exchange(parsed, poolKey,
1259-
build_request(request, parsed), guard);
1282+
build_request(request, parsed), guard, stop);
12601283
if (!exchange.error.empty()) {
12611284
result.error = exchange.error;
1285+
result.cancelled = exchange.cancelled;
12621286
return result;
12631287
}
12641288

@@ -1283,14 +1307,20 @@ private:
12831307
} else {
12841308
guard.drop();
12851309
}
1310+
if (stop.stop_requested()) {
1311+
result.error = "cancelled";
1312+
result.cancelled = true;
1313+
return result;
1314+
}
12861315

12871316
if (location.starts_with("/")) {
12881317
location = parsed.scheme + "://" + parsed.host +
12891318
(parsed.port != 443 ? ":" + std::to_string(parsed.port) : "") +
12901319
location;
12911320
}
12921321
return download_to_file_impl(location, destFile, std::move(onProgress),
1293-
std::move(isCancelled), redirectCount + 1);
1322+
std::move(isCancelled), redirectCount + 1,
1323+
stop);
12941324
}
12951325
}
12961326

@@ -1383,7 +1413,10 @@ private:
13831413

13841414
written += size;
13851415
if (onProgress) onProgress(totalBytes, written);
1386-
if (isCancelled && isCancelled()) { cancelled = true; return false; }
1416+
if ((isCancelled && isCancelled()) || stop.stop_requested()) {
1417+
cancelled = true;
1418+
return false;
1419+
}
13871420
return true;
13881421
});
13891422

@@ -1401,16 +1434,20 @@ private:
14011434
result.error = cancelled ? "cancelled" : "stopped";
14021435
break;
14031436
case BodyEnd::Truncated:
1404-
result.error = outcome.error;
1437+
// A wait that the stop token ended reads as a truncated body.
1438+
cancelled = stop.stop_requested();
1439+
result.error = cancelled ? "cancelled" : outcome.error;
14051440
break;
14061441
}
1442+
result.cancelled = cancelled;
14071443

14081444
// The destination's failure outranks whatever the switch recorded: a
14091445
// stop that the file asked for is not a cancellation, and a body that
14101446
// arrived whole into a file that lost it has not succeeded.
14111447
if (writeFailed) {
14121448
result.error = write_failure();
14131449
result.writeFailed = true;
1450+
result.cancelled = false;
14141451
}
14151452
return result;
14161453
}

‎src/proxy.cppm‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -388,11 +388,11 @@ export ProxyTunnel proxy_tunnel(const ProxyConfig& proxy,
388388
// say why it failed; `proxy_tunnel` can.
389389
export Socket proxy_connect(std::string_view proxyHost, int proxyPort,
390390
std::string_view targetHost, int targetPort,
391-
int timeoutMs, std::stop_token stop = {}) {
391+
int timeoutMs) {
392392
ProxyConfig proxy;
393393
proxy.host = std::string(proxyHost);
394394
proxy.port = proxyPort;
395-
auto tunnel = proxy_tunnel(proxy, targetHost, targetPort, timeoutMs, true, stop);
395+
auto tunnel = proxy_tunnel(proxy, targetHost, targetPort, timeoutMs);
396396
return std::move(tunnel.socket);
397397
}
398398

‎src/tls.cppm‎

Lines changed: 44 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,21 @@ module;
1515
#pragma comment(lib, "bcrypt.lib")
1616
#endif
1717

18+
// WHERE THE C LIBRARY IS musl, THE ENTROPY IS THE C LIBRARY'S.
19+
//
20+
// mbedTLS's own platform source asks the kernel through `getrandom` only where
21+
// it recognises glibc (`entropy_poll.c`: `__linux__ && __GLIBC__`); with musl it
22+
// opens `/dev/urandom`. Above openkal that file is not a thing the C library can
23+
// promise — openkal reaches entropy through its own `openkal.random` interface,
24+
// and openkal-musl answers `getrandom` with it — so on x86_64-windows-musl every
25+
// handshake failed before it began, with "CTR_DRBG - The entropy source failed".
26+
// musl's `getentropy` is `getrandom` underneath on every system it runs on,
27+
// which is the source mbedTLS would have chosen had it known. The macro comes
28+
// from mcpp.toml's `cfg(c-abi = "musl")`, because musl defines none of its own.
29+
#ifdef TINYHTTPS_GETENTROPY
30+
#include <unistd.h>
31+
#endif
32+
1833
export module mcpplibs.tinyhttps:tls;
1934

2035
import :socket;
@@ -90,6 +105,19 @@ static int bio_recv(void* ctx, unsigned char* buf, size_t len) {
90105
return ret; // 0 means end of stream; mbedtls turns it into SSL_CONN_EOF
91106
}
92107

108+
#ifdef TINYHTTPS_GETENTROPY
109+
// `getentropy` gives at most 256 bytes a call.
110+
static int c_library_entropy(void*, unsigned char* out, size_t len) {
111+
while (len > 0) {
112+
const size_t n = len < 256 ? len : 256;
113+
if (::getentropy(out, n) != 0) return MBEDTLS_ERR_ENTROPY_SOURCE_FAILED;
114+
out += n;
115+
len -= n;
116+
}
117+
return 0;
118+
}
119+
#endif
120+
93121
static std::string mbedtls_message(int ret) {
94122
char buf[200];
95123
mbedtls_strerror(ret, buf, sizeof buf);
@@ -165,7 +193,18 @@ public:
165193
[[nodiscard]] const std::string& error() const { return error_; }
166194

167195
// Call before connect(); it covers the handshake and every later wait.
168-
void set_stop(std::stop_token stop) { socket_.set_stop(std::move(stop)); }
196+
//
197+
// THE TOKEN BELONGS TO THE SOCKET AT THE BOTTOM, AND ONLY THERE. Every wait
198+
// ends in a `Socket` — this session's own, or, inside an https:// proxy's
199+
// tunnel, the one beneath `lower_` — so that is where it is set. Setting
200+
// only `socket_` left a reused tunnel holding the token of the call that
201+
// opened it: a later call could not be cancelled, and once that first
202+
// token was stopped every wait failed at once and the stale-connection
203+
// retry sent a POST that had already arrived a second time.
204+
void set_stop(std::stop_token stop) {
205+
if (lower_) lower_->set_stop(stop);
206+
socket_.set_stop(std::move(stop));
207+
}
169208

170209
// Connect over an already-established Socket (e.g. a proxy tunnel).
171210
// Takes ownership of the socket and performs TLS handshake on top of it.
@@ -336,7 +375,11 @@ private:
336375
state_ = std::make_unique<TlsState>();
337376

338377
int ret = mbedtls_ctr_drbg_seed(
378+
#ifdef TINYHTTPS_GETENTROPY
379+
&state_->ctr_drbg, c_library_entropy, nullptr,
380+
#else
339381
&state_->ctr_drbg, mbedtls_entropy_func, &state_->entropy,
382+
#endif
340383
nullptr, 0);
341384
if (ret != 0) return fail(mbedtls_message(ret));
342385

0 commit comments

Comments
 (0)