diff --git a/src/patchdl_net.c b/src/patchdl_net.c index 44548e6..37e4d7a 100644 --- a/src/patchdl_net.c +++ b/src/patchdl_net.c @@ -198,14 +198,14 @@ dns_lookup(const char *host, char *ip_out, size_t ip_sz) { return 0; } } - pthread_mutex_unlock(&dns_cache_mtx); - + /* Cold miss: resolve while HOLDING the cache lock (single-flight). N pool + workers needing the same CDN host would otherwise each blast Sony's + rate-limited resolver; this way one resolves and the rest get the cache. + It also serializes dns_resolve so its diagnostic globals can't be raced. + (This is the DNS lock, independent of the pool lock.) */ for (int attempt = 0; attempt < 4 && rc; attempt++) rc = dns_resolve(host, ip_out, ip_sz); - if (rc) return -1; - - pthread_mutex_lock(&dns_cache_mtx); - if (dns_cache_n < (int)(sizeof(dns_cache) / sizeof(dns_cache[0]))) { + if (!rc && dns_cache_n < (int)(sizeof(dns_cache) / sizeof(dns_cache[0]))) { strncpy(dns_cache[dns_cache_n].host, host, sizeof(dns_cache[0].host) - 1); strncpy(dns_cache[dns_cache_n].ip, ip_out, @@ -213,7 +213,7 @@ dns_lookup(const char *host, char *ip_out, size_t ip_sz) { dns_cache_n++; } pthread_mutex_unlock(&dns_cache_mtx); - return 0; + return rc; } /* ---------- HTTP GET via curl ------------------------------------------- */ @@ -678,6 +678,204 @@ patchdl_http_download_manifest(const char *manifest_url, const char *dest_path, bytes_out, NULL, NULL, 0, 0); } +/* ---- global init + parallel piece download (connection pool) ----------- */ + +void patchdl_net_global_init(void) { curl_global_init(CURL_GLOBAL_ALL); } +void patchdl_net_global_cleanup(void) { curl_global_cleanup(); } + +/* Write sink for one piece: pwrite at a fixed base offset (concurrent + non-overlapping pieces of the same fd are safe), tee into SHA-256 if asked, + and publish bytes-so-far for live progress. */ +typedef struct { + int fd; + long long base; + long long written; + EVP_MD_CTX *md; + volatile long long *bytes_slot; +} piece_sink_t; + +static size_t +piece_write_cb(void *ptr, size_t size, size_t nmemb, void *ud) { + piece_sink_t *s = (piece_sink_t *)ud; + size_t n = size * nmemb; + ssize_t w; + if (n == 0) return 0; + w = pwrite(s->fd, ptr, n, (off_t)(s->base + s->written)); + if (w < 0 || (size_t)w != n) return 0; /* short write -> curl errors out */ + if (s->md) EVP_DigestUpdate(s->md, ptr, n); + s->written += (long long)n; + if (s->bytes_slot) *s->bytes_slot = s->written; + return n; +} + +static int +piece_xfer_cb(void *clientp, curl_off_t dltotal, curl_off_t dlnow, + curl_off_t ultotal, curl_off_t ulnow) { + volatile int *abort_flag = (volatile int *)clientp; + (void)dltotal; (void)dlnow; (void)ultotal; (void)ulnow; + return (abort_flag && *abort_flag) ? 1 : 0; /* non-zero aborts the transfer */ +} + +int +patchdl_http_download_piece(const char *url, int fd, + long long file_offset, long long file_size, + const char *expected_sha256_or_null, + patchdl_piece_ctx_t *ctx) { + CURL *curl; + CURLcode res; + long http_code = 0; + char host[256], ip[INET_ADDRSTRLEN], rs443[512], rs80[512]; + struct curl_slist *rl = NULL; + struct curl_blob ca_blob; + piece_sink_t sink; + int verify = (expected_sha256_or_null && expected_sha256_or_null[0]); + + if (url_host(url, host, sizeof(host))) return -1; + if (!host_allowed(host)) return -1; + if (dns_lookup(host, ip, sizeof(ip))) return -1; + + memset(&sink, 0, sizeof(sink)); + sink.fd = fd; + sink.base = file_offset; + sink.bytes_slot = ctx ? ctx->bytes_slot : NULL; + if (verify) { + sink.md = EVP_MD_CTX_new(); + if (sink.md) EVP_DigestInit_ex(sink.md, EVP_sha256(), NULL); + } + + snprintf(rs443, sizeof(rs443), "%s:443:%s", host, ip); + rl = curl_slist_append(NULL, rs443); + snprintf(rs80, sizeof(rs80), "%s:80:%s", host, ip); + rl = curl_slist_append(rl, rs80); + + ca_blob.data = (void *)PATCHDL_SCEI_DNAS_ROOT_PEM; + ca_blob.len = strlen(PATCHDL_SCEI_DNAS_ROOT_PEM); + ca_blob.flags = CURL_BLOB_COPY; + + curl = curl_easy_init(); + if (!curl) { + curl_slist_free_all(rl); + if (sink.md) EVP_MD_CTX_free(sink.md); + return -1; + } + + curl_easy_setopt(curl, CURLOPT_URL, url); + curl_easy_setopt(curl, CURLOPT_RESOLVE, rl); + curl_easy_setopt(curl, CURLOPT_WRITEFUNCTION, piece_write_cb); + curl_easy_setopt(curl, CURLOPT_WRITEDATA, &sink); + curl_easy_setopt(curl, CURLOPT_CAINFO_BLOB, &ca_blob); + curl_easy_setopt(curl, CURLOPT_SSL_VERIFYPEER, 1L); + curl_easy_setopt(curl, CURLOPT_SSL_VERIFYHOST, 2L); + curl_easy_setopt(curl, CURLOPT_SSL_CIPHER_LIST, "DEFAULT@SECLEVEL=0"); + curl_easy_setopt(curl, CURLOPT_FOLLOWLOCATION, 1L); + curl_easy_setopt(curl, CURLOPT_MAXREDIRS, 5L); + curl_easy_setopt(curl, CURLOPT_FAILONERROR, 1L); /* 4xx/5xx -> error, no body written */ + curl_easy_setopt(curl, CURLOPT_CONNECTTIMEOUT, 20L); + curl_easy_setopt(curl, CURLOPT_LOW_SPEED_LIMIT, 1024L); + curl_easy_setopt(curl, CURLOPT_LOW_SPEED_TIME, 30L); + curl_easy_setopt(curl, CURLOPT_USERAGENT, "patchdl/1.0"); + if (ctx && ctx->abort) { + curl_easy_setopt(curl, CURLOPT_NOPROGRESS, 0L); + curl_easy_setopt(curl, CURLOPT_XFERINFOFUNCTION, piece_xfer_cb); + curl_easy_setopt(curl, CURLOPT_XFERINFODATA, (void *)ctx->abort); + } + + res = curl_easy_perform(curl); + curl_easy_getinfo(curl, CURLINFO_RESPONSE_CODE, &http_code); + curl_easy_cleanup(curl); + curl_slist_free_all(rl); + + if (res != CURLE_OK) { + if (sink.md) EVP_MD_CTX_free(sink.md); + return -1; /* network error / abort */ + } + if (file_size > 0 && sink.written != file_size) { + if (sink.md) EVP_MD_CTX_free(sink.md); + return -1; /* short or over-long -> failed */ + } + if (sink.md) { + unsigned char dig[EVP_MAX_MD_SIZE]; + unsigned int dl = 0; + char hex[2 * EVP_MAX_MD_SIZE + 1]; + EVP_DigestFinal_ex(sink.md, dig, &dl); + EVP_MD_CTX_free(sink.md); + hex_encode(dig, dl, hex, sizeof(hex)); + if (strcasecmp(hex, expected_sha256_or_null) != 0) + return -2; /* integrity mismatch */ + } + fdatasync(fd); /* durable before the caller sets the done bit */ + return 0; +} + +void +patchdl_manifest_free(patchdl_manifest_t *m) { + if (!m || !m->pieces) return; + for (int i = 0; i < m->count; i++) free(m->pieces[i].url); + free(m->pieces); + m->pieces = NULL; + m->count = 0; +} + +int +patchdl_fetch_manifest(const char *manifest_url, patchdl_manifest_t *out) { + patchdl_buf_t buf; + const char *pieces, *pieces_end, *p; + int cap = 0, n = 0; + long long running = 0; + + memset(out, 0, sizeof(*out)); + if (patchdl_http_get(manifest_url, &buf)) return -1; + if (!buf.data || !buf.size) { free(buf.data); return -1; } + + pieces = strstr(buf.data, "\"pieces\""); + if (!pieces || !(pieces = strchr(pieces, '['))) { free(buf.data); return -1; } + pieces_end = strchr(pieces, ']'); + + for (p = pieces; (p = strstr(p, "\"url\"")) && (!pieces_end || p < pieces_end); p += 5) + cap++; + if (cap <= 0) { free(buf.data); return -1; } + out->pieces = calloc((size_t)cap, sizeof(patchdl_piece_t)); + if (!out->pieces) { free(buf.data); return -1; } + + p = pieces; + while ((p = strstr(p, "\"url\"")) && (!pieces_end || p < pieces_end) && n < cap) { + char url[768] = {0}; + unsigned long long sz = 0, off = 0; + const char *obj_end = strchr(p, '}'); + + if (json_string_after(p, "url", url, sizeof(url))) + break; + json_u64_after(p, "fileSize", &sz); + if (json_u64_after(p, "fileOffset", &off) != 0) + off = (unsigned long long)running; /* no offset -> assume contiguous */ + + /* Validate tiling: pieces must be in order, contiguous, non-empty. */ + if ((long long)off != running || sz == 0) { + patchdl_manifest_free(out); + free(buf.data); + return -1; + } + out->pieces[n].url = strdup(url); + out->pieces[n].offset = (long long)off; + out->pieces[n].size = (long long)sz; + json_string_after(p, "hashValue", out->pieces[n].hash, + sizeof(out->pieces[n].hash)); + if (!out->pieces[n].url) { + patchdl_manifest_free(out); + free(buf.data); + return -1; + } + running += (long long)sz; + n++; + out->count = n; /* keep current so manifest_free frees exactly n */ + p = obj_end ? obj_end + 1 : p + 5; + } + free(buf.data); + if (n == 0) { patchdl_manifest_free(out); return -1; } + out->total = running; /* authoritative assembled size */ + return 0; +} + void patchdl_net_diag(const char *url, char *out_json, size_t sz) { char host[256] = {0}, ip[INET_ADDRSTRLEN] = {0}; diff --git a/src/patchdl_net.h b/src/patchdl_net.h index 94698fd..fda4f23 100644 --- a/src/patchdl_net.h +++ b/src/patchdl_net.h @@ -11,8 +11,52 @@ typedef struct { patchdl_buf_t *patchdl_buf_new(void); void patchdl_buf_free(patchdl_buf_t *b); +/* Call once, single-threaded, before any concurrent download worker starts / + after they have all joined. curl's global/OpenSSL init is otherwise lazy and + races across threads. */ +void patchdl_net_global_init(void); +void patchdl_net_global_cleanup(void); + int patchdl_http_get(const char *url, patchdl_buf_t *out); +/* ---- parallel piece download (used by the connection pool) ------------- */ + +/* One piece of a split manifest package. `url` is heap-allocated. */ +typedef struct { + char *url; + long long offset; /* byte offset of this piece in the assembled file */ + long long size; /* exact length of this piece */ + char hash[80]; /* manifest SHA-256 hex, or "" */ +} patchdl_piece_t; + +typedef struct { + patchdl_piece_t *pieces; + int count; + long long total; /* assembled file size = sum of piece sizes */ +} patchdl_manifest_t; + +/* Fetch + parse a Sony JSON manifest into a validated, contiguously-tiled + piece list. Returns 0 on success (caller frees with patchdl_manifest_free), + -1 on fetch/parse/tiling failure. */ +int patchdl_fetch_manifest(const char *manifest_url, patchdl_manifest_t *out); +void patchdl_manifest_free(patchdl_manifest_t *m); + +/* Live state shared with one in-flight piece download. The worker owns these; + the curl callbacks read `abort` (set elsewhere) and publish progress into + `bytes_slot` (single-writer per worker slot). */ +typedef struct { + volatile long long *bytes_slot; /* bytes written so far for this piece */ + volatile int *abort; /* non-zero -> stop this transfer */ +} patchdl_piece_ctx_t; + +/* Download one whole piece and pwrite it into `fd` at `file_offset`. Concurrent + non-overlapping pieces of the same fd are safe. Returns 0 on success (and + fdatasyncs fd), -1 on network/IO/abort, -2 on a SHA-256 mismatch. */ +int patchdl_http_download_piece(const char *url, int fd, + long long file_offset, long long file_size, + const char *expected_sha256_or_null, + patchdl_piece_ctx_t *ctx); + /* Progress callback. Return non-zero to ABORT the in-flight download (used to cancel large patch downloads); return 0 to continue. */ typedef int (*patchdl_download_progress_cb)(void *ctx, diff --git a/src/patchdl_websrv.c b/src/patchdl_websrv.c index 05aec8c..2222ec3 100644 --- a/src/patchdl_websrv.c +++ b/src/patchdl_websrv.c @@ -26,6 +26,11 @@ static patchdl_title_t *g_titles; static size_t g_title_count; static pthread_mutex_t g_mutex = PTHREAD_MUTEX_INITIALIZER; static char *g_debug_json; /* built once at startup */ +/* Background version.xml fetch thread. Joinable so shutdown can wait for it + before patchdl_scan_free frees g_titles out from under it. */ +static pthread_t g_verxml_tid; +static int g_verxml_started; +static volatile int g_verxml_stop; static unsigned long g_pkg_hits; static char g_pkg_diag_json[768] = "{\"hits\":0}"; @@ -33,25 +38,77 @@ static char g_pkg_diag_json[768] = "{\"hits\":0}"; title structs (t->enabled) so it travels with the scan; the global fields live here. Both are saved to PATCHDL_CFG_PATH (homebrew data dir, never a system file) and reloaded on the next start. Guarded by g_mutex. */ +#define POOL_MAX_CONN 16 +#define POOL_MAX_JOBS 16 + static struct { char default_policy[8]; /* "deny" | "allow" */ int install_after_download; int delete_pkg_after_install; int verify_downloads; /* SHA-256 each manifest piece (default off) */ int home_shortcut; /* user wants a home-screen browser tile */ -} g_cfg = { "deny", 0, 1, 0, 1 }; + int max_connections; /* parallel download connections (1..POOL_MAX_CONN) */ +} g_cfg = { "deny", 0, 1, 0, 1, 4 }; + +/* ---------- download connection pool ----------------------------------- */ +/* N worker threads share one job at a time (admit_max=1): a single download + spreads its pieces across all N connections; extra titles queue. Resume uses + a per-piece done-bitmap in the sidecar so an out-of-order parallel download + survives a reboot. All pool state is guarded by g_pool.mtx (NOT g_mutex). */ + +typedef enum { PC_PENDING = 0, PC_INFLIGHT, PC_DONE, PC_FAILED } pc_state_t; + +typedef struct { + pc_state_t state; + int attempts; + int slot; /* worker slot while INFLIGHT, else -1 */ +} dl_pstate_t; + +typedef enum { + JOB_QUEUED = 0, JOB_ADMITTING, JOB_ACTIVE, JOB_PAUSING, JOB_CANCELLING, + JOB_PAUSED, JOB_DONE, JOB_FAILED +} job_state_t; + +typedef struct dl_job { + char title_id[32]; + char name[128]; + char version[16]; + char manifest_url[768]; + char dir[288]; + char dest[320]; + int is_manifest; + int verify; + job_state_t state; + unsigned int seq; /* bumped on cancel; stale completions discarded */ + volatile int abort; /* read by the piece transfer callback */ + long long total; /* assembled size (preallocated) */ + long long done_bytes; /* sum of DONE piece sizes (committed) */ + long long speed_bps; + long long sample_bytes; + double sample_ts; + patchdl_manifest_t mf; /* pieces (url/offset/size/hash) */ + dl_pstate_t *ps; /* per-piece runtime state, mf.count entries */ + unsigned char *bitmap; /* (count+7)/8; bit set == DONE + durable */ + int pieces_done; + int pieces_failed; + int inflight; + int unpersisted; /* DONE pieces since last sidecar flush */ + int fd; /* O_RDWR dest fd, -1 until admit */ + int rc; /* 0 ok, -1 net/io, -2 verify */ + struct dl_job *next; +} dl_job_t; static struct { - int active; - int cancel; /* set by a cancel request; worker aborts + deletes */ - int pause; /* set by a pause request; worker aborts, keeps partial */ - char title_id[32]; - char name[128]; - char version[16]; - char path[320]; - long long downloaded; - long long total; -} g_dl; + pthread_mutex_t mtx; + pthread_cond_t cv; /* workers wait; broadcast on any change */ + dl_job_t *jobs; /* linked list, queue order */ + dl_job_t *active; /* the one admitted job, or NULL */ + int admitting; /* a worker is fetching/preallocating next job */ + int stopping; + int n_workers; + pthread_t workers[POOL_MAX_CONN]; + volatile long long inflight_bytes[POOL_MAX_CONN]; /* per-slot live counters */ +} g_pool = { .mtx = PTHREAD_MUTEX_INITIALIZER, .cv = PTHREAD_COND_INITIALIZER }; #define PATCHDL_DL_DIR "/data/patchdl" #define PATCHDL_CFG_PATH "/data/patchdl/config.json" @@ -89,6 +146,20 @@ json_get_bool(const char *s, const char *key, int dflt) { return dflt; } +static int +json_get_int(const char *s, const char *key, int dflt) { + char pat[64]; + snprintf(pat, sizeof(pat), "\"%s\"", key); + const char *p = strstr(s, pat); + if (!p) return dflt; + p = strchr(p + strlen(pat), ':'); + if (!p) return dflt; + p++; + while (*p == ' ' || *p == '\t' || *p == '\n' || *p == '\r') p++; + if (*p < '0' || *p > '9') return dflt; + return (int)strtol(p, NULL, 10); +} + /* Find `"key":"value"` and copy value into out. */ static void json_get_str(const char *s, const char *key, char *out, size_t sz) { @@ -433,7 +504,7 @@ build_status_json(void) { are appended from config_tail_json. Caller owns the result (queue_json_owned). */ static char * build_config_json(void) { - char head[192]; + char head[224]; pthread_mutex_lock(&g_mutex); snprintf(head, sizeof(head), @@ -441,12 +512,14 @@ build_config_json(void) { "\"install_after_download\":%s," "\"delete_pkg_after_install\":%s," "\"verify_downloads\":%s," - "\"home_shortcut\":%s,", + "\"home_shortcut\":%s," + "\"max_connections\":%d,", g_cfg.default_policy[0] ? g_cfg.default_policy : "deny", g_cfg.install_after_download ? "true" : "false", g_cfg.delete_pkg_after_install ? "true" : "false", g_cfg.verify_downloads ? "true" : "false", - g_cfg.home_shortcut ? "true" : "false"); + g_cfg.home_shortcut ? "true" : "false", + g_cfg.max_connections); pthread_mutex_unlock(&g_mutex); char *out = malloc(strlen(head) + sizeof(config_tail_json)); @@ -472,12 +545,14 @@ save_config(void) { "\"install_after_download\":%s,\n" "\"delete_pkg_after_install\":%s,\n" "\"verify_downloads\":%s,\n" - "\"home_shortcut\":%s,\n\"titles\":{", + "\"home_shortcut\":%s,\n" + "\"max_connections\":%d,\n\"titles\":{", g_cfg.default_policy[0] ? g_cfg.default_policy : "deny", g_cfg.install_after_download ? "true" : "false", g_cfg.delete_pkg_after_install ? "true" : "false", g_cfg.verify_downloads ? "true" : "false", - g_cfg.home_shortcut ? "true" : "false"); + g_cfg.home_shortcut ? "true" : "false", + g_cfg.max_connections); for (size_t i = 0; i < g_title_count; i++) fprintf(f, "%s\"%s\":%s", i ? "," : "", g_titles[i].title_id, g_titles[i].enabled ? "true" : "false"); @@ -515,6 +590,10 @@ load_config(void) { json_get_bool(buf, "verify_downloads", g_cfg.verify_downloads); g_cfg.home_shortcut = json_get_bool(buf, "home_shortcut", g_cfg.home_shortcut); + g_cfg.max_connections = + json_get_int(buf, "max_connections", g_cfg.max_connections); + 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; /* per-title overrides live under "titles": { "": true|false, ... } */ const char *titles = strstr(buf, "\"titles\""); @@ -605,42 +684,55 @@ build_titles_json(void) { return j.buf; /* caller owns; use queue_json_owned */ } +static const char * +job_state_str(job_state_t s) { + switch (s) { + case JOB_QUEUED: return "queued"; + case JOB_ADMITTING: return "active"; + case JOB_ACTIVE: return "active"; + case JOB_PAUSING: return "active"; + case JOB_CANCELLING: return "active"; + case JOB_PAUSED: return "paused"; + case JOB_DONE: return "done"; + default: return "error"; + } +} + +/* List every download in the pool (queued/active/paused/done/error) with live + bytes = committed done_bytes + in-flight piece bytes. */ static char * build_downloads_json(void) { jbuf_t j = {0}; - long long downloaded, total; - int progress = 0; - char detail[128]; - pthread_mutex_lock(&g_mutex); + pthread_mutex_lock(&g_pool.mtx); jbuf_append(&j, "["); - if (g_dl.active) { - downloaded = g_dl.downloaded; - total = g_dl.total; - if (total > 0 && downloaded >= 0) - progress = (int)((downloaded * 100) / total); - if (progress < 0) progress = 0; - if (progress > 100) progress = 100; - if (total > 0) - snprintf(detail, sizeof(detail), "%.1f GB of %.1f GB", - downloaded / 1073741824.0, total / 1073741824.0); - else - snprintf(detail, sizeof(detail), "%.1f MB downloaded", - downloaded / 1048576.0); - + int first = 1; + for (dl_job_t *job = g_pool.jobs; job; job = job->next) { + long long bytes = job->done_bytes, total = job->total; + int progress; + if (job == g_pool.active) + for (int s = 0; s < g_pool.n_workers; s++) + bytes += g_pool.inflight_bytes[s]; + if (total > 0) { + progress = (int)((bytes * 100) / total); + if (progress < 0) progress = 0; + if (progress > 100) progress = 100; + } else { + progress = (job->state == JOB_DONE) ? 100 : 0; + } + if (!first) jbuf_append(&j, ","); + first = 0; jbuf_append(&j, "{"); - jbuf_append(&j, "\"title_id\":"); jbuf_append_str(&j, g_dl.title_id); - jbuf_append(&j, ",\"name\":"); jbuf_append_str(&j, g_dl.name); - jbuf_append(&j, ",\"version\":"); jbuf_append_str(&j, g_dl.version); + jbuf_append(&j, "\"title_id\":"); jbuf_append_str(&j, job->title_id); + jbuf_append(&j, ",\"name\":"); jbuf_append_str(&j, job->name); + jbuf_append(&j, ",\"version\":"); jbuf_append_str(&j, job->version); + jbuf_append(&j, ",\"state\":"); jbuf_append_str(&j, job_state_str(job->state)); jbuf_appendf(&j, ",\"progress\":%d", progress); - jbuf_append(&j, ",\"detail\":"); jbuf_append_str(&j, detail); - jbuf_appendf(&j, ",\"bytes\":%lld,\"total_bytes\":%lld", - downloaded, total); - jbuf_append(&j, ",\"path\":"); jbuf_append_str(&j, g_dl.path); + jbuf_appendf(&j, ",\"bytes\":%lld,\"total_bytes\":%lld", bytes, total); jbuf_append(&j, "}"); } jbuf_append(&j, "]"); - pthread_mutex_unlock(&g_mutex); + pthread_mutex_unlock(&g_pool.mtx); return j.buf; } @@ -656,6 +748,15 @@ cleanup_installed_download(const char *title_id, const char *patch_url) { const char *base = strrchr(patch_url, '/'); if (!path_segment_safe(title_id)) return; + /* Never delete a directory the download pool still owns (a queued/active/ + paused job for this title would have its fd/sidecar pulled out). */ + pthread_mutex_lock(&g_pool.mtx); + for (dl_job_t *j = g_pool.jobs; j; j = j->next) + if (!strcmp(j->title_id, title_id)) { + pthread_mutex_unlock(&g_pool.mtx); + return; + } + pthread_mutex_unlock(&g_pool.mtx); base = base ? base + 1 : "patch.pkg"; snprintf(dir, sizeof(dir), "/data/patchdl/%s", title_id); snprintf(path, sizeof(path), "%s/%s", dir, base); @@ -672,6 +773,7 @@ verxml_fetch_thread(void *arg) { char url[512] = {0}; patchdl_verinfo_t info; + if (g_verxml_stop) break; /* shutdown: stop before the next blocking query */ memset(&info, 0, sizeof(info)); /* The version.xml URL comes straight from the app.db (UUID included); @@ -789,35 +891,28 @@ url_is_manifest(const char *url) { return n > 5 && !strcmp(url + n - 5, ".json"); } -/* Returns non-zero to abort the download when a cancel OR pause is requested. */ -static int -download_progress_cb(void *ctx, long long downloaded, long long total) { - int abort_now; - (void)ctx; - pthread_mutex_lock(&g_mutex); - if (g_dl.active) { - g_dl.downloaded = downloaded; - g_dl.total = total; - } - abort_now = g_dl.cancel || g_dl.pause; - pthread_mutex_unlock(&g_mutex); - return abort_now; +/* ---------- resume sidecar + pool helpers ------------------------------- */ + +static void +title_state_path(const char *title_id, char *out, size_t sz) { + snprintf(out, sz, "%s/%s/state.json", PATCHDL_DL_DIR, title_id); } -/* Delete a title's internal download directory and its contents (a partial or - a finished package). Best-effort; only ever touches /data/patchdl. */ +static void +remove_dl_state(const char *title_id) { + char path[320]; + title_state_path(title_id, path, sizeof(path)); + unlink(path); +} + +/* Delete a title's internal download directory and its contents. Best-effort; + only ever touches /data/patchdl (self-protecting against a bad segment). */ static void remove_title_dir(const char *title_id) { char dir[288], path[560]; DIR *d; struct dirent *e; - - /* Self-protecting: never build a delete path from an unsafe segment, even - if a future caller forgets the upstream check. A title_id of ".." would - otherwise resolve dir to /data and wipe it. */ - if (!path_segment_safe(title_id)) - return; - + if (!path_segment_safe(title_id)) return; snprintf(dir, sizeof(dir), "%s/%s", PATCHDL_DL_DIR, title_id); if ((d = opendir(dir))) { while ((e = readdir(d))) { @@ -831,103 +926,8 @@ remove_title_dir(const char *title_id) { rmdir(dir); } -static void set_title_resumable(const char *title_id, int resumable, long long bytes); - -/* Cancel an in-progress download for this title (the worker aborts and removes - the partial), or delete an already-downloaded package if nothing is running. */ -static enum MHD_Result -do_cancel(struct MHD_Connection *conn, const char *title_id) { - char resp[96]; - int was_active = 0; - - pthread_mutex_lock(&g_mutex); - if (g_dl.active && !strcmp(g_dl.title_id, title_id)) { - g_dl.cancel = 1; - was_active = 1; - } - pthread_mutex_unlock(&g_mutex); - - if (!was_active) { - remove_title_dir(title_id); - set_title_resumable(title_id, 0, 0); - } - - snprintf(resp, sizeof(resp), - "{\"ok\":true,\"cancelled\":%s,\"deleted\":true}", - was_active ? "true" : "false"); - return queue_json_owned(conn, MHD_HTTP_OK, strdup(resp)); -} - -/* Pause an in-progress download: the worker aborts but the partial is kept on - disk (and marked resumable), so it can be continued later. */ -static enum MHD_Result -do_pause(struct MHD_Connection *conn, const char *title_id) { - int active = 0; - pthread_mutex_lock(&g_mutex); - if (g_dl.active && !strcmp(g_dl.title_id, title_id)) { - g_dl.pause = 1; - active = 1; - } - pthread_mutex_unlock(&g_mutex); - if (!active) - return queue_json(conn, MHD_HTTP_CONFLICT, - "{\"ok\":false,\"reason\":\"not_downloading\"}"); - return queue_json(conn, MHD_HTTP_OK, "{\"ok\":true,\"paused\":true}"); -} - -/* ---------- resume sidecar (/data/patchdl//state.json) ----------- */ - -static long long -file_size(const char *path) { - struct stat st; - if (stat(path, &st) == 0 && S_ISREG(st.st_mode)) return (long long)st.st_size; - return -1; -} - -static void -title_state_path(const char *title_id, char *out, size_t sz) { - snprintf(out, sz, "%s/%s/state.json", PATCHDL_DL_DIR, title_id); -} - -/* Record which manifest a partial download belongs to, so a resume only - continues a partial that matches the current patch. Sony CDN URLs contain no - characters that need JSON escaping, so the value is embedded verbatim. */ -static void -write_dl_state(const char *title_id, const char *manifest_url) { - char path[320]; - FILE *f; - title_state_path(title_id, path, sizeof(path)); - f = fopen(path, "w"); - if (!f) return; - fprintf(f, "{\"manifest_url\":\"%s\"}\n", manifest_url); - fclose(f); -} - -static int -read_dl_state_url(const char *title_id, char *out, size_t sz) { - char path[320], buf[1280]; - FILE *f; - size_t n; - out[0] = '\0'; - title_state_path(title_id, path, sizeof(path)); - f = fopen(path, "r"); - if (!f) return -1; - n = fread(buf, 1, sizeof(buf) - 1, f); - fclose(f); - buf[n] = '\0'; - json_get_str(buf, "manifest_url", out, sz); - return out[0] ? 0 : -1; -} - -static void -remove_dl_state(const char *title_id) { - char path[320]; - title_state_path(title_id, path, sizeof(path)); - unlink(path); -} - -/* Keep the resumable flag (set at startup) in sync after a download finishes or - its partial is deleted, so /api/titles stops reporting a stale partial. */ +/* Mirror the resumable flag into g_titles (read by /api/titles). Takes g_mutex, + so it must be called with the POOL lock NOT held (lock order: g_mutex first). */ static void set_title_resumable(const char *title_id, int resumable, long long bytes) { pthread_mutex_lock(&g_mutex); @@ -940,131 +940,587 @@ set_title_resumable(const char *title_id, int resumable, long long bytes) { pthread_mutex_unlock(&g_mutex); } +static long long +json_get_ll(const char *s, const char *key, long long dflt) { + char pat[64]; + snprintf(pat, sizeof(pat), "\"%s\"", key); + const char *p = strstr(s, pat); + if (!p) return dflt; + p = strchr(p + strlen(pat), ':'); + if (!p) return dflt; + p++; + while (*p == ' ' || *p == '\t' || *p == '\n' || *p == '\r') p++; + if (*p < '0' || *p > '9') return dflt; + return strtoll(p, NULL, 10); +} + +static void +bitmap_to_hex(const unsigned char *bm, int nbytes, char *out, size_t sz) { + static const char hx[] = "0123456789abcdef"; + int i = 0; + for (; i < nbytes && (size_t)(2 * i + 2) < sz; i++) { + out[2 * i] = hx[(bm[i] >> 4) & 0xf]; + out[2 * i + 1] = hx[bm[i] & 0xf]; + } + out[2 * i] = '\0'; +} + +static void +hex_to_bitmap(const char *hex, unsigned char *bm, int nbytes) { + for (int i = 0; i < nbytes; i++) { + char a = hex[2 * i], b = a ? hex[2 * i + 1] : 0; + int hi, lo; + if (!a || !b) break; + hi = (a <= '9') ? a - '0' : (a | 32) - 'a' + 10; + lo = (b <= '9') ? b - '0' : (b | 32) - 'a' + 10; + bm[i] = (unsigned char)((hi << 4) | lo); + } +} + +/* Write the job's resume sidecar atomically. Each set bit's bytes were already + fdatasync'd in patchdl_http_download_piece, so a set bit implies durable data. + Called with the pool lock held (the write is tiny). */ +static void +write_job_state(const dl_job_t *job) { + char path[320], tmp[340], *hex; + int nbytes = (job->mf.count + 7) / 8; + FILE *f; + if (!job->bitmap) return; + title_state_path(job->title_id, path, sizeof(path)); + snprintf(tmp, sizeof(tmp), "%s.tmp", path); + hex = malloc((size_t)nbytes * 2 + 1); + if (!hex) return; + bitmap_to_hex(job->bitmap, nbytes, hex, (size_t)nbytes * 2 + 1); + f = fopen(tmp, "w"); + if (f) { + fprintf(f, + "{\"manifest_url\":\"%s\",\"total\":%lld,\"piece_count\":%d," + "\"done_bytes\":%lld,\"done\":\"%s\"}\n", + job->manifest_url, job->total, job->mf.count, job->done_bytes, hex); + fclose(f); + rename(tmp, path); + } + free(hex); +} + +/* Seed a freshly-admitted job's piece states from its resume sidecar — only if + it belongs to THIS manifest (same url + size + piece count). Sets DONE bits, + ps[] states and done_bytes. Called with the pool lock held. */ +static void +seed_from_sidecar(dl_job_t *job) { + char path[320], buf[16384], murl[768]; + FILE *f; + size_t n; + int nbytes = (job->mf.count + 7) / 8; + + title_state_path(job->title_id, path, sizeof(path)); + f = fopen(path, "r"); + if (!f) return; + n = fread(buf, 1, sizeof(buf) - 1, f); + fclose(f); + buf[n] = '\0'; + + json_get_str(buf, "manifest_url", murl, sizeof(murl)); + if (strcmp(murl, job->manifest_url) || + json_get_ll(buf, "total", -1) != job->total || + json_get_int(buf, "piece_count", -1) != job->mf.count) + return; /* stale / different patch -> start fresh */ + + { + size_t hexsz = (size_t)nbytes * 2 + 1; + char *hex = malloc(hexsz); + if (!hex) return; + json_get_str(buf, "done", hex, hexsz); + hex_to_bitmap(hex, job->bitmap, nbytes); + free(hex); + } + for (int i = 0; i < job->mf.count; i++) + if (job->bitmap[i / 8] & (1 << (i % 8))) { + job->ps[i].state = PC_DONE; + job->done_bytes += job->mf.pieces[i].size; + } +} + +/* ---------- pool scheduler + workers ------------------------------------ */ + +static int +pick_pending_piece(const dl_job_t *job) { + /* pieces are stored in ascending offset order -> first PENDING is lowest */ + for (int i = 0; i < job->mf.count; i++) + if (job->ps[i].state == PC_PENDING) return i; + return -1; +} + +/* Unlink a job from the list and free it. Pool lock held. */ +static void +free_job_locked(dl_job_t *job) { + dl_job_t **pp = &g_pool.jobs; + while (*pp && *pp != job) pp = &(*pp)->next; + if (*pp) *pp = job->next; + if (job->fd >= 0) close(job->fd); + patchdl_manifest_free(&job->mf); + free(job->ps); + free(job->bitmap); + free(job); +} + +/* Cancel finalize: delete the file + sidecar and drop the job. Pool lock held + on entry/exit; dropped briefly for set_title_resumable (g_mutex). */ +static void +finalize_cancel_locked(dl_job_t *job) { + char title_id[32], dir[288], dest[320]; + snprintf(title_id, sizeof title_id, "%s", job->title_id); + snprintf(dir, sizeof dir, "%s", job->dir); + snprintf(dest, sizeof dest, "%s", job->dest); + if (g_pool.active == job) g_pool.active = NULL; + free_job_locked(job); /* closes fd, frees job */ + unlink(dest); + remove_dl_state(title_id); + rmdir(dir); + pthread_mutex_unlock(&g_pool.mtx); + set_title_resumable(title_id, 0, 0); + pthread_mutex_lock(&g_pool.mtx); +} + +/* Pause finalize: keep the partial, persist the bitmap, mark resumable. The job + stays in the list as JOB_PAUSED. Pool lock held on entry/exit. */ +static void +finalize_pause_locked(dl_job_t *job) { + char title_id[32]; + long long done; + snprintf(title_id, sizeof title_id, "%s", job->title_id); + job->abort = 0; + write_job_state(job); + if (job->fd >= 0) { close(job->fd); job->fd = -1; } + job->state = JOB_PAUSED; + if (g_pool.active == job) g_pool.active = NULL; + done = job->done_bytes; + pthread_mutex_unlock(&g_pool.mtx); + set_title_resumable(title_id, done > 0, done); + pthread_mutex_lock(&g_pool.mtx); +} + +/* Settle an active job whose pieces are all done or some failed. Pool lock held + on entry/exit. */ +static void +settle_active_locked(dl_job_t *job) { + char title_id[32], dir[288], dest[320]; + snprintf(title_id, sizeof title_id, "%s", job->title_id); + snprintf(dir, sizeof dir, "%s", job->dir); + snprintf(dest, sizeof dest, "%s", job->dest); + if (g_pool.active == job) g_pool.active = NULL; + if (job->fd >= 0) { close(job->fd); job->fd = -1; } + + if (job->pieces_failed == 0 && job->pieces_done == job->mf.count) { + /* complete: keep the .pkg, drop the resume sidecar */ + job->state = JOB_DONE; + remove_dl_state(title_id); + pthread_mutex_unlock(&g_pool.mtx); + set_title_resumable(title_id, 0, 0); + pthread_mutex_lock(&g_pool.mtx); + } else if (job->rc == -2) { + /* corrupt (failed SHA-256): delete, not resumable */ + job->state = JOB_FAILED; + write_job_state(job); /* harmless; overwritten by delete below */ + unlink(dest); + remove_dl_state(title_id); + rmdir(dir); + pthread_mutex_unlock(&g_pool.mtx); + set_title_resumable(title_id, 0, 0); + pthread_mutex_lock(&g_pool.mtx); + } else { + /* network failure on a piece after retries: keep partial, resumable */ + long long done = job->done_bytes; + job->state = JOB_FAILED; + write_job_state(job); + pthread_mutex_unlock(&g_pool.mtx); + set_title_resumable(title_id, done > 0, done); + pthread_mutex_lock(&g_pool.mtx); + } +} + +/* Promote the head QUEUED job to ACTIVE: fetch+parse manifest, preallocate the + file, seed the resume bitmap. Network/disk I/O runs with the pool lock + DROPPED. Pool lock held on entry/exit. */ +static void +admit_next(void) { + dl_job_t *q = NULL; + patchdl_manifest_t mf; + char murl[768], dest[320], dir[288]; + int is_manifest, fd = -1, ok = 0; + long long total = 0, old_size = 0; + dl_pstate_t *ps = NULL; + unsigned char *bitmap = NULL; + + for (dl_job_t *j = g_pool.jobs; j; j = j->next) + if (j->state == JOB_QUEUED) { q = j; break; } + if (!q) return; + + g_pool.admitting = 1; + q->state = JOB_ADMITTING; /* NOT claimable; a cancel/pause mid-admit only + flags it — admit_next is the sole finalizer. + g_pool.active stays NULL until fully built. */ + is_manifest = q->is_manifest; + snprintf(murl, sizeof murl, "%s", q->manifest_url); + snprintf(dest, sizeof dest, "%s", q->dest); + snprintf(dir, sizeof dir, "%s", q->dir); + pthread_mutex_unlock(&g_pool.mtx); + + /* ---- network + disk, NO pool lock ---- */ + memset(&mf, 0, sizeof mf); + mkdir(PATCHDL_DL_DIR, 0777); + mkdir(dir, 0777); + if (is_manifest) { + if (patchdl_fetch_manifest(murl, &mf) == 0 && mf.count > 0) { + total = mf.total; + fd = open(dest, O_RDWR | O_CREAT, 0666); + if (fd >= 0) { + struct stat st; + if (fstat(fd, &st) == 0) old_size = (long long)st.st_size; + if (ftruncate(fd, (off_t)total) == 0) ok = 1; + else { close(fd); fd = -1; } + } + } + } else { + mf.pieces = calloc(1, sizeof(patchdl_piece_t)); + if (mf.pieces) { + mf.pieces[0].url = strdup(murl); + mf.pieces[0].offset = 0; + mf.pieces[0].size = 0; /* unknown -> size check skipped */ + mf.count = 1; mf.total = 0; total = 0; + fd = open(dest, O_RDWR | O_CREAT | O_TRUNC, 0666); + if (fd >= 0 && mf.pieces[0].url) ok = 1; + else if (fd >= 0) { close(fd); fd = -1; } + } + } + if (ok) { + ps = calloc((size_t)mf.count, sizeof(dl_pstate_t)); + bitmap = calloc((size_t)((mf.count + 7) / 8), 1); + if (!ps || !bitmap) ok = 0; + } + + pthread_mutex_lock(&g_pool.mtx); + g_pool.admitting = 0; + /* During the unlocked I/O window do_cancel/do_pause may have flagged q + (it stays JOB_ADMITTING otherwise). The job was never published as + g_pool.active, so no worker touched it and q is still valid here. */ + if (!ok || q->state == JOB_CANCELLING) { + if (fd >= 0) close(fd); + patchdl_manifest_free(&mf); + free(ps); + free(bitmap); + if (q->state == JOB_CANCELLING) finalize_cancel_locked(q); + else { q->state = JOB_FAILED; q->rc = -1; } + pthread_cond_broadcast(&g_pool.cv); + return; + } + /* Resume of a previously paused job reuses the same dl_job_t, which still + carries the prior mf/ps/bitmap and a committed done_bytes. Free the stale + buffers and zero the committed counters before attaching the fresh ones so + the sidecar bitmap is the single source of truth (no leak, no double-count). + Harmless on a fresh job: mf is zeroed, ps/bitmap are NULL, free(NULL) is ok. */ + patchdl_manifest_free(&q->mf); + free(q->ps); q->ps = NULL; + free(q->bitmap); q->bitmap = NULL; + q->done_bytes = 0; q->pieces_failed = 0; q->unpersisted = 0; q->rc = 0; + + q->mf = mf; q->ps = ps; q->bitmap = bitmap; q->fd = fd; q->total = total; + for (int i = 0; i < mf.count; i++) q->ps[i].slot = -1; + seed_from_sidecar(q); + q->pieces_done = 0; + for (int i = 0; i < mf.count; i++) + if (q->ps[i].state == PC_DONE) q->pieces_done++; + /* One-time migration: an old sequential partial (written in order, no + per-piece bitmap) is a contiguous prefix on disk. If the bitmap seed found + nothing and the file is a partial (smaller than the full preallocated + size), mark every piece fully within the on-disk bytes as done so we don't + re-download what's already there. The partial frontier piece is excluded. */ + if (q->pieces_done == 0 && old_size > 0 && old_size < total) { + for (int i = 0; i < mf.count; i++) + if (mf.pieces[i].offset + mf.pieces[i].size <= old_size) { + q->ps[i].state = PC_DONE; + q->bitmap[i / 8] |= (unsigned char)(1 << (i % 8)); + q->done_bytes += mf.pieces[i].size; + q->pieces_done++; + } + if (q->pieces_done > 0) write_job_state(q); /* persist the migrated bitmap */ + } + /* A pause that landed during admit: persist the seeded/migrated bitmap and + park as PAUSED (keeping the on-disk partial) rather than starting the + transfer. finalize_pause_locked writes the sidecar and closes the fd. */ + if (q->state == JOB_PAUSING) { + finalize_pause_locked(q); + pthread_cond_broadcast(&g_pool.cv); + return; + } + q->state = JOB_ACTIVE; + g_pool.active = q; + pthread_cond_broadcast(&g_pool.cv); +} + +static void * +dl_worker(void *arg) { + int slot = (int)(intptr_t)arg; + + for (;;) { + dl_job_t *job; + int pidx = -1; + char url[768], hash[80]; + long long off = 0, sz = 0; + unsigned int my_seq = 0; + int fd = -1, verify = 0, rc; + volatile int *abort_ptr = NULL; + + pthread_mutex_lock(&g_pool.mtx); + for (;;) { + if (g_pool.stopping) break; + job = g_pool.active; + if (!job) { + if (!g_pool.admitting) { + int queued = 0; + for (dl_job_t *j = g_pool.jobs; j; j = j->next) + if (j->state == JOB_QUEUED) { queued = 1; break; } + if (queued) { admit_next(); continue; } + } + pthread_cond_wait(&g_pool.cv, &g_pool.mtx); + continue; + } + if (job->state == JOB_ACTIVE) { + pidx = pick_pending_piece(job); + if (pidx >= 0) break; /* claim below */ + if (job->inflight == 0) { settle_active_locked(job); pthread_cond_broadcast(&g_pool.cv); continue; } + pthread_cond_wait(&g_pool.cv, &g_pool.mtx); + continue; + } + if (job->state == JOB_PAUSING || job->state == JOB_CANCELLING) { + if (job->inflight == 0) { + if (job->state == JOB_PAUSING) finalize_pause_locked(job); + else finalize_cancel_locked(job); + pthread_cond_broadcast(&g_pool.cv); + continue; + } + pthread_cond_wait(&g_pool.cv, &g_pool.mtx); + continue; + } + pthread_cond_wait(&g_pool.cv, &g_pool.mtx); + } + if (g_pool.stopping) { pthread_mutex_unlock(&g_pool.mtx); break; } + + /* claim piece pidx */ + job->ps[pidx].state = PC_INFLIGHT; + job->ps[pidx].slot = slot; + job->inflight++; + my_seq = job->seq; + fd = job->fd; + verify = job->verify; + off = job->mf.pieces[pidx].offset; + sz = job->mf.pieces[pidx].size; + abort_ptr = &job->abort; + snprintf(url, sizeof url, "%s", job->mf.pieces[pidx].url); + snprintf(hash, sizeof hash, "%s", job->mf.pieces[pidx].hash); + g_pool.inflight_bytes[slot] = 0; + pthread_mutex_unlock(&g_pool.mtx); + + /* ---- download the piece, NO lock ---- */ + { + patchdl_piece_ctx_t ctx = { &g_pool.inflight_bytes[slot], abort_ptr }; + rc = patchdl_http_download_piece(url, fd, off, sz, + verify && hash[0] ? hash : NULL, &ctx); + } + + pthread_mutex_lock(&g_pool.mtx); + g_pool.inflight_bytes[slot] = 0; + job->inflight--; + job->ps[pidx].slot = -1; + if (my_seq != job->seq) { + /* job cancelled/torn down under us: discard result (do not touch + bitmap/counters); the finalizer runs once inflight hits 0 */ + } else if (rc == 0) { + job->ps[pidx].state = PC_DONE; + job->pieces_done++; + job->done_bytes += sz; + job->bitmap[pidx / 8] |= (unsigned char)(1 << (pidx % 8)); + if (++job->unpersisted >= 8) { write_job_state(job); job->unpersisted = 0; } + } else if (job->abort && + (job->state == JOB_PAUSING || job->state == JOB_CANCELLING)) { + job->ps[pidx].state = PC_PENDING; /* aborted by pause -> redo on resume */ + } else if (rc == -2) { + job->ps[pidx].state = PC_FAILED; job->pieces_failed++; job->rc = -2; + } else { + job->ps[pidx].attempts++; + if (job->ps[pidx].attempts < 4) job->ps[pidx].state = PC_PENDING; + else { job->ps[pidx].state = PC_FAILED; job->pieces_failed++; job->rc = -1; } + } + pthread_cond_broadcast(&g_pool.cv); + pthread_mutex_unlock(&g_pool.mtx); + } + return NULL; +} + +/* ---------- HTTP handlers (enqueue/pause/cancel) ------------------------ */ + +/* Enqueue a download. Returns 202 immediately; the pool admits + downloads it + across N connections. De-dupes by title_id; resumes a paused job. */ static enum MHD_Result do_download(struct MHD_Connection *conn, const char *title_id, patchdl_source_t src, const char *patch_url, const char *name, const char *version, int enabled) { - char dir[256], dest[320], resp[640]; - long long bytes = 0; - int verify = 0, dlrc, is_manifest, resume = 0; + dl_job_t *job, *existing = NULL; + int count = 0, verify; if (!enabled) return queue_json(conn, MHD_HTTP_FORBIDDEN, "{\"ok\":false,\"reason\":\"title_disabled\"}"); - if (src == PATCHDL_SOURCE_UNKNOWN) return queue_json(conn, MHD_HTTP_FORBIDDEN, "{\"ok\":false,\"reason\":\"source_unknown\"}"); - if (!patch_url[0]) return queue_json(conn, MHD_HTTP_CONFLICT, "{\"ok\":false,\"reason\":\"no_compatible_patch\"}"); - /* /data/patchdl/<title_id>/<pkg-basename> — homebrew data dir, not system */ - mkdir(PATCHDL_DL_DIR, 0777); - snprintf(dir, sizeof(dir), "%s/%s", PATCHDL_DL_DIR, title_id); - mkdir(dir, 0777); - title_pkg_path(title_id, patch_url, dest, sizeof(dest)); - - /* Resume: keep an existing partial only if it belongs to THIS manifest - (recorded in the sidecar); a partial from a different/older patch is - dropped. The partial survives a reboot because a killed process runs no - cleanup. Only manifest downloads resume. */ - is_manifest = url_is_manifest(patch_url); - if (is_manifest && file_size(dest) > 0) { - char prev_url[1024] = {0}; - if (read_dl_state_url(title_id, prev_url, sizeof(prev_url)) == 0 && - !strcmp(prev_url, patch_url)) - resume = 1; - else - unlink(dest); /* stale partial from a different patch */ - } - if (is_manifest) - write_dl_state(title_id, patch_url); - + /* Snapshot the g_mutex-guarded flag before taking the pool lock (lock order + is g_mutex-before-g_pool, so we cannot read it while holding g_pool.mtx). */ pthread_mutex_lock(&g_mutex); - if (g_dl.active) { - pthread_mutex_unlock(&g_mutex); - return queue_json(conn, MHD_HTTP_CONFLICT, - "{\"ok\":false,\"reason\":\"download_in_progress\"}"); - } - memset(&g_dl, 0, sizeof(g_dl)); - g_dl.active = 1; - strncpy(g_dl.title_id, title_id, sizeof(g_dl.title_id) - 1); - strncpy(g_dl.name, name && name[0] ? name : title_id, sizeof(g_dl.name) - 1); - strncpy(g_dl.version, version ? version : "", sizeof(g_dl.version) - 1); - strncpy(g_dl.path, dest, sizeof(g_dl.path) - 1); verify = g_cfg.verify_downloads; pthread_mutex_unlock(&g_mutex); - dlrc = is_manifest - ? patchdl_http_download_manifest_progress(patch_url, dest, &bytes, - download_progress_cb, NULL, verify, resume) - : patchdl_http_download_progress(patch_url, dest, &bytes, - download_progress_cb, NULL); - if (dlrc) { - int was_cancel, was_pause; - long long have; - pthread_mutex_lock(&g_mutex); - was_cancel = g_dl.cancel; - was_pause = g_dl.pause; - g_dl.active = 0; - g_dl.cancel = 0; - g_dl.pause = 0; - pthread_mutex_unlock(&g_mutex); - - if (was_cancel) { - /* user cancelled (stop): drop the partial + its resume sidecar */ - unlink(dest); - remove_dl_state(title_id); - rmdir(dir); - set_title_resumable(title_id, 0, 0); + pthread_mutex_lock(&g_pool.mtx); + for (job = g_pool.jobs; job; job = job->next) { + count++; + if (!strcmp(job->title_id, title_id)) existing = job; + } + if (existing) { + job_state_t st = existing->state; + if (st == JOB_QUEUED || st == JOB_ADMITTING || + st == JOB_ACTIVE || st == JOB_PAUSING) { + pthread_mutex_unlock(&g_pool.mtx); return queue_json(conn, MHD_HTTP_OK, - "{\"ok\":false,\"cancelled\":true," - "\"reason\":\"download_cancelled\"}"); + "{\"ok\":true,\"queued\":true,\"already\":true}"); } - if (was_pause) { - /* user paused: keep the partial + sidecar; it is now resumable. */ - have = file_size(dest); - if (have < 0) have = 0; - set_title_resumable(title_id, have > 0, have); - snprintf(resp, sizeof(resp), - "{\"ok\":false,\"paused\":true,\"reason\":\"download_paused\"," - "\"bytes\":%lld}", have); - return queue_json_owned(conn, MHD_HTTP_OK, strdup(resp)); + if (st == JOB_CANCELLING) { /* mid-cancel teardown — don't touch it */ + pthread_mutex_unlock(&g_pool.mtx); + return queue_json(conn, MHD_HTTP_CONFLICT, + "{\"ok\":false,\"reason\":\"cancelling\"}"); } - if (dlrc == -2) { - /* corrupt data (failed SHA-256): don't keep it for resume */ - unlink(dest); - remove_dl_state(title_id); - rmdir(dir); - set_title_resumable(title_id, 0, 0); - } else { - /* network failure: the partial + sidecar are kept, and the title is - now resumable (also survives a reboot). */ - set_title_resumable(title_id, 1, file_size(dest)); + if (st == JOB_PAUSED) { /* resume */ + existing->state = JOB_QUEUED; + pthread_cond_broadcast(&g_pool.cv); + pthread_mutex_unlock(&g_pool.mtx); + return queue_json(conn, MHD_HTTP_OK, + "{\"ok\":true,\"queued\":true,\"resumed\":true}"); } - snprintf(resp, sizeof(resp), "{\"ok\":false,\"reason\":\"%s\"}", - dlrc == -2 ? "piece_verify_failed" : "download_failed"); - return queue_json_owned(conn, MHD_HTTP_BAD_GATEWAY, strdup(resp)); + if (st == JOB_DONE) { /* already downloaded */ + pthread_mutex_unlock(&g_pool.mtx); + return queue_json(conn, MHD_HTTP_OK, + "{\"ok\":true,\"downloaded\":true,\"already\":true}"); + } + /* FAILED: drop it and re-enqueue fresh below */ + free_job_locked(existing); + count--; + } + if (count >= POOL_MAX_JOBS) { + pthread_mutex_unlock(&g_pool.mtx); + return queue_json(conn, MHD_HTTP_SERVICE_UNAVAILABLE, + "{\"ok\":false,\"reason\":\"queue_full\"}"); } - pthread_mutex_lock(&g_mutex); - g_dl.downloaded = bytes; - g_dl.total = bytes; - g_dl.active = 0; - pthread_mutex_unlock(&g_mutex); + job = calloc(1, sizeof(*job)); + if (!job) { + pthread_mutex_unlock(&g_pool.mtx); + return queue_text(conn, MHD_HTTP_INTERNAL_SERVER_ERROR, "oom"); + } + snprintf(job->title_id, sizeof job->title_id, "%s", title_id); + snprintf(job->name, sizeof job->name, "%s", name && name[0] ? name : title_id); + snprintf(job->version, sizeof job->version, "%s", version ? version : ""); + snprintf(job->manifest_url, sizeof job->manifest_url, "%s", patch_url); + snprintf(job->dir, sizeof job->dir, "%s/%s", PATCHDL_DL_DIR, title_id); + title_pkg_path(title_id, patch_url, job->dest, sizeof job->dest); + job->is_manifest = url_is_manifest(patch_url); + job->verify = verify; + job->state = JOB_QUEUED; + job->fd = -1; + /* append to keep FIFO queue order */ + { + dl_job_t **pp = &g_pool.jobs; + while (*pp) pp = &(*pp)->next; + *pp = job; + } + pthread_cond_broadcast(&g_pool.cv); + pthread_mutex_unlock(&g_pool.mtx); - remove_dl_state(title_id); /* complete: drop sidecar, keep the pkg */ - set_title_resumable(title_id, 0, 0); /* no longer a partial */ + return queue_json(conn, MHD_HTTP_ACCEPTED, "{\"ok\":true,\"queued\":true}"); +} - /* shadowmount: download allowed, install is not (per source policy). */ - snprintf(resp, sizeof(resp), - "{\"ok\":true,\"downloaded\":true,\"bytes\":%lld,\"path\":\"%s\"," - "\"install_allowed\":%s}", - bytes, dest, - src == PATCHDL_SOURCE_SHADOWMOUNT ? "false" : "true"); - return queue_json_owned(conn, MHD_HTTP_OK, strdup(resp)); +/* Pause: keep the partial (resumable). The active job aborts; a queued job is + just parked. */ +static enum MHD_Result +do_pause(struct MHD_Connection *conn, const char *title_id) { + dl_job_t *job; + int acted = 0; + + pthread_mutex_lock(&g_pool.mtx); + for (job = g_pool.jobs; job; job = job->next) { + if (strcmp(job->title_id, title_id)) continue; + if (job->state == JOB_ACTIVE || job->state == JOB_ADMITTING) { + /* ADMITTING: admit_next sees JOB_PAUSING on relock and parks it + PAUSED, keeping the partial — no transfer is started. */ + job->state = JOB_PAUSING; + job->abort = 1; + pthread_cond_broadcast(&g_pool.cv); + acted = 1; + } else if (job->state == JOB_QUEUED) { + job->state = JOB_PAUSED; + acted = 1; + } + break; + } + pthread_mutex_unlock(&g_pool.mtx); + if (!acted) + return queue_json(conn, MHD_HTTP_CONFLICT, + "{\"ok\":false,\"reason\":\"not_downloading\"}"); + return queue_json(conn, MHD_HTTP_OK, "{\"ok\":true,\"paused\":true}"); +} + +/* Cancel: stop AND delete. The active job's workers abort and a worker finalizes + the delete once in-flight pieces drain; an idle/queued/paused/done job is + deleted directly; a leftover on-disk partial (no job) is removed too. */ +static enum MHD_Result +do_cancel(struct MHD_Connection *conn, const char *title_id) { + dl_job_t *job; + int had_job = 0; + + pthread_mutex_lock(&g_pool.mtx); + for (job = g_pool.jobs; job; job = job->next) { + if (strcmp(job->title_id, title_id)) continue; + had_job = 1; + if (job->state == JOB_ADMITTING) { + /* being admitted (lock dropped for manifest I/O): only flag it. + admit_next holds a raw pointer to this job and is the sole + finalizer once it relocks — freeing here would be a UAF. */ + job->state = JOB_CANCELLING; + job->seq++; + job->abort = 1; + pthread_cond_broadcast(&g_pool.cv); + } else if (job->state == JOB_ACTIVE || job->state == JOB_PAUSING) { + job->state = JOB_CANCELLING; + job->seq++; /* discard late piece completions */ + job->abort = 1; + pthread_cond_broadcast(&g_pool.cv); + if (job->inflight == 0) finalize_cancel_locked(job); + } else { + /* QUEUED / PAUSED / DONE / FAILED: no in-flight pieces */ + finalize_cancel_locked(job); + } + break; + } + pthread_mutex_unlock(&g_pool.mtx); + + if (!had_job) { + /* a partial left on disk from a previous boot has no live job */ + remove_title_dir(title_id); + set_title_resumable(title_id, 0, 0); + } + return queue_json(conn, MHD_HTTP_OK, "{\"ok\":true,\"cancelled\":true}"); } static enum MHD_Result @@ -1237,6 +1693,10 @@ handle_config_post(struct MHD_Connection *conn, const char *body) { json_get_bool(body, "verify_downloads", g_cfg.verify_downloads); g_cfg.home_shortcut = json_get_bool(body, "home_shortcut", g_cfg.home_shortcut); + 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 > POOL_MAX_CONN) g_cfg.max_connections = POOL_MAX_CONN; pthread_mutex_unlock(&g_mutex); save_config(); @@ -1339,32 +1799,30 @@ on_request(void *cls, struct MHD_Connection *conn, const char *url, return queue_asset(conn, url); } -/* A title is resumable when its download dir holds both a resume sidecar and a - partial .pkg. Records the partial size for the UI. Single-threaded startup. */ +/* A title is resumable when its resume sidecar reports committed bytes between + 0 and total (the .pkg is preallocated to the full size, so its on-disk size is + not progress — done_bytes from the bitmap sidecar is). Single-threaded. */ static void detect_resumable_partials(void) { - char dir[288], state[320], pkg[576]; - DIR *d; - struct dirent *e; + char state[320], buf[16384]; + FILE *f; + size_t n; for (size_t i = 0; i < g_title_count; i++) { patchdl_title_t *t = &g_titles[i]; + long long total, done; title_state_path(t->title_id, state, sizeof(state)); - if (file_size(state) < 0) continue; /* no sidecar -> not resumable */ - snprintf(dir, sizeof(dir), "%s/%s", PATCHDL_DL_DIR, t->title_id); - d = opendir(dir); - if (!d) continue; - while ((e = readdir(d))) { - size_t nl = strlen(e->d_name); - if (nl > 4 && !strcmp(e->d_name + nl - 4, ".pkg")) { - long long sz; - snprintf(pkg, sizeof(pkg), "%s/%s", dir, e->d_name); - sz = file_size(pkg); - if (sz > 0) { t->resumable = 1; t->partial_bytes = sz; } - break; - } + f = fopen(state, "r"); + if (!f) continue; + n = fread(buf, 1, sizeof(buf) - 1, f); + fclose(f); + buf[n] = '\0'; + total = json_get_ll(buf, "total", 0); + done = json_get_ll(buf, "done_bytes", 0); + if (done > 0 && (total <= 0 || done < total)) { + t->resumable = 1; + t->partial_bytes = done; } - closedir(d); } } @@ -1372,9 +1830,6 @@ detect_resumable_partials(void) { int patchdl_websrv_start(unsigned short port) { - pthread_t tid; - pthread_attr_t attr; - if (web_daemon) return 0; /* Collect real FW version and installed titles synchronously. @@ -1398,6 +1853,23 @@ patchdl_websrv_start(unsigned short port) { swap it performs is unsafe once MHD worker threads are running. */ g_debug_json = patchdl_scan_debug_json(); + /* curl's global/OpenSSL init MUST run once, single-threaded, before the + download workers race their first curl_easy_init. */ + patchdl_net_global_init(); + + /* Spawn the N-connection download pool BEFORE the web server accepts work. */ + g_pool.n_workers = g_cfg.max_connections; + if (g_pool.n_workers < 1) g_pool.n_workers = 1; + if (g_pool.n_workers > POOL_MAX_CONN) g_pool.n_workers = POOL_MAX_CONN; + g_pool.stopping = 0; + for (int s = 0; s < g_pool.n_workers; s++) { + if (pthread_create(&g_pool.workers[s], NULL, dl_worker, + (void *)(intptr_t)s)) { + g_pool.n_workers = s; /* only the threads that started exist */ + break; + } + } + web_daemon = MHD_start_daemon( MHD_USE_INTERNAL_POLLING_THREAD | MHD_USE_THREAD_PER_CONNECTION, port, NULL, NULL, &on_request, NULL, @@ -1405,17 +1877,24 @@ patchdl_websrv_start(unsigned short port) { MHD_OPTION_END); if (!web_daemon) { + pthread_mutex_lock(&g_pool.mtx); + g_pool.stopping = 1; + pthread_cond_broadcast(&g_pool.cv); + pthread_mutex_unlock(&g_pool.mtx); + for (int s = 0; s < g_pool.n_workers; s++) + pthread_join(g_pool.workers[s], NULL); + patchdl_net_global_cleanup(); patchdl_scan_free(g_titles, g_title_count); g_titles = NULL; g_title_count = 0; return -1; } - /* Start background verxml fetch — detached, runs until complete */ - pthread_attr_init(&attr); - pthread_attr_setdetachstate(&attr, PTHREAD_CREATE_DETACHED); - pthread_create(&tid, &attr, verxml_fetch_thread, NULL); - pthread_attr_destroy(&attr); + /* Start background verxml fetch — joinable so patchdl_websrv_stop can wait + for it before freeing g_titles (it reads g_titles across blocking queries). */ + g_verxml_stop = 0; + if (pthread_create(&g_verxml_tid, NULL, verxml_fetch_thread, NULL) == 0) + g_verxml_started = 1; return 0; } @@ -1426,6 +1905,24 @@ patchdl_websrv_stop(void) { MHD_stop_daemon(web_daemon); web_daemon = NULL; } + /* Stop and join the background version.xml thread before anything frees + g_titles — it dereferences g_titles across blocking network queries. */ + if (g_verxml_started) { + g_verxml_stop = 1; + pthread_join(g_verxml_tid, NULL); + g_verxml_started = 0; + } + /* Stop the pool: signal, join every worker (each finishes its current piece + write then exits), THEN tear curl + the title list down. */ + pthread_mutex_lock(&g_pool.mtx); + g_pool.stopping = 1; + pthread_cond_broadcast(&g_pool.cv); + pthread_mutex_unlock(&g_pool.mtx); + for (int s = 0; s < g_pool.n_workers; s++) + pthread_join(g_pool.workers[s], NULL); + g_pool.n_workers = 0; + patchdl_net_global_cleanup(); + pthread_mutex_lock(&g_mutex); patchdl_scan_free(g_titles, g_title_count); g_titles = NULL; diff --git a/web/app.js b/web/app.js index 981623c..2fc7fd3 100644 --- a/web/app.js +++ b/web/app.js @@ -29,6 +29,7 @@ const fallback = { delete_pkg_after_install: true, verify_downloads: false, home_shortcut: true, + max_connections: 4, source_policy: { official: { allow_check: true, allow_download: true, allow_install: true }, external: { allow_check: true, allow_download: true, allow_install: true }, @@ -116,6 +117,7 @@ function bindElements() { deleteAfterInstall: document.getElementById("deleteAfterInstall"), verifyDownloads: document.getElementById("verifyDownloads"), homeShortcut: document.getElementById("homeShortcut"), + maxConnections: document.getElementById("maxConnections"), refreshBtn: document.getElementById("refreshBtn"), saveBtn: document.getElementById("saveBtn"), clearLogBtn: document.getElementById("clearLogBtn"), @@ -173,27 +175,26 @@ async function loadInitialData() { // /api/titles carries no client-only progress flags, so preserve them across a // refresh — otherwise an in-flight download/install flips back to a clickable // button mid-operation. + // Carry client-only flags across a refresh; the pool's job list is the source + // of truth for download state and is reconciled right after. const prev = new Map(state.titles.map((g) => [g.title_id, g])); titles.forEach((g) => { const old = prev.get(g.title_id); - // Carry client-only progress flags across a refresh — unless the server now - // reports the title up to date (the patch applied), in which case drop them. if (old && g.status !== "up_to_date") { - if (old.downloading) g.downloading = true; if (old.downloaded) g.downloaded = true; if (old.installing) g.installing = true; + if (old._autoInstalled) g._autoInstalled = true; if (old._localDownloading) g._localDownloading = true; } }); - // Reconcile active downloads the server reports, so progress + Cancel show even - // after a hard reload or a download started from another session/device. - const activeDl = new Set(downloads.map((d) => d.title_id)); - titles.forEach((g) => { if (activeDl.has(g.title_id)) g.downloading = true; }); state = { ...state, status, config, titles, downloads }; + reconcileFromJobs(downloads); render(); - if (state.downloads.length) startDownloadPolling(); + if (downloads.some((j) => j.state === "active" || j.state === "queued") || + state.titles.some((g) => g._localDownloading)) + startDownloadPolling(); showToast(state.usingFallback ? "Demo data loaded. API is not reachable yet." : "Data refreshed."); } @@ -234,6 +235,7 @@ function renderSettings() { els.deleteAfterInstall.checked = Boolean(state.config.delete_pkg_after_install); if (els.verifyDownloads) els.verifyDownloads.checked = Boolean(state.config.verify_downloads); if (els.homeShortcut) els.homeShortcut.checked = state.config.home_shortcut !== false; + if (els.maxConnections) els.maxConnections.value = state.config.max_connections || 4; els.allowlistHosts.replaceChildren(...(state.config.cdn_allowlist || []).map((host) => { const chip = document.createElement("span"); chip.className = "host-chip"; @@ -441,6 +443,7 @@ function progressMetaHtml(d) { const done = Number(d.bytes) || 0; const total = Number(d.total_bytes) || 0; const speed = Number(d._speed) || 0; + if (d.state === "queued") return `<span>Queued — waiting for a free slot</span>`; const parts = []; parts.push(`<span><b>${formatBytes(done)}</b>${total > 0 ? ` / ${formatBytes(total)}` : ""}</span>`); if (speed > 0) parts.push(`<span><b>${formatBytes(speed)}/s</b></span>`); @@ -463,55 +466,99 @@ function stopDownloadPolling() { downloadPollTimer = null; } +// Map the pool's job list onto per-title flags, and auto-install once a job +// finishes if "install after download" is on. +function reconcileFromJobs(jobs) { + const byId = new Map((jobs || []).map((j) => [j.title_id, j])); + state.titles.forEach((g) => { + const j = byId.get(g.title_id); + if (!j) { + // The pool never produced a job for our local intent: time it out so the + // card can't wedge in "Downloading" forever (server restart between POST + // and poll, or an unexpected response shape). + if (g._localDownloading && g._localSince && + Date.now() - g._localSince > 12000) { + g._localDownloading = false; + } + if (g.downloading && !g._localDownloading) g.downloading = false; + return; + } + g._localDownloading = false; // the pool now tracks it + if (j.state === "active" || j.state === "queued") { + g.downloading = true; + g.resumable = false; + g._wasActive = true; + } else if (j.state === "paused") { + g.downloading = false; + g.resumable = true; + g.partial_bytes = Number(j.bytes) || g.partial_bytes || 0; + g._wasActive = false; + } else if (j.state === "done") { + g.downloading = false; + g.resumable = false; + g.downloaded = true; + g._wasActive = false; + if (state.config.install_after_download && !g.installing && !g._autoInstalled) { + g._autoInstalled = true; + doInstall(g); + } + } else if (j.state === "error") { + // The server keeps the partial + sidecar (resumable) on a post-retry + // network failure. Reflect that so primaryButton shows "Resume" and the + // red Cancel stays available to delete the kept partial. + const bytes = Number(j.bytes) || 0; + if (g._wasActive) { + showToast(`${g.name}: download failed${bytes > 0 ? " — partial kept, press Resume to continue" : "."}`); + } + g._wasActive = false; + g.downloading = false; + g.resumable = bytes > 0; + g.partial_bytes = bytes || g.partial_bytes || 0; + } + }); +} + async function refreshDownloads() { - let downloads; + let jobs; try { const response = await fetch(API.downloads, { cache: "no-store" }); if (!response.ok) throw new Error(`HTTP ${response.status}`); - downloads = await response.json(); + jobs = await response.json(); } catch (error) { return; // keep last state on a transient failure } const now = Date.now(); - const activeIds = new Set(); - downloads.forEach((d) => { - activeIds.add(d.title_id); - const done = Number(d.bytes) || 0; - const prev = dlMeta[d.title_id]; + const ids = new Set(); + jobs.forEach((j) => { + ids.add(j.title_id); + const done = Number(j.bytes) || 0; + const prev = dlMeta[j.title_id]; if (prev && now > prev.t) { if (done >= prev.bytes) { const inst = ((done - prev.bytes) * 1000) / (now - prev.t); // bytes/s prev.speed = prev.speed ? prev.speed * 0.5 + inst * 0.5 : inst; // smoothed } else { - prev.speed = 0; // counter went backwards -> re-baseline, no stale speed + prev.speed = 0; // counter went backwards -> re-baseline } } - const meta = prev || (dlMeta[d.title_id] = { speed: 0 }); + const meta = prev || (dlMeta[j.title_id] = { speed: 0 }); meta.bytes = done; meta.t = now; - d._speed = meta.speed || 0; + j._speed = meta.speed || 0; }); - Object.keys(dlMeta).forEach((id) => { if (!activeIds.has(id)) delete dlMeta[id]; }); + Object.keys(dlMeta).forEach((id) => { if (!ids.has(id)) delete dlMeta[id]; }); - state.downloads = downloads; - - // Reconcile downloading flags with the server. A structural change (a download - // appeared or finished) needs a full re-render to add/remove the progress block - // and morph the button; otherwise update the bar in place. + state.downloads = jobs; const before = downloadingIds(); - state.titles.forEach((g) => { - if (activeIds.has(g.title_id)) g.downloading = true; - else if (g.downloading && !g._localDownloading) g.downloading = false; - }); + reconcileFromJobs(jobs); if (downloadingIds() !== before) renderGames(); else applyDownloadProgress(); - if (!downloads.length && !state.titles.some((g) => g.downloading)) { - if (++emptyPolls >= 3) stopDownloadPolling(); - } else { - emptyPolls = 0; - } + const busy = jobs.some((j) => j.state === "active" || j.state === "queued") || + state.titles.some((g) => g._localDownloading); + if (!busy) { if (++emptyPolls >= 3) stopDownloadPolling(); } + else emptyPolls = 0; } function downloadingIds() { @@ -522,6 +569,11 @@ function downloadingIds() { function applyDownloadProgress() { let needRender = false; state.downloads.forEach((d) => { + // Only titles currently downloading render a progress tile; paused/done/error + // jobs linger in the pool list but have no .progress bar — skip them so a + // missing tile for a non-downloading title doesn't force a full rebuild. + const g = state.titles.find((t) => t.title_id === d.title_id); + if (!g || !g.downloading) return; const card = els.gameGrid.querySelector(`[data-title-id="${d.title_id}"]`); const bar = card && card.querySelector(".card-progress .progress > i"); const meta = card && card.querySelector(".card-progress .progress-meta"); @@ -563,6 +615,9 @@ async function saveConfig() { delete_pkg_after_install: els.deleteAfterInstall.checked, verify_downloads: els.verifyDownloads ? els.verifyDownloads.checked : Boolean(state.config.verify_downloads), home_shortcut: els.homeShortcut ? els.homeShortcut.checked : state.config.home_shortcut !== false, + max_connections: els.maxConnections + ? Math.max(1, Math.min(16, parseInt(els.maxConnections.value, 10) || 4)) + : (state.config.max_connections || 4), }; try { await postJson(API.config, config); @@ -578,19 +633,31 @@ async function saveConfig() { /* ---------------- actions (data layer) ---------------- */ +// Enqueue a download. The pool returns immediately (202); progress, completion +// and (if configured) auto-install are driven by reconcileFromJobs() on poll. async function doDownload(game) { game.downloading = true; - game._localDownloading = true; // this client owns it; don't let a poll clear it + game._localDownloading = true; // until the pool reports a job for this title + game._localSince = Date.now(); // bounded in reconcileFromJobs if no job appears + game._autoInstalled = false; state.downloads = state.downloads.filter((i) => i.title_id !== game.title_id); - state.downloads.push({ title_id: game.title_id, name: game.name, version: game.compatible_version || "", progress: 0, bytes: 0, total_bytes: 0 }); - state.logs.push(`[${timeNow()}] Download started: ${game.title_id} ${game.compatible_version}`); + state.downloads.push({ title_id: game.title_id, name: game.name, + version: game.compatible_version || "", + state: "queued", progress: 0, bytes: 0, total_bytes: 0 }); renderGames(); - renderLogs(); startDownloadPolling(); - let r; try { - r = await postJson(API.action(game.title_id, "download"), {}); + const r = await postJson(API.action(game.title_id, "download"), {}); + if (r && r.downloaded && r.already) { + game.downloading = false; + game._localDownloading = false; + game.downloaded = true; + renderGames(); + } else { + state.logs.push(`[${timeNow()}] Download queued: ${game.title_id} ${game.compatible_version || ""}`); + renderLogs(); + } } catch (error) { game.downloading = false; game._localDownloading = false; @@ -598,40 +665,8 @@ async function doDownload(game) { const why = reasonText(error); state.logs.push(`[${timeNow()}] download ${game.title_id} blocked: ${why}`); showToast(`${game.name}: ${why}`); - renderGames(); renderLogs(); stopDownloadPolling(); - return false; + renderGames(); renderLogs(); } - - game.downloading = false; - game._localDownloading = false; - state.downloads = state.downloads.filter((i) => i.title_id !== game.title_id); - - // Cancel / pause / soft failure returns HTTP 200 with ok:false (not thrown). - if (!r || r.ok === false) { - game.downloaded = false; - if (r && r.paused) { - // Paused: keep the partial, mark resumable so the tile shows Resume. - game.resumable = true; - game.partial_bytes = r.bytes || game.partial_bytes || 0; - state.logs.push(`[${timeNow()}] Download paused: ${game.title_id} (${formatBytes(game.partial_bytes)})`); - showToast(`${game.name}: paused at ${formatBytes(game.partial_bytes)}.`); - } else { - const what = r && r.cancelled ? "cancelled" : "failed"; - state.logs.push(`[${timeNow()}] Download ${what}: ${game.title_id}`); - showToast(`${game.name}: download ${what}.`); - } - renderGames(); renderLogs(); stopDownloadPolling(); - return false; - } - - game.downloaded = true; - game.resumable = false; - game.partial_bytes = 0; - const sz = r && r.bytes ? formatBytes(r.bytes) : "?"; - state.logs.push(`[${timeNow()}] Downloaded ${game.title_id} ${game.compatible_version} (${sz}, internal)`); - showToast(`${game.name}: downloaded ${sz}.`); - renderGames(); renderLogs(); stopDownloadPolling(); - return true; } async function doInstall(game) { @@ -680,11 +715,12 @@ async function cancelDownload(titleId) { async function runTitleAction(titleId, action) { const game = state.titles.find((i) => i.title_id === titleId); if (!game) return; - if (action === "download") await doDownload(game); + // "update" and "download" both just enqueue; for "update" (install-after- + // download on) the reconciler auto-installs once the pool reports it done. + if (action === "download" || action === "update") await doDownload(game); else if (action === "install") await doInstall(game); else if (action === "pause") await doPause(game); else if (action === "cancel") await cancelDownload(titleId); - else if (action === "update") { if (await doDownload(game)) await doInstall(game); } } // Pause only sends the signal; the in-flight doDownload() request returns its diff --git a/web/index.html b/web/index.html index c4abd94..7cffcfc 100644 --- a/web/index.html +++ b/web/index.html @@ -186,6 +186,11 @@ <span class="track"></span> </span> </label> + <label class="field"> + <span>Parallel download connections</span> + <input id="maxConnections" type="number" min="1" max="16" step="1" + title="Connections used to download a patch in parallel — applies after a restart" /> + </label> </div> <div class="allowlist">