diff --git a/core/src/task/download_task.cpp b/core/src/task/download_task.cpp index a261990..486f8fb 100644 --- a/core/src/task/download_task.cpp +++ b/core/src/task/download_task.cpp @@ -64,6 +64,19 @@ std::string lower(std::string s) { return s; } +// scheme://host[:port] of `url`, with no path/query/fragment -- what docs/04 §7's "403 +// after redirect: retry once with the original referrer" retries with as the Referer +// header. Empty on an unparseable URL (the caller just won't get a referrer retry). +std::string origin_of(std::string_view url) { + auto s = net::split_url(url); + if (!s.valid) + return {}; + std::string out = s.scheme + "://" + s.host; + if (s.port) + out += ":" + std::to_string(*s.port); + return out; +} + } // namespace // What to do once every worker has drained (see DownloadTaskState::begin_drain_locked). @@ -88,6 +101,7 @@ struct SegWorker { bool needs_auth = false; bool wrong_status = false; bool range_bad = false; + bool forbidden = false; // 403 -- docs/04 §7's "retry once with the original referrer" bool auth_handshake = false; // saw a 401/407 and let libcurl resend with credentials std::string resp_etag, resp_last_modified; // captured on a wrong_status 200, for demote std::optional flush_error; @@ -135,6 +149,16 @@ struct DownloadTaskState : std::enable_shared_from_this { std::optional total_size; std::string origin_host; + // docs/04 §7's "403 after redirect: retry once with the original referrer" -- many + // CDNs 403 a bare/foreign Referer. Starts as spec.referrer (the browser's, verbatim); + // start_worker_locked() sends this, not spec.referrer directly, so a 403 retry can + // override it (to the download URL's own origin) without touching what the caller + // actually asked for. referrer_retried bounds it to exactly once per task -- a second + // 403 with a same-origin Referer already set is a real, honest failure + // (Error::forbidden), not something a referrer swap can fix. + std::string effective_referrer; + bool referrer_retried = false; + std::unique_ptr seg; std::unique_ptr file; std::unordered_map> workers; @@ -164,7 +188,7 @@ struct DownloadTaskState : std::enable_shared_from_this { std::vector> deferred; DownloadTaskState(TaskHost &h, TaskId i, DownloadSpec s, DownloadCallbacks c) - : host(h), id(i), spec(std::move(s)), cbs(std::move(c)) {} + : host(h), id(i), spec(std::move(s)), cbs(std::move(c)), effective_referrer(spec.referrer) {} // --- deferred callbacks ------------------------------------------------------------- void defer(std::function fn) { @@ -286,7 +310,7 @@ void DownloadTaskState::restart_probe(bool with_auth) { pr.url = spec.url; pr.headers = spec.headers; pr.cookies = spec.cookies; - pr.referrer = spec.referrer; + pr.referrer = effective_referrer; pr.user_agent = spec.user_agent; pr.proxy = spec.proxy; if (with_auth) @@ -299,12 +323,37 @@ void DownloadTaskState::restart_probe(bool with_auth) { } void DownloadTaskState::on_probe_result(Result r) { + bool retry_probe_with_referrer = false; { std::unique_lock lk(mu); if (retired.load() || is_terminal(state)) return; if (!r.has_value()) { - fail_locked(std::move(r).error()); + ErrorInfo e = std::move(r).error(); + // docs/04 §7's referrer retry applies here too: a probe (HEAD, or the + // ranged-GET fallback when HEAD is refused -- probe.cpp) can be the request + // that actually gets 403'd, before any segment worker exists to retry it + // (net::Prober builds its own request from ProbeRequest::referrer, not + // through start_worker_locked() -- see restart_probe()'s use of + // effective_referrer below). Same one-shot bound via referrer_retried as the + // worker-level retry (seg_finished's w->forbidden branch) shares. + if (e.code == Error::forbidden && !referrer_retried) { + referrer_retried = true; + effective_referrer = origin_of(spec.url); + retry_probe_with_referrer = true; + } else if (e.code == Error::forbidden) { + // Already retried with the origin referrer and still 403 -- not something + // another blind retry fixes (an expired signed URL, a private resource). + // Ask rather than fail outright, the same "ask, don't just fail" shape as + // wrong_status/range_bad/the worker-level 403 branch: refresh_url() is a + // no-op once the task is terminal, and tools/testserver's expiring-signed- + // url mode (also a bare 403, indistinguishable from any other without + // parsing the body -- CLAUDE.md §3, core never does) is meant to be + // recovered exactly that way. + auto_pause_locked(std::move(e), false, true); + } else { + fail_locked(std::move(e)); + } } else { probe = std::move(r).value(); have_probe = true; @@ -319,6 +368,8 @@ void DownloadTaskState::on_probe_result(Result r) { } } flush_deferred(); + if (retry_probe_with_referrer) + restart_probe(false); } void DownloadTaskState::finish_probe_locked() { @@ -472,7 +523,7 @@ void DownloadTaskState::start_worker_locked(std::uint32_t seg_idx) { req.url = current_url(); req.headers = spec.headers; req.cookies = spec.cookies; - req.referrer = spec.referrer; + req.referrer = effective_referrer; req.user_agent = spec.user_agent; req.proxy = spec.proxy; req.auth = spec.auth; @@ -552,6 +603,10 @@ net::DataAction DownloadTaskState::seg_head(std::uint32_t seg_idx, const net::Re w->range_bad = true; return net::DataAction::abort; } + if (h.status == 403) { + w->forbidden = true; + return net::DataAction::abort; + } if (h.status >= 400) return net::DataAction::abort; seg->set_segment_state(seg_idx, segment::SegState::downloading); @@ -637,7 +692,7 @@ void DownloadTaskState::seg_finished(std::uint32_t seg_idx, Resultneeds_auth && !w->wrong_status && - !w->flush_error && !w->range_bad) { + !w->flush_error && !w->range_bad && !w->forbidden) { const std::uint64_t len = seg->segment_end(seg_idx) - seg->segment_start(seg_idx) + 1; if (len != 0 && seg->segment_completed(seg_idx) >= len) { r = Result(net::TransferStats{}); @@ -756,6 +811,35 @@ void DownloadTaskState::seg_finished(std::uint32_t seg_idx, Resultforbidden) { + // docs/04 §7: "403 after redirect: retry once with the original referrer -- many + // CDNs require it." Bare/foreign Referer is the common cause; origin_of() rebuilds + // it from the (possibly redirected) URL the response actually came from. Exactly + // once per task, not a backoff series -- a second 403 with a same-origin Referer + // already set isn't something another blind retry can fix (a private/expired + // resource, an expiring signed URL past its window, ...). That's not necessarily + // terminal, though: ask (same "ask, don't just fail outright" shape as + // wrong_status/range_bad above) rather than fail_locked() outright, specifically + // so DownloadHandle::refresh_url() -- do_refresh_url() is a no-op once the task is + // terminal -- stays usable for the case tools/testserver's README pairs it with: + // a caller that gets a fresh signed URL and hands it back. + release_slot(); + if (!referrer_retried) { + referrer_retried = true; + effective_referrer = origin_of(current_url()); + seg->set_segment_state(seg_idx, segment::SegState::stalled); + auto wp = weak_from_this(); + host.schedule(std::chrono::steady_clock::now(), [wp, seg_idx] { + if (auto s = wp.lock()) + s->retry_worker(seg_idx); + }); + if (workers.empty()) + transition(EngineState::retry_wait, std::nullopt); + } else { + auto_pause_locked(ErrorInfo(Error::forbidden, "403", w->http_status), false, true); + } + return done(); + } if (!r.has_value()) { ErrorInfo e = std::move(r).error(); @@ -1215,6 +1299,7 @@ void DownloadTaskState::do_refresh_url(std::string url, std::vector r) { @@ -1224,10 +1309,43 @@ void DownloadTaskState::do_refresh_url(std::string url, std::vectormu); if (s->retired.load() || is_terminal(s->state)) return; - if (r.has_value()) { - s->probe.effective_url = r.value().effective_url; - s->probe.etag = r.value().etag; - s->probe.last_modified = r.value().last_modified; + if (!r.has_value()) { + lk.unlock(); + s->flush_deferred(); + return; // still paused; the caller can retry refresh_url() or decide() + } + if (!s->have_probe) { + // The task's *first* probe never succeeded (e.g. this session's own + // expiring-signed-url path: 403, one referrer retry, still 403 -> ask rather + // than fail outright -- see on_probe_result() -- specifically so this branch + // exists to recover it). finish_probe_locked() is what actually registers the + // task with the budget and builds its Segmenter; nothing downstream of a + // partial field copy would ever start a worker without it. + s->probe = std::move(r).value(); + s->have_probe = true; + s->awaiting_auth = false; + s->awaiting_decision = false; + s->finish_probe_locked(); + lk.unlock(); + s->flush_deferred(); + return; + } + s->probe.effective_url = r.value().effective_url; + s->probe.etag = r.value().etag; + s->probe.last_modified = r.value().last_modified; + // refresh_url()'s own contract is "on a live OR PAUSED task, without losing + // progress" -- distinct from do_decide(restart), which discards progress. A task + // can be paused here for any of three reasons (a plain user pause, awaiting_auth, + // or awaiting_decision -- e.g. this session's own 403-after-referrer-retry path, + // or the pre-existing wrong_status/range_bad ones); apply_slot_target()'s guard + // blocks on awaiting_auth/awaiting_decision specifically, so leaving either set + // would have set_want() below recompute a target that nothing ever acts on -- + // the caller's new URL re-probed successfully and then the task just sat there. + // Clear both and leave `paused` the same way do_decide(restart) does. + if (s->state == EngineState::paused) { + s->awaiting_auth = false; + s->awaiting_decision = false; + s->transition(EngineState::connecting, std::nullopt); } if (s->registered) s->host.budget().set_want(s->id, s->want_slots()); diff --git a/core/tests/net/testserver_fixture.hpp b/core/tests/net/testserver_fixture.hpp index 28a4bc7..6d8e0a4 100644 --- a/core/tests/net/testserver_fixture.hpp +++ b/core/tests/net/testserver_fixture.hpp @@ -28,7 +28,14 @@ namespace vdm::testing { class TestServer { public: - TestServer() { + TestServer() : TestServer(1.0) {} + + // loris_seconds overrides the dribble duration slow-loris mode uses (default matches the + // no-arg ctor's long-standing 1s). A test that needs curl's stall detector + // (CURLOPT_LOW_SPEED_TIME, hardcoded to 30s in download_task.cpp) to actually fire needs a + // dribble that outlasts that threshold, not the short one every other test relies on to + // keep runtime down. + explicit TestServer(double loris_seconds) { const char *script = VDM_TESTSERVER_PY; if (!script || !*script || ::access(script, R_OK) != 0) return; @@ -50,8 +57,9 @@ class TestServer { int devnull = ::open("/dev/null", O_WRONLY); if (devnull >= 0) ::dup2(devnull, STDERR_FILENO); + std::string loris_str = std::to_string(loris_seconds); ::execlp("python3", "python3", script, "--port", "0", "--seed", "9", "--loris-seconds", - "1", "--throttle-bps", "131072", static_cast(nullptr)); + loris_str.c_str(), "--throttle-bps", "131072", static_cast(nullptr)); ::_exit(127); } ::close(pipefd[1]); diff --git a/core/tests/task/engine_test.cpp b/core/tests/task/engine_test.cpp index 1dadb29..94a903c 100644 --- a/core/tests/task/engine_test.cpp +++ b/core/tests/task/engine_test.cpp @@ -124,6 +124,31 @@ std::string server_sha(TestServer &srv, const std::string &mode, const std::stri return out.substr(open + 1, close - open - 1); } +// Small, deliberately identical extraction to server_sha's: GET //sign/?ttl=N +// and pull the "url" field's value out of the {"url":..., "exp":...} JSON body. +std::string sign_url(TestServer &srv, const std::string &mode, const std::string &size, + int ttl_seconds) { + std::string url = + srv.url("/" + mode + "/sign/" + size + "?ttl=" + std::to_string(ttl_seconds)); + std::string cmd = "curl -s '" + url + "'"; + std::string out; + if (FILE *f = ::popen(cmd.c_str(), "r")) { + char buf[1024]; + while (std::fgets(buf, sizeof buf, f)) + out += buf; + ::pclose(f); + } + auto q = out.find("\"url\""); + if (q == std::string::npos) + return {}; + auto colon = out.find(':', q); + auto open = out.find('"', colon); + auto close = out.find('"', open + 1); + if (open == std::string::npos || close == std::string::npos) + return {}; + return out.substr(open + 1, close - open - 1); +} + DownloadSpec spec_for(TestServer &srv, const std::string &urlpath, const std::string &save) { DownloadSpec s; s.url = srv.url(urlpath); @@ -328,6 +353,30 @@ VT_TEST(engine_401_then_provide_auth_completes) { VT_CHECK_EQ(file_size(td.file("au.bin")), 1u * 1024 * 1024); } +VT_TEST(engine_401_digest_then_provide_auth_completes) { + // Same shape as engine_401_then_provide_auth_completes, but the challenge is HTTP + // Digest (qop=auth) rather than Basic. provide_auth() doesn't know or care which -- + // http_client.cpp always asks libcurl for CURLAUTH_ANY (net::AuthScheme::any) and lets + // curl negotiate against whatever WWW-Authenticate the server actually sent -- so this + // exists purely to prove that's true end-to-end, not just at the unit level. + TestServer srv; + VT_REQUIRE(srv.available()); + TmpDir td; + Recorder rec; + Engine eng; + DownloadHandle h; + auto cbs = rec.cbs(&h, "test", "test"); + h = eng.start(spec_for(srv, "/401-digest/file/1M", td.file("dg.bin")), std::move(cbs)); + rec.arm(h); + + auto r = rec.wait(); + VT_REQUIRE(r.has_value()); + VT_CHECK(rec.auth_calls.load() >= 1); + VT_CHECK_EQ(file_size(td.file("dg.bin")), 1u * 1024 * 1024); + auto got = hash_file(td.file("dg.bin"), Checksum::Algo::sha256); + VT_CHECK_EQ(got.value(), server_sha(srv, "401-digest", "1M")); +} + // --- hostile-mode matrix: the four where a bug is silent corruption, not a visible // failure (docs/04 §5 "ask, never silently corrupt" / §7's failure-policy table). --- @@ -458,6 +507,142 @@ VT_TEST(engine_content_length_mismatch_fails_honestly) { VT_CHECK_EQ(::access(td.file("clm.bin").c_str(), F_OK), -1); // never renamed into place } +// --- remaining hostile-mode matrix (tools/testserver/README.md's mode table). --- + +VT_TEST(engine_expiring_signed_url_recovers_via_refresh_url) { + // A signed URL past its ttl 403s (tools/testserver's own JSON body distinguishes + // "expired" from "bad signature", but core never parses response bodies -- CLAUDE.md + // §3 -- so both just read as a 403). The one automatic referrer retry (see + // engine_403_without_referer_retries_with_origin, below) can't fix an expired + // signature, so the second 403 asks -- via the same auto_pause_locked(..., false, + // true) "ask, don't just fail" path as wrong_status/range_bad -- rather than + // terminally failing outright, specifically so DownloadHandle::refresh_url() (its own + // contract: works "on a live or paused task", never on a terminal one) stays usable: + // the README pairs this mode with exactly that recovery. + TestServer srv; + VT_REQUIRE(srv.available()); + TmpDir td; + Recorder rec; + Engine eng; + + std::string expired = sign_url(srv, "expiring-signed-url", "64K", /*ttl=*/1); + VT_REQUIRE(!expired.empty()); + std::this_thread::sleep_for(1500ms); // let the ttl actually pass before the first request + + DownloadSpec s; + s.url = expired; + s.save_path = td.file("exp.bin"); + auto h = eng.start(std::move(s), rec.cbs()); + + for (int i = 0; i < 300 && rec.decision_calls.load() == 0; ++i) + std::this_thread::sleep_for(20ms); + VT_REQUIRE(rec.decision_calls.load() >= 1); + VT_CHECK_EQ(h.state(), EngineState::paused); + + std::string fresh = sign_url(srv, "expiring-signed-url", "64K", /*ttl=*/60); + VT_REQUIRE(!fresh.empty()); + h.refresh_url(fresh); + + auto r = rec.wait(60s); + VT_REQUIRE(r.has_value()); + VT_CHECK_EQ(file_size(td.file("exp.bin")), 64u * 1024); + auto got = hash_file(td.file("exp.bin"), Checksum::Algo::sha256); + VT_CHECK_EQ(got.value(), server_sha(srv, "expiring-signed-url", "64K")); +} + +VT_TEST(engine_403_without_referer_retries_with_origin) { + // docs/04 §7: "403 after redirect: retry once with the original referrer -- many CDNs + // require it." No spec.referrer is set here (the common case for anything not + // initiated from a browser page, e.g. `velox add `), so the first attempt 403s; + // the engine's own retry supplies the download URL's own origin as Referer, which + // this mode accepts, and the download completes with no decision ever asked. + TestServer srv; + VT_REQUIRE(srv.available()); + TmpDir td; + Recorder rec; + Engine eng; + auto h = eng.start(spec_for(srv, "/403-without-referer/file/128K", td.file("ref.bin")), + rec.cbs()); + auto r = rec.wait(30s); + VT_REQUIRE(r.has_value()); + VT_CHECK_EQ(rec.decision_calls.load(), 0); // recovered automatically, not asked + VT_CHECK_EQ(file_size(td.file("ref.bin")), 128u * 1024); + auto got = hash_file(td.file("ref.bin"), Checksum::Algo::sha256); + VT_CHECK_EQ(got.value(), server_sha(srv, "403-without-referer", "128K")); +} + +VT_TEST(engine_redirect_chain_follows_to_completion) { + // 5 hops (tools/testserver's own --redirect-depth default) of a plain 302, query + // string preserved across each. No CORE-side logic needed for this one -- libcurl's + // own CURLOPT_FOLLOWLOCATION (RequestOptions::follow_redirects, already on) and + // CURLOPT_MAXREDIRS (default 20, well over 5) do the whole thing -- this is here as + // the end-to-end check that they're actually wired through both the probe and every + // segment worker's own request, not just one of the two. + TestServer srv; + VT_REQUIRE(srv.available()); + TmpDir td; + Recorder rec; + Engine eng; + auto h = eng.start(spec_for(srv, "/redirect-chain/file/1M", td.file("rc.bin")), rec.cbs()); + auto r = rec.wait(30s); + VT_REQUIRE(r.has_value()); + VT_CHECK_EQ(file_size(td.file("rc.bin")), 1u * 1024 * 1024); + auto got = hash_file(td.file("rc.bin"), Checksum::Algo::sha256); + VT_CHECK_EQ(got.value(), server_sha(srv, "redirect-chain", "1M")); +} + +VT_TEST(engine_slow_loris_stall_timeout_fires) { + // Status line, headers, and body dribbled out one byte at a time for --loris-seconds, + // then (if the dribble hasn't already been cut off) normal streaming -- a connection + // that's technically alive (bytes ARE arriving, just far too slowly) but must not be + // allowed to hang the task forever. http_client.cpp sets CURLOPT_LOW_SPEED_LIMIT/_TIME + // (RequestOptions::low_speed_bytes_per_sec/low_speed_secs, hardcoded in + // download_task.cpp to 1024 B/s for 30s) for exactly this. + // + // Every other test in this file uses TestServer's default 1s loris dribble to keep + // runtime down, but 1s is far shorter than curl's 30s low_speed_time: a 1s trickle + // followed by full-speed streaming never accumulates 30 CONSECUTIVE seconds under the + // floor, so curl would never actually abort it -- the download would just complete + // slightly late, which would make this test pass for the wrong reason (or not exercise + // the stall timeout at all). Explicitly ask for a dribble that outlasts the 30s + // threshold so the stall timeout is the thing actually observed firing, not assumed. + TestServer srv(40.0); + VT_REQUIRE(srv.available()); + TmpDir td; + Recorder rec; + Engine eng; + auto s = spec_for(srv, "/slow-loris/file/64K", td.file("sl.bin")); + s.segments = 1; + s.max_retries = 1; + auto h = eng.start(std::move(s), rec.cbs()); + auto r = rec.wait(60s); // stall timeout fires ~30s in; must resolve, not hang to 60s + VT_REQUIRE(!r.has_value()); + VT_CHECK(is_retryable(r.error().code) || r.error().code == Error::max_retries_exhausted); +} + +VT_TEST(engine_chunked_no_length_completes_single_segment) { + // No Content-Length anywhere (HEAD gets none either, since it's the same handler path) + // -- the probe can't know total_size or prove resumability, so this should take the + // exact same "unknown size, one plain-GET segment" path as engine_non_resumable_single_ + // segment, just arriving there via a chunked body instead of a server that plainly + // refuses Range. No core-side work needed if that demotion is already size-agnostic; + // this is here to prove it, since every other test's server tells the probe the size + // up front. + TestServer srv; + VT_REQUIRE(srv.available()); + TmpDir td; + Recorder rec; + Engine eng; + auto h = eng.start(spec_for(srv, "/chunked-no-length/file/2M", td.file("ch.bin")), + rec.cbs()); + auto r = rec.wait(); + VT_REQUIRE(r.has_value()); + VT_CHECK_EQ(rec.decision_calls.load(), 0); + VT_CHECK_EQ(file_size(td.file("ch.bin")), 2u * 1024 * 1024); + auto got = hash_file(td.file("ch.bin"), Checksum::Algo::sha256); + VT_CHECK_EQ(got.value(), server_sha(srv, "chunked-no-length", "2M")); +} + // --- DAEMON-reported bug: Progress.speed_bps reads 0 for the whole life of a live // download while downloaded bytes visibly advance. DAEMON reads progress by polling // DownloadHandle::progress() (engine_port_core.hpp), not the on_progress push callback --