mirror of
https://github.com/knutwurst/patchdl.git
synced 2026-10-06 05:00:22 +02:00
Apply the connection-count setting live, without a payload restart
The pool now spawns the full worker set at startup and gates each worker by its slot against a live active_conns limit, instead of spawning exactly max_connections threads once. Saving a new value in Settings updates the limit and broadcasts: idle workers wake to pull pieces, and a lowered limit parks the extra workers after they finish their current piece. No restart, and no thread creation/teardown at runtime. Verified on device: max_connections changed 4 -> 8 -> 16 -> 4 through the API while a download stayed active throughout. (Throughput did not scale with connections on this CDN, which caps aggregate bandwidth per source IP; 4 is a sensible default.)
This commit is contained in:
1 parent
21f6bd61ad
commit
6024dbfc5d
2 files changed
+49
-10
No files matched your search
+48
-9
@@ -105,7 +105,10 @@ static struct {
|
|||||||
dl_job_t *active; /* the one admitted job, or NULL */
|
dl_job_t *active; /* the one admitted job, or NULL */
|
||||||
int admitting; /* a worker is fetching/preallocating next job */
|
int admitting; /* a worker is fetching/preallocating next job */
|
||||||
int stopping;
|
int stopping;
|
||||||
int n_workers;
|
int n_workers; /* threads actually spawned (== POOL_MAX_CONN) */
|
||||||
|
int active_conns; /* live connection limit (1..n_workers); a worker
|
||||||
|
with slot >= active_conns idles. Changing this
|
||||||
|
+ broadcast resizes the pool with no restart. */
|
||||||
pthread_t workers[POOL_MAX_CONN];
|
pthread_t workers[POOL_MAX_CONN];
|
||||||
volatile long long inflight_bytes[POOL_MAX_CONN]; /* per-slot live counters */
|
volatile long long inflight_bytes[POOL_MAX_CONN]; /* per-slot live counters */
|
||||||
} g_pool = { .mtx = PTHREAD_MUTEX_INITIALIZER, .cv = PTHREAD_COND_INITIALIZER };
|
} g_pool = { .mtx = PTHREAD_MUTEX_INITIALIZER, .cv = PTHREAD_COND_INITIALIZER };
|
||||||
@@ -1274,6 +1277,14 @@ dl_worker(void *arg) {
|
|||||||
pthread_mutex_lock(&g_pool.mtx);
|
pthread_mutex_lock(&g_pool.mtx);
|
||||||
for (;;) {
|
for (;;) {
|
||||||
if (g_pool.stopping) break;
|
if (g_pool.stopping) break;
|
||||||
|
/* Live connection limit: a worker above the configured count idles
|
||||||
|
here (claims no new piece, doesn't admit). Lowering active_conns
|
||||||
|
lets in-flight pieces finish first, then those workers park here;
|
||||||
|
raising it wakes them to start pulling pieces — no restart. */
|
||||||
|
if (slot >= g_pool.active_conns) {
|
||||||
|
pthread_cond_wait(&g_pool.cv, &g_pool.mtx);
|
||||||
|
continue;
|
||||||
|
}
|
||||||
job = g_pool.active;
|
job = g_pool.active;
|
||||||
if (!job) {
|
if (!job) {
|
||||||
if (!g_pool.admitting) {
|
if (!g_pool.admitting) {
|
||||||
@@ -1677,6 +1688,7 @@ request_completed(void *cls, struct MHD_Connection *conn, void **con_cls,
|
|||||||
static enum MHD_Result
|
static enum MHD_Result
|
||||||
handle_config_post(struct MHD_Connection *conn, const char *body) {
|
handle_config_post(struct MHD_Connection *conn, const char *body) {
|
||||||
char pol[8];
|
char pol[8];
|
||||||
|
int mc;
|
||||||
|
|
||||||
json_get_str(body, "default_policy", pol, sizeof(pol));
|
json_get_str(body, "default_policy", pol, sizeof(pol));
|
||||||
|
|
||||||
@@ -1697,8 +1709,20 @@ handle_config_post(struct MHD_Connection *conn, const char *body) {
|
|||||||
json_get_int(body, "max_connections", g_cfg.max_connections);
|
json_get_int(body, "max_connections", g_cfg.max_connections);
|
||||||
if (g_cfg.max_connections < 1) g_cfg.max_connections = 1;
|
if (g_cfg.max_connections < 1) g_cfg.max_connections = 1;
|
||||||
if (g_cfg.max_connections > POOL_MAX_CONN) g_cfg.max_connections = POOL_MAX_CONN;
|
if (g_cfg.max_connections > POOL_MAX_CONN) g_cfg.max_connections = POOL_MAX_CONN;
|
||||||
|
mc = g_cfg.max_connections;
|
||||||
pthread_mutex_unlock(&g_mutex);
|
pthread_mutex_unlock(&g_mutex);
|
||||||
|
|
||||||
|
/* Apply the new connection count to the live pool — no restart. Raising it
|
||||||
|
wakes idle workers to pull pieces immediately; lowering it lets the extra
|
||||||
|
workers finish their current piece, then they park. (g_mutex released
|
||||||
|
first to keep the g_mutex-before-g_pool lock order.) */
|
||||||
|
pthread_mutex_lock(&g_pool.mtx);
|
||||||
|
g_pool.active_conns = mc;
|
||||||
|
if (g_pool.active_conns < 1) g_pool.active_conns = 1;
|
||||||
|
if (g_pool.active_conns > g_pool.n_workers) g_pool.active_conns = g_pool.n_workers;
|
||||||
|
pthread_cond_broadcast(&g_pool.cv);
|
||||||
|
pthread_mutex_unlock(&g_pool.mtx);
|
||||||
|
|
||||||
save_config();
|
save_config();
|
||||||
return queue_json_owned(conn, MHD_HTTP_OK, build_config_json());
|
return queue_json_owned(conn, MHD_HTTP_OK, build_config_json());
|
||||||
}
|
}
|
||||||
@@ -1857,16 +1881,31 @@ patchdl_websrv_start(unsigned short port) {
|
|||||||
download workers race their first curl_easy_init. */
|
download workers race their first curl_easy_init. */
|
||||||
patchdl_net_global_init();
|
patchdl_net_global_init();
|
||||||
|
|
||||||
/* Spawn the N-connection download pool BEFORE the web server accepts work. */
|
/* Spawn the full set of workers BEFORE the web server accepts work, then
|
||||||
g_pool.n_workers = g_cfg.max_connections;
|
cap how many are *active* with g_pool.active_conns. Spawning all of them
|
||||||
if (g_pool.n_workers < 1) g_pool.n_workers = 1;
|
up front lets the connection count be raised live (up to POOL_MAX_CONN)
|
||||||
if (g_pool.n_workers > POOL_MAX_CONN) g_pool.n_workers = POOL_MAX_CONN;
|
without creating threads at runtime; idle workers just wait on the cv.
|
||||||
|
active_conns/n_workers are set before the first pthread_create so every
|
||||||
|
worker reads a valid limit from its gate (no startup data race). */
|
||||||
|
g_pool.n_workers = POOL_MAX_CONN;
|
||||||
|
g_pool.active_conns = g_cfg.max_connections;
|
||||||
|
if (g_pool.active_conns < 1) g_pool.active_conns = 1;
|
||||||
|
if (g_pool.active_conns > POOL_MAX_CONN) g_pool.active_conns = POOL_MAX_CONN;
|
||||||
g_pool.stopping = 0;
|
g_pool.stopping = 0;
|
||||||
for (int s = 0; s < g_pool.n_workers; s++) {
|
{
|
||||||
|
int spawned = 0;
|
||||||
|
for (int s = 0; s < POOL_MAX_CONN; s++) {
|
||||||
if (pthread_create(&g_pool.workers[s], NULL, dl_worker,
|
if (pthread_create(&g_pool.workers[s], NULL, dl_worker,
|
||||||
(void *)(intptr_t)s)) {
|
(void *)(intptr_t)s))
|
||||||
g_pool.n_workers = s; /* only the threads that started exist */
|
break; /* only the threads that started exist */
|
||||||
break;
|
spawned++;
|
||||||
|
}
|
||||||
|
if (spawned < g_pool.n_workers) {
|
||||||
|
pthread_mutex_lock(&g_pool.mtx);
|
||||||
|
g_pool.n_workers = spawned;
|
||||||
|
if (g_pool.active_conns > g_pool.n_workers)
|
||||||
|
g_pool.active_conns = g_pool.n_workers;
|
||||||
|
pthread_mutex_unlock(&g_pool.mtx);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+1
-1
@@ -189,7 +189,7 @@
|
|||||||
<label class="field">
|
<label class="field">
|
||||||
<span>Parallel download connections</span>
|
<span>Parallel download connections</span>
|
||||||
<input id="maxConnections" type="number" min="1" max="16" step="1"
|
<input id="maxConnections" type="number" min="1" max="16" step="1"
|
||||||
title="Connections used to download a patch in parallel — applies after a restart" />
|
title="Connections used to download a patch in parallel — applies live on Save (raising is instant; lowering settles as in-flight pieces finish)" />
|
||||||
</label>
|
</label>
|
||||||
</div>
|
</div>
|
||||||
|
|
||||||
|
|||||||
Reference in new issue
Block a user