// End-to-end: a real Engine against tools/testserver, covering the CORE M1 DoD paths. #include "vdm/engine.hpp" #include #include #include #include #include #include #include #include #include "task/digest.hpp" #include "testserver_fixture.hpp" #include "vtest.hpp" using namespace vdm; using namespace vdm::task; using vdm::testing::TestServer; using namespace std::chrono_literals; namespace { struct TmpDir { std::string path; TmpDir() { const char *d = std::getenv("TMPDIR"); path = (d ? d : "/tmp"); path += "/vdm_engine_XXXXXX"; path = ::mkdtemp(path.data()) ? path : ""; } ~TmpDir() { // best-effort recursive cleanup of our flat dir if (path.empty()) return; std::string cmd = "rm -rf '" + path + "'"; (void)std::system(cmd.c_str()); } std::string file(const std::string &name) const { return path + "/" + name; } }; struct Recorder { std::promise> done; std::future> fut = done.get_future(); std::atomic fired{false}; std::vector states; std::mutex mu; std::atomic auth_calls{0}; std::atomic decision_calls{0}; // A probe callback can fire before the caller has stored the handle returned by // eng.start(). Callbacks that reach back into the handle wait on this. std::atomic handle_ready{false}; void arm(DownloadHandle &) { handle_ready.store(true, std::memory_order_release); } DownloadCallbacks cbs(DownloadHandle *h = nullptr, std::string user = "", std::string pass = "") { DownloadCallbacks c; c.on_state = [this](EngineState, EngineState to, const std::optional &) { std::lock_guard lk(mu); states.push_back(to); }; c.on_decision_needed = [this](const DecisionRequest &) { decision_calls.fetch_add(1); }; c.on_finished = [this](Result r) { if (!fired.exchange(true)) done.set_value(std::move(r)); }; if (h) { c.on_auth_required = [this, h, user, pass](const AuthChallenge &) { auth_calls.fetch_add(1); while (!handle_ready.load(std::memory_order_acquire)) std::this_thread::sleep_for(1ms); h->provide_auth(user, pass, false); }; } return c; } Result wait(std::chrono::seconds to = 40s) { if (fut.wait_for(to) != std::future_status::ready) return Err{Error::timeout, "engine test wait"}; return fut.get(); } bool saw(EngineState s) { std::lock_guard lk(mu); for (auto x : states) if (x == s) return true; return false; } }; std::uint64_t file_size(const std::string &p) { int fd = ::open(p.c_str(), O_RDONLY); if (fd < 0) return ~0ull; off_t e = ::lseek(fd, 0, SEEK_END); ::close(fd); return e < 0 ? ~0ull : static_cast(e); } // The reference SHA-256 the testserver will report for a given path. std::string server_sha(TestServer &srv, const std::string &mode, const std::string &size) { // one-shot GET //sha256/ via a throwaway Engine::probe? no — use raw curl // through a small helper. Simplest: shell out. std::string url = srv.url("/" + mode + "/sha256/" + size); std::string cmd = "curl -s '" + url + "'"; std::string out; if (FILE *f = ::popen(cmd.c_str(), "r")) { char buf[512]; while (std::fgets(buf, sizeof buf, f)) out += buf; ::pclose(f); } auto q = out.find("\"sha256\""); 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); s.save_path = save; return s; } } // namespace VT_TEST(engine_plain_multisegment_download) { TestServer srv; VT_REQUIRE(srv.available()); TmpDir td; VT_REQUIRE(!td.path.empty()); Recorder rec; Engine eng; auto h = eng.start(spec_for(srv, "/plain/file/4M", td.file("a.bin")), rec.cbs()); auto r = rec.wait(); VT_REQUIRE(r.has_value()); VT_CHECK_EQ(r.value().final_path, td.file("a.bin")); VT_CHECK_EQ(r.value().bytes, 4u * 1024 * 1024); VT_CHECK_EQ(file_size(td.file("a.bin")), 4u * 1024 * 1024); VT_CHECK(rec.saw(EngineState::downloading)); VT_CHECK(rec.saw(EngineState::complete)); auto got = hash_file(td.file("a.bin"), Checksum::Algo::sha256); VT_REQUIRE(got.has_value()); VT_CHECK_EQ(got.value(), server_sha(srv, "plain", "4M")); // the sidecar is gone on success VT_CHECK_EQ(::access((td.file("a.bin") + ".veloxpart.meta").c_str(), F_OK), -1); } VT_TEST(engine_checksum_pass_and_mismatch) { TestServer srv; VT_REQUIRE(srv.available()); TmpDir td; Engine eng; std::string want = server_sha(srv, "plain", "1M"); VT_REQUIRE(!want.empty()); { Recorder rec; auto s = spec_for(srv, "/plain/file/1M", td.file("ok.bin")); s.checksum = Checksum{Checksum::Algo::sha256, want}; auto h = eng.start(std::move(s), rec.cbs()); auto r = rec.wait(); VT_REQUIRE(r.has_value()); VT_CHECK(rec.saw(EngineState::verifying)); } { Recorder rec; auto s = spec_for(srv, "/plain/file/1M", td.file("bad.bin")); s.checksum = Checksum{Checksum::Algo::sha256, std::string(64, 'a')}; auto h = eng.start(std::move(s), rec.cbs()); auto r = rec.wait(); VT_REQUIRE(!r.has_value()); VT_CHECK_EQ(r.error().code, Error::checksum_mismatch); } } VT_TEST(engine_non_resumable_single_segment) { TestServer srv; VT_REQUIRE(srv.available()); TmpDir td; Recorder rec; Engine eng; auto h = eng.start(spec_for(srv, "/no-range/file/2M", td.file("nr.bin")), rec.cbs()); auto r = rec.wait(); VT_REQUIRE(r.has_value()); VT_CHECK_EQ(file_size(td.file("nr.bin")), 2u * 1024 * 1024); auto got = hash_file(td.file("nr.bin"), Checksum::Algo::sha256); VT_CHECK_EQ(got.value(), server_sha(srv, "no-range", "2M")); } VT_TEST(engine_404_is_an_error) { TestServer srv; VT_REQUIRE(srv.available()); TmpDir td; Recorder rec; Engine eng; auto h = eng.start(spec_for(srv, "/plain/nope", td.file("x.bin")), rec.cbs()); auto r = rec.wait(); VT_REQUIRE(!r.has_value()); VT_CHECK_EQ(r.error().code, Error::not_found); VT_CHECK(rec.saw(EngineState::failed)); } VT_TEST(engine_cancel_mid_download) { TestServer srv; VT_REQUIRE(srv.available()); TmpDir td; Recorder rec; Engine eng; auto h = eng.start(spec_for(srv, "/throttled/file/8M", td.file("c.bin")), rec.cbs()); for (int i = 0; i < 200 && !rec.saw(EngineState::downloading); ++i) std::this_thread::sleep_for(10ms); VT_REQUIRE(rec.saw(EngineState::downloading)); h.cancel(/*discard_partial=*/true); auto r = rec.wait(); VT_REQUIRE(!r.has_value()); VT_CHECK_EQ(r.error().code, Error::canceled); VT_CHECK(rec.saw(EngineState::cancelled)); VT_CHECK_EQ(::access((td.file("c.bin") + ".veloxpart").c_str(), F_OK), -1); // discarded } VT_TEST(engine_pause_resume_completes) { TestServer srv; VT_REQUIRE(srv.available()); TmpDir td; Recorder rec; Engine eng; auto h = eng.start(spec_for(srv, "/throttled/file/2M", td.file("pr.bin")), rec.cbs()); for (int i = 0; i < 200 && !rec.saw(EngineState::downloading); ++i) std::this_thread::sleep_for(10ms); h.pause(); for (int i = 0; i < 100 && h.state() != EngineState::paused; ++i) std::this_thread::sleep_for(20ms); VT_CHECK_EQ(h.state(), EngineState::paused); h.resume(); auto r = rec.wait(90s); VT_REQUIRE(r.has_value()); VT_CHECK_EQ(file_size(td.file("pr.bin")), 2u * 1024 * 1024); auto got = hash_file(td.file("pr.bin"), Checksum::Algo::sha256); VT_CHECK_EQ(got.value(), server_sha(srv, "throttled", "2M")); } VT_TEST(engine_resume_after_a_fresh_task) { // Simulates kill -9: cancel WITHOUT discard, then a new task with allow_resume picks // up the .veloxpart[.meta] and finishes with a byte-identical file. TestServer srv; VT_REQUIRE(srv.available()); TmpDir td; std::string save = td.file("resume.bin"); std::string want = server_sha(srv, "throttled", "3M"); { Recorder rec; Engine eng; auto h = eng.start(spec_for(srv, "/throttled/file/3M", save), rec.cbs()); for (int i = 0; i < 800 && (h.progress().downloaded < 512u * 1024); ++i) std::this_thread::sleep_for(10ms); VT_REQUIRE(h.progress().downloaded >= 512u * 1024); h.cancel(/*discard_partial=*/false); (void)rec.wait(); VT_CHECK_EQ(::access((save + ".veloxpart.meta").c_str(), F_OK), 0); // sidecar kept } { Recorder rec; Engine eng; auto s = spec_for(srv, "/throttled/file/3M", save); // same source -> sidecar validates s.allow_resume = true; auto h = eng.start(std::move(s), rec.cbs()); auto r = rec.wait(120s); VT_REQUIRE(r.has_value()); VT_CHECK_EQ(file_size(save), 3u * 1024 * 1024); auto got = hash_file(save, Checksum::Algo::sha256); VT_CHECK_EQ(got.value(), want); // byte-identical after resume } } VT_TEST(engine_flaky_reset_retries_to_completion) { TestServer srv; VT_REQUIRE(srv.available()); TmpDir td; Recorder rec; Engine eng; // flaky-reset RSTs the first two attempts per (path,range); a single-segment request // therefore needs the retry loop. auto s = spec_for(srv, "/flaky-reset/file/256K", td.file("fl.bin")); s.segments = 1; s.max_retries = 40; // the server RSTs at the halfway point every attempt -> ~18 halvings auto h = eng.start(std::move(s), rec.cbs()); auto r = rec.wait(120s); VT_REQUIRE(r.has_value()); VT_CHECK_EQ(file_size(td.file("fl.bin")), 256u * 1024); auto got = hash_file(td.file("fl.bin"), Checksum::Algo::sha256); VT_CHECK_EQ(got.value(), server_sha(srv, "flaky-reset", "256K")); VT_CHECK(rec.saw(EngineState::retry_wait) || rec.saw(EngineState::connecting)); } VT_TEST(engine_401_then_provide_auth_completes) { 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-basic/file/1M", td.file("au.bin")), std::move(cbs)); rec.arm(h); // publish h to the auth callback (which may already be waiting) auto r = rec.wait(); VT_REQUIRE(r.has_value()); VT_CHECK(rec.auth_calls.load() >= 1); VT_CHECK_EQ(file_size(td.file("au.bin")), 1u * 1024 * 1024); } // --- 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). --- VT_TEST(engine_etag_changes_asks_instead_of_splicing) { // A server that revalidates with a different ETag on every response fails an If-Range // on any retry or resume. That must surface as "ask the user" (server_file_changed), // never as a silent restart-from-offset-0 spliced onto bytes already on disk. TestServer srv; VT_REQUIRE(srv.available()); TmpDir td; Recorder rec; Engine eng; auto h = eng.start(spec_for(srv, "/throttled+etag-changes/file/2M", td.file("ec.bin")), rec.cbs()); // Get real progress on at least one segment before pausing, so resume's If-Range (only // sent once a segment has completed > 0) actually fires. for (int i = 0; i < 300 && h.progress().downloaded < 64u * 1024; ++i) std::this_thread::sleep_for(10ms); VT_REQUIRE(h.progress().downloaded >= 64u * 1024); h.pause(); for (int i = 0; i < 200 && h.state() != EngineState::paused; ++i) std::this_thread::sleep_for(20ms); VT_REQUIRE(h.state() == EngineState::paused); h.resume(); 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); h.decide(Decision::restart); auto r = rec.wait(90s); VT_REQUIRE(r.has_value()); VT_CHECK_EQ(file_size(td.file("ec.bin")), 2u * 1024 * 1024); auto got = hash_file(td.file("ec.bin"), Checksum::Algo::sha256); VT_CHECK_EQ(got.value(), server_sha(srv, "throttled+etag-changes", "2M")); } VT_TEST(engine_416_mid_download_asks_instead_of_exhausting_retries) { // 416-always 416s every ranged request, including the probe's own -- a live probe // correctly concludes "not resumable" and a plain-GET download never touches Range // (that path is the same shape as engine_non_resumable_single_segment). The failure // mode docs/04 means -- a server that *was* proven resumable dropping Range support // mid-download -- needs a worker to actually send Range against it, so force the // resumable, multi-segment assumption directly via probe_hint. TestServer srv; VT_REQUIRE(srv.available()); TmpDir td; Recorder rec; Engine eng; net::ProbeResult hint; hint.total_size = 256u * 1024; hint.last_modified = "Wed, 01 Jan 2025 00:00:00 GMT"; hint.accept_ranges = true; hint.resumable = true; auto s = spec_for(srv, "/416-always/file/256K", td.file("rb.bin")); s.probe_hint = hint; s.segments = 2; 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); // 416-always never recovers -- re-probing would just 416 again -- so the only sound // resolution is to stop, honestly, rather than retry the stale range until exhaustion. h.decide(Decision::abort); auto r = rec.wait(); VT_REQUIRE(!r.has_value()); VT_CHECK_EQ(r.error().code, Error::range_not_satisfiable); VT_CHECK_EQ(::access(td.file("rb.bin").c_str(), F_OK), -1); // never declared complete } VT_TEST(engine_lies_about_accept_ranges_demotes_without_asking) { // Ranges are always ignored (a plain 200, full body) but ETag/Last-Modified are stable // and honest -- unlike etag-changes, this is provably the *same* file, just a Range- // blind connection. docs/04 §7: demote to 1 segment and continue, automatically, no // user round-trip. As with 416-always, a live probe already gets this right up front // (proven non-resumable), so probe_hint forces the interesting mid-download case. TestServer srv; VT_REQUIRE(srv.available()); TmpDir td; Recorder rec; Engine eng; std::string want = server_sha(srv, "lies-about-accept-ranges", "512K"); VT_REQUIRE(!want.empty()); net::ProbeResult hint; hint.total_size = 512u * 1024; hint.last_modified = "Wed, 01 Jan 2025 00:00:00 GMT"; // testserver sends this verbatim hint.accept_ranges = true; hint.resumable = true; auto s = spec_for(srv, "/lies-about-accept-ranges/file/512K", td.file("lar.bin")); s.probe_hint = hint; s.segments = 4; auto h = eng.start(std::move(s), rec.cbs()); auto r = rec.wait(60s); VT_REQUIRE(r.has_value()); VT_CHECK_EQ(rec.decision_calls.load(), 0); // demoted automatically, not asked VT_CHECK_EQ(file_size(td.file("lar.bin")), 512u * 1024); auto got = hash_file(td.file("lar.bin"), Checksum::Algo::sha256); VT_CHECK_EQ(got.value(), want); } VT_TEST(engine_content_length_mismatch_fails_honestly) { // Content-Length promises the true size but the connection always closes short of it. // There is no recovery (unlike flaky-reset, this never "heals" on a later attempt), so // the segment's remaining range shrinks every retry until it stalls at zero progress. // The only correct outcome is a real, visible failure -- never a rename to save_path // built from a file that is quietly missing bytes. TestServer srv; VT_REQUIRE(srv.available()); TmpDir td; Recorder rec; Engine eng; auto s = spec_for(srv, "/content-length-mismatch/file/16K", td.file("clm.bin")); s.segments = 1; s.max_retries = 4; auto h = eng.start(std::move(s), rec.cbs()); auto r = rec.wait(60s); VT_REQUIRE(!r.has_value()); VT_CHECK_EQ(r.error().code, Error::max_retries_exhausted); VT_CHECK(rec.saw(EngineState::failed)); VT_CHECK_EQ(::access(td.file("clm.bin").c_str(), F_OK), -1); // never renamed into place } // --- 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 -- // this exercises exactly that path. --- VT_TEST(engine_polled_progress_reports_nonzero_speed) { TestServer srv; VT_REQUIRE(srv.available()); TmpDir td; Recorder rec; Engine eng; auto h = eng.start(spec_for(srv, "/throttled/file/4M", td.file("sp.bin")), rec.cbs()); // Give it real, sustained progress: the speed estimate only updates on a >=0.5s // sample window (seg_data), so a snapshot taken too early would legitimately read 0 // even with the bug fixed. Poll until downloaded has clearly advanced twice over. std::uint64_t speed = 0; for (int i = 0; i < 400 && speed == 0; ++i) { std::this_thread::sleep_for(20ms); auto p = h.progress(); if (p.downloaded >= 256u * 1024) speed = p.speed_bps; } VT_CHECK(speed > 0); h.cancel(/*discard_partial=*/true); auto r = rec.wait(); VT_REQUIRE(!r.has_value()); }