ratelimit: redesign

Description of how this works in `docs/internal/RATELIMITS.ms`.

Notable implementation changes:
- KEEP_SEND_PAUSE/KEEP_SEND_HOLD and KEEP_RECV_PAUSE/KEEP_RECV_HOLD
  no longer exist. Pausing is down via blocked the new rlimits.
- KEEP_SEND_TIMED no longer exists. Pausing "100-continue" transfers
  is done in the new `Curl_http_perform_pollset()` method.
- HTTP/2 rate limiting implemented via window updates. When
  transfer initiaiting connection has a ratelimit, adjust the
  initial window size
- HTTP/3 ngtcp2 rate limitin implemnented via ack updates
- HTTP/3 quiche does not seem to support this via its API
- the default progress-meter has been improved for accuracy
  in "current speed" results.

pytest speed tests have been improved.

Closes #19384
This commit is contained in:
Stefan Eissing 2025-11-11 14:26:48 +01:00 committed by Daniel Stenberg
parent bfde781121
commit 24b36fdd15
No known key found for this signature in database
GPG key ID: 5CC908FDB71E12C2
48 changed files with 1146 additions and 675 deletions

View file

@ -43,7 +43,6 @@
#include "select.h"
#include "curlx/warnless.h"
#include "curlx/wait.h"
#include "speedcheck.h"
#include "conncache.h"
#include "multihandle.h"
#include "sigpipe.h"
@ -923,79 +922,145 @@ void Curl_attach_connection(struct Curl_easy *data,
conn->handler->attach(data, conn);
}
/* adjust pollset for rate limits/pauses */
static CURLcode multi_adjust_pollset(struct Curl_easy *data,
struct easy_pollset *ps)
{
CURLcode result = CURLE_OK;
if(ps->n) {
struct curltime now = curlx_now();
bool send_blocked, recv_blocked;
recv_blocked = (Curl_rlimit_avail(&data->progress.dl.rlimit, now) <= 0);
send_blocked = (Curl_rlimit_avail(&data->progress.ul.rlimit, now) <= 0);
if(send_blocked || recv_blocked) {
int i;
for(i = 0; i <= SECONDARYSOCKET; ++i) {
curl_socket_t sock = data->conn->sock[i];
if(sock == CURL_SOCKET_BAD)
continue;
if(recv_blocked && Curl_pollset_want_recv(data, ps, sock)) {
result = Curl_pollset_remove_in(data, ps, sock);
if(result)
break;
}
if(send_blocked && Curl_pollset_want_send(data, ps, sock)) {
result = Curl_pollset_remove_out(data, ps, sock);
if(result)
break;
}
}
}
/* Not blocked and wanting to receive. If there is data pending
* in the connection filters, make transfer run again. */
if(!recv_blocked &&
((Curl_pollset_want_recv(data, ps, data->conn->sock[FIRSTSOCKET]) &&
Curl_conn_data_pending(data, FIRSTSOCKET)) ||
(Curl_pollset_want_recv(data, ps, data->conn->sock[SECONDARYSOCKET]) &&
Curl_conn_data_pending(data, SECONDARYSOCKET)))) {
CURL_TRC_M(data, "pollset[] has POLLIN, but there is still "
"buffered input -> mark as dirty");
Curl_multi_mark_dirty(data);
}
}
return result;
}
static CURLcode mstate_connecting_pollset(struct Curl_easy *data,
struct easy_pollset *ps)
{
if(data->conn) {
curl_socket_t sockfd = Curl_conn_get_first_socket(data);
if(sockfd != CURL_SOCKET_BAD) {
/* Default is to wait to something from the server */
return Curl_pollset_change(data, ps, sockfd, CURL_POLL_IN, 0);
}
struct connectdata *conn = data->conn;
curl_socket_t sockfd;
CURLcode result = CURLE_OK;
if(Curl_xfer_recv_is_paused(data))
return CURLE_OK;
/* If a socket is set, receiving is default. If the socket
* has not been determined yet (eyeballing), always ask the
* connection filters for what to monitor. */
sockfd = Curl_conn_get_first_socket(data);
if(sockfd != CURL_SOCKET_BAD) {
result = Curl_pollset_change(data, ps, sockfd, CURL_POLL_IN, 0);
if(!result)
result = multi_adjust_pollset(data, ps);
}
return CURLE_OK;
if(!result)
result = Curl_conn_adjust_pollset(data, conn, ps);
return result;
}
static CURLcode mstate_protocol_pollset(struct Curl_easy *data,
struct easy_pollset *ps)
{
struct connectdata *conn = data->conn;
if(conn) {
curl_socket_t sockfd;
if(conn->handler->proto_pollset)
return conn->handler->proto_pollset(data, ps);
sockfd = conn->sock[FIRSTSOCKET];
CURLcode result = CURLE_OK;
if(conn->handler->proto_pollset)
result = conn->handler->proto_pollset(data, ps);
else {
curl_socket_t sockfd = conn->sock[FIRSTSOCKET];
if(sockfd != CURL_SOCKET_BAD) {
/* Default is to wait to something from the server */
return Curl_pollset_change(data, ps, sockfd, CURL_POLL_IN, 0);
result = Curl_pollset_change(data, ps, sockfd, CURL_POLL_IN, 0);
}
}
return CURLE_OK;
if(!result)
result = multi_adjust_pollset(data, ps);
if(!result)
result = Curl_conn_adjust_pollset(data, conn, ps);
return result;
}
static CURLcode mstate_do_pollset(struct Curl_easy *data,
struct easy_pollset *ps)
{
struct connectdata *conn = data->conn;
if(conn) {
if(conn->handler->doing_pollset)
return conn->handler->doing_pollset(data, ps);
else if(CONN_SOCK_IDX_VALID(conn->send_idx)) {
/* Default is that we want to send something to the server */
return Curl_pollset_add_out(
data, ps, conn->sock[conn->send_idx]);
}
CURLcode result = CURLE_OK;
if(conn->handler->doing_pollset)
result = conn->handler->doing_pollset(data, ps);
else if(CONN_SOCK_IDX_VALID(conn->send_idx)) {
/* Default is that we want to send something to the server */
result = Curl_pollset_add_out(data, ps, conn->sock[conn->send_idx]);
}
return CURLE_OK;
if(!result)
result = multi_adjust_pollset(data, ps);
if(!result)
result = Curl_conn_adjust_pollset(data, conn, ps);
return result;
}
static CURLcode mstate_domore_pollset(struct Curl_easy *data,
struct easy_pollset *ps)
{
struct connectdata *conn = data->conn;
if(conn) {
if(conn->handler->domore_pollset)
return conn->handler->domore_pollset(data, ps);
else if(CONN_SOCK_IDX_VALID(conn->send_idx)) {
/* Default is that we want to send something to the server */
return Curl_pollset_add_out(
data, ps, conn->sock[conn->send_idx]);
}
CURLcode result = CURLE_OK;
if(conn->handler->domore_pollset)
result = conn->handler->domore_pollset(data, ps);
else if(CONN_SOCK_IDX_VALID(conn->send_idx)) {
/* Default is that we want to send something to the server */
result = Curl_pollset_add_out(data, ps, conn->sock[conn->send_idx]);
}
return CURLE_OK;
if(!result)
result = multi_adjust_pollset(data, ps);
if(!result)
result = Curl_conn_adjust_pollset(data, conn, ps);
return result;
}
static CURLcode mstate_perform_pollset(struct Curl_easy *data,
struct easy_pollset *ps)
{
struct connectdata *conn = data->conn;
if(!conn)
return CURLE_OK;
else if(conn->handler->perform_pollset)
return conn->handler->perform_pollset(data, ps);
CURLcode result = CURLE_OK;
if(conn->handler->perform_pollset)
result = conn->handler->perform_pollset(data, ps);
else {
/* Default is to obey the data->req.keepon flags for send/recv */
CURLcode result = CURLE_OK;
if(CURL_WANT_RECV(data) && CONN_SOCK_IDX_VALID(conn->recv_idx)) {
result = Curl_pollset_add_in(
data, ps, conn->sock[conn->recv_idx]);
@ -1006,19 +1071,21 @@ static CURLcode mstate_perform_pollset(struct Curl_easy *data,
result = Curl_pollset_add_out(
data, ps, conn->sock[conn->send_idx]);
}
return result;
}
if(!result)
result = multi_adjust_pollset(data, ps);
if(!result)
result = Curl_conn_adjust_pollset(data, conn, ps);
return result;
}
/* Initializes `poll_set` with the current socket poll actions needed
* for transfer `data`. */
CURLMcode Curl_multi_pollset(struct Curl_easy *data,
struct easy_pollset *ps,
const char *caller)
struct easy_pollset *ps)
{
CURLMcode mresult = CURLM_OK;
CURLcode result = CURLE_OK;
bool expect_sockets = TRUE;
/* If the transfer has no connection, this is fine. Happens when
called via curl_multi_remove_handle() => Curl_multi_ev_assess() =>
@ -1033,70 +1100,49 @@ CURLMcode Curl_multi_pollset(struct Curl_easy *data,
case MSTATE_SETUP:
case MSTATE_CONNECT:
/* nothing to poll for yet */
expect_sockets = FALSE;
break;
case MSTATE_RESOLVING:
result = Curl_resolv_pollset(data, ps);
/* connection filters are not involved in this phase. It is OK if we get no
* sockets to wait for. Resolving can wake up from other sources. */
expect_sockets = FALSE;
break;
case MSTATE_CONNECTING:
case MSTATE_TUNNELING:
if(!Curl_xfer_recv_is_paused(data)) {
result = mstate_connecting_pollset(data, ps);
if(!result)
result = Curl_conn_adjust_pollset(data, data->conn, ps);
}
else
expect_sockets = FALSE;
result = mstate_connecting_pollset(data, ps);
break;
case MSTATE_PROTOCONNECT:
case MSTATE_PROTOCONNECTING:
result = mstate_protocol_pollset(data, ps);
if(!result)
result = Curl_conn_adjust_pollset(data, data->conn, ps);
break;
case MSTATE_DO:
case MSTATE_DOING:
result = mstate_do_pollset(data, ps);
if(!result)
result = Curl_conn_adjust_pollset(data, data->conn, ps);
break;
case MSTATE_DOING_MORE:
result = mstate_domore_pollset(data, ps);
if(!result)
result = Curl_conn_adjust_pollset(data, data->conn, ps);
break;
case MSTATE_DID: /* same as PERFORMING in regard to polling */
case MSTATE_PERFORMING:
result = mstate_perform_pollset(data, ps);
if(!result)
result = Curl_conn_adjust_pollset(data, data->conn, ps);
break;
case MSTATE_RATELIMITING:
/* we need to let time pass, ignore socket(s) */
expect_sockets = FALSE;
break;
case MSTATE_DONE:
case MSTATE_COMPLETED:
case MSTATE_MSGSENT:
/* nothing more to poll for */
expect_sockets = FALSE;
break;
default:
failf(data, "multi_getsock: unexpected multi state %d", data->mstate);
DEBUGASSERT(0);
expect_sockets = FALSE;
break;
}
@ -1110,39 +1156,27 @@ CURLMcode Curl_multi_pollset(struct Curl_easy *data,
goto out;
}
/* Unblocked and waiting to receive with buffered input.
* Make transfer run again at next opportunity. */
if(!Curl_xfer_is_blocked(data) && !Curl_xfer_is_too_fast(data) &&
((Curl_pollset_want_read(data, ps, data->conn->sock[FIRSTSOCKET]) &&
Curl_conn_data_pending(data, FIRSTSOCKET)) ||
(Curl_pollset_want_read(data, ps, data->conn->sock[SECONDARYSOCKET]) &&
Curl_conn_data_pending(data, SECONDARYSOCKET)))) {
CURL_TRC_M(data, "%s pollset[] has POLLIN, but there is still "
"buffered input to consume -> mark as dirty", caller);
Curl_multi_mark_dirty(data);
}
#ifndef CURL_DISABLE_VERBOSE_STRINGS
if(CURL_TRC_M_is_verbose(data)) {
size_t timeout_count = Curl_llist_count(&data->state.timeoutlist);
switch(ps->n) {
case 0:
CURL_TRC_M(data, "%s pollset[], timeouts=%zu, paused %d/%d (r/w)",
caller, timeout_count,
CURL_TRC_M(data, "pollset[], timeouts=%zu, paused %d/%d (r/w)",
timeout_count,
Curl_xfer_send_is_paused(data),
Curl_xfer_recv_is_paused(data));
break;
case 1:
CURL_TRC_M(data, "%s pollset[fd=%" FMT_SOCKET_T " %s%s], timeouts=%zu",
caller, ps->sockets[0],
CURL_TRC_M(data, "pollset[fd=%" FMT_SOCKET_T " %s%s], timeouts=%zu",
ps->sockets[0],
(ps->actions[0] & CURL_POLL_IN) ? "IN" : "",
(ps->actions[0] & CURL_POLL_OUT) ? "OUT" : "",
timeout_count);
break;
case 2:
CURL_TRC_M(data, "%s pollset[fd=%" FMT_SOCKET_T " %s%s, "
CURL_TRC_M(data, "pollset[fd=%" FMT_SOCKET_T " %s%s, "
"fd=%" FMT_SOCKET_T " %s%s], timeouts=%zu",
caller, ps->sockets[0],
ps->sockets[0],
(ps->actions[0] & CURL_POLL_IN) ? "IN" : "",
(ps->actions[0] & CURL_POLL_OUT) ? "OUT" : "",
ps->sockets[1],
@ -1151,27 +1185,14 @@ CURLMcode Curl_multi_pollset(struct Curl_easy *data,
timeout_count);
break;
default:
CURL_TRC_M(data, "%s pollset[fds=%u], timeouts=%zu",
caller, ps->n, timeout_count);
CURL_TRC_M(data, "pollset[fds=%u], timeouts=%zu",
ps->n, timeout_count);
break;
}
CURL_TRC_EASY_TIMERS(data);
}
#endif
if(expect_sockets && !ps->n && data->multi &&
!Curl_uint_bset_contains(&data->multi->dirty, data->mid) &&
!Curl_llist_count(&data->state.timeoutlist) &&
!Curl_cwriter_is_paused(data) && !Curl_creader_is_paused(data) &&
Curl_conn_is_ip_connected(data, FIRSTSOCKET)) {
/* We expected sockets for POLL monitoring, but none are set.
* We are not dirty (and run anyway).
* We are not waiting on any timer.
* None of the READ/WRITE directions are paused.
* We are connected to the server on IP level, at least. */
infof(data, "WARNING: no socket in pollset or timer, transfer may stall!");
DEBUGASSERT(0);
}
out:
return mresult;
}
@ -1205,7 +1226,7 @@ CURLMcode curl_multi_fdset(CURLM *m,
continue;
}
Curl_multi_pollset(data, &ps, "curl_multi_fdset");
Curl_multi_pollset(data, &ps);
for(i = 0; i < ps.n; i++) {
if(!FDSET_SOCK(ps.sockets[i]))
/* pretend it does not exist */
@ -1268,7 +1289,7 @@ CURLMcode curl_multi_waitfds(CURLM *m,
Curl_uint_bset_remove(&multi->dirty, mid);
continue;
}
Curl_multi_pollset(data, &ps, "curl_multi_waitfds");
Curl_multi_pollset(data, &ps);
need += Curl_waitfds_add_ps(&cwfds, &ps);
}
while(Curl_uint_bset_next(&multi->process, mid, &mid));
@ -1354,7 +1375,7 @@ static CURLMcode multi_wait(struct Curl_multi *multi,
Curl_uint_bset_remove(&multi->dirty, mid);
continue;
}
Curl_multi_pollset(data, &ps, "multi_wait");
Curl_multi_pollset(data, &ps);
if(Curl_pollfds_add_ps(&cpfds, &ps)) {
result = CURLM_OUT_OF_MEMORY;
goto out;
@ -1907,35 +1928,28 @@ static CURLcode multi_follow(struct Curl_easy *data,
}
static CURLcode mspeed_check(struct Curl_easy *data,
struct curltime *nowp)
struct curltime now)
{
timediff_t recv_wait_ms = 0;
timediff_t send_wait_ms = 0;
/* check if over send speed */
if(data->set.max_send_speed)
send_wait_ms = Curl_pgrsLimitWaitTime(&data->progress.ul,
data->set.max_send_speed,
*nowp);
/* check if over recv speed */
if(data->set.max_recv_speed)
recv_wait_ms = Curl_pgrsLimitWaitTime(&data->progress.dl,
data->set.max_recv_speed,
*nowp);
/* check if our send/recv limits require idle waits */
send_wait_ms = Curl_rlimit_wait_ms(&data->progress.ul.rlimit, now);
recv_wait_ms = Curl_rlimit_wait_ms(&data->progress.dl.rlimit, now);
if(send_wait_ms || recv_wait_ms) {
if(data->mstate != MSTATE_RATELIMITING) {
Curl_ratelimit(data, *nowp);
multistate(data, MSTATE_RATELIMITING);
}
Curl_expire(data, CURLMAX(send_wait_ms, recv_wait_ms), EXPIRE_TOOFAST);
Curl_multi_clear_dirty(data);
CURL_TRC_M(data, "[RLIMIT] waiting %" FMT_TIMEDIFF_T "ms",
CURLMAX(send_wait_ms, recv_wait_ms));
return CURLE_AGAIN;
}
else if(data->mstate != MSTATE_PERFORMING) {
CURL_TRC_M(data, "[RLIMIT] wait over, continue");
multistate(data, MSTATE_PERFORMING);
Curl_ratelimit(data, *nowp);
}
return CURLE_OK;
}
@ -1951,7 +1965,7 @@ static CURLMcode state_performing(struct Curl_easy *data,
CURLcode result = *resultp = CURLE_OK;
*stream_errorp = FALSE;
if(mspeed_check(data, nowp) == CURLE_AGAIN)
if(mspeed_check(data, *nowp) == CURLE_AGAIN)
return CURLM_OK;
/* read/write data if it is ready to do so */
@ -2073,7 +2087,8 @@ static CURLMcode state_performing(struct Curl_easy *data,
}
}
else { /* not errored, not done */
mspeed_check(data, nowp);
*nowp = curlx_now();
mspeed_check(data, *nowp);
}
free(newurl);
*resultp = result;
@ -2228,10 +2243,7 @@ static CURLMcode state_ratelimiting(struct Curl_easy *data,
CURLMcode rc = CURLM_OK;
DEBUGASSERT(data->conn);
/* if both rates are within spec, resume transfer */
if(Curl_pgrsUpdate(data))
result = CURLE_ABORTED_BY_CALLBACK;
else
result = Curl_speedcheck(data, *nowp);
result = Curl_pgrsCheck(data);
if(result) {
if(!(data->conn->handler->flags & PROTOPT_DUAL) &&
@ -2242,7 +2254,7 @@ static CURLMcode state_ratelimiting(struct Curl_easy *data,
multi_done(data, result, TRUE);
}
else {
if(!mspeed_check(data, nowp))
if(!mspeed_check(data, *nowp))
rc = CURLM_CALL_MULTI_PERFORM;
}
*resultp = result;
@ -2387,6 +2399,8 @@ static CURLMcode multi_runsingle(struct Curl_multi *multi,
(HTTP/2), or the full connection for older protocols */
bool stream_error = FALSE;
rc = CURLM_OK;
/* update at start for continuous increase when looping */
*nowp = curlx_now();
if(multi_ischanged(multi, TRUE)) {
CURL_TRC_M(data, "multi changed, check CONNECT_PEND queue");
@ -2704,16 +2718,18 @@ statemachine_end:
rc = CURLM_CALL_MULTI_PERFORM;
}
/* if there is still a connection to use, call the progress function */
else if(data->conn && Curl_pgrsUpdate(data)) {
/* aborted due to progress callback return code must close the
connection */
result = CURLE_ABORTED_BY_CALLBACK;
streamclose(data->conn, "Aborted by callback");
else if(data->conn) {
result = Curl_pgrsUpdate(data);
if(result) {
/* aborted due to progress callback return code must close the
connection */
streamclose(data->conn, "Aborted by callback");
/* if not yet in DONE state, go there, otherwise COMPLETED */
multistate(data, (data->mstate < MSTATE_DONE) ?
MSTATE_DONE : MSTATE_COMPLETED);
rc = CURLM_CALL_MULTI_PERFORM;
/* if not yet in DONE state, go there, otherwise COMPLETED */
multistate(data, (data->mstate < MSTATE_DONE) ?
MSTATE_DONE : MSTATE_COMPLETED);
rc = CURLM_CALL_MULTI_PERFORM;
}
}
}