core: add the download engine — task machine, DownloadHandle, Engine
Stage 8 of the CORE build order: the bodies behind the DownloadSpec / callback API reviewed in core/docs/engine-api-m1.md. Wires probe -> segment workers -> WriteBuffer -> SparseFile -> .veloxpart.meta -> retry/backoff -> SegmentBudget -> RateLimiter -> callbacks into one event-driven machine. - Engine (src/engine.cpp): owns HttpClient, Prober, SegmentBudget, RateLimiter and one timer jthread (min-heap of scheduled fns). start() builds a task and returns a DownloadHandle; ~Impl quiesces every task before joining the timer so no callback fires during teardown. - DownloadTaskState (src/task/download_task.cpp): one `mu` task lock; a shared_mutex over the worker map for the curl write path; callbacks collected under `mu` and fired after release via a separate deferred queue; weak_from_this() in every async hop. State machine over the CORE-owned EngineState subset, auto-pause on 401/407 and on a 200 where 206 was expected, validated resume via If-Range. - digest (src/task/digest.cpp): OpenSSL EVP hash_file() for the optional post-download checksum; links OpenSSL::Crypto PRIVATE. - Segmenter::release_segment(): hand a paused segment back to the pool unassigned so resume's assign_slot() picks it up instead of splitting a still-"assigned" range and orphaning its front half. - DownloadHandle now names the real control block (vdm::task:: DownloadTaskState, defined only in the engine TU) via a namespace-scope fwd decl and a public-but-effectively-engine-only ctor, replacing the nested State/friend pair. Every public signature is unchanged; DAEMON (vdm-79) confirmed sched/ names only the public API. Fixes found while building the end-to-end suite (tests/task/engine_test.cpp, 9 cases against tools/testserver, green under ASan/UBSan and TSan): - a dropped connection lost its unflushed WriteBuffer tail while advance() had already counted those bytes as done -> a retry resumed past an unwritten hole. Flush on the failure path. - when the byte counters hit total while other workers were still live, teardown dropped their buffered tails. Now: cancel them and let each worker's own seg_finished drain it (the `assembling` state), last one starts verification -- no cross-thread buffer access. - seg_head() let a 401 with credentials present abort before libcurl's resend; now it proceeds once and acts on the final status. Co-Authored-By: Claude Sonnet 5 <[email protected]> Claude-Session: https://claude.ai/code/session_01HPPSGhiArbvQgwC2DNiURS
This commit is contained in:
@@ -0,0 +1,169 @@
|
||||
// vdm/engine.cpp
|
||||
|
||||
#include "vdm/engine.hpp"
|
||||
|
||||
#include <atomic>
|
||||
#include <condition_variable>
|
||||
#include <cstdint>
|
||||
#include <map>
|
||||
#include <mutex>
|
||||
#include <thread>
|
||||
#include <unordered_map>
|
||||
|
||||
#include "task/download_task.hpp"
|
||||
|
||||
namespace vdm {
|
||||
|
||||
struct Engine::Impl : task::TaskHost {
|
||||
explicit Impl(Config c)
|
||||
: cfg_(c),
|
||||
http_(net::HttpClient::Options{.workers = c.http_workers}),
|
||||
prober_(c.probe_pool_size ? c.probe_pool_size : 4),
|
||||
budget_(segment::SegmentBudget::Options{.max_active_segments = c.max_active_segments}) {
|
||||
timer_ = std::jthread([this](std::stop_token st) { timer_loop(st); });
|
||||
}
|
||||
|
||||
~Impl() override {
|
||||
// Quiesce every task first so no callback fires during or after teardown: mark it
|
||||
// retired and cancel its transfers. Then stop the timer thread (a fn already
|
||||
// running holds a shared_ptr and finishes, but seg_finished early-returns on
|
||||
// `retired`). Then drop the registry; http_/prober_/budget_ destruct after.
|
||||
{
|
||||
std::lock_guard lk(reg_mu_);
|
||||
for (auto &[id, t] : tasks_)
|
||||
task::quiesce_task(t);
|
||||
}
|
||||
timer_.request_stop();
|
||||
timer_cv_.notify_all();
|
||||
if (timer_.joinable())
|
||||
timer_.join();
|
||||
{
|
||||
std::lock_guard lk(reg_mu_);
|
||||
tasks_.clear();
|
||||
}
|
||||
}
|
||||
|
||||
// --- TaskHost -------------------------------------------------------------------
|
||||
net::HttpClient &http() override { return http_; }
|
||||
segment::SegmentBudget &budget() override { return budget_; }
|
||||
rate::RateLimiter &limiter() override { return limiter_; }
|
||||
const Config &config() override { return cfg_; }
|
||||
|
||||
task::TimerId schedule(SteadyTime at, std::function<void()> fn) override {
|
||||
task::TimerId id;
|
||||
{
|
||||
std::lock_guard lk(timer_mu_);
|
||||
id = ++next_timer_;
|
||||
timers_.emplace(at, Entry{id, std::move(fn)});
|
||||
}
|
||||
timer_cv_.notify_all();
|
||||
return id;
|
||||
}
|
||||
void cancel_timer(task::TimerId id) override {
|
||||
std::lock_guard lk(timer_mu_);
|
||||
for (auto it = timers_.begin(); it != timers_.end(); ++it)
|
||||
if (it->second.id == id) {
|
||||
timers_.erase(it);
|
||||
return;
|
||||
}
|
||||
}
|
||||
void probe(net::ProbeRequest req, std::function<void(Result<net::ProbeResult>)> done) override {
|
||||
prober_.probe(std::move(req), std::move(done));
|
||||
}
|
||||
void task_retired(TaskId id) override {
|
||||
std::lock_guard lk(reg_mu_);
|
||||
tasks_.erase(id);
|
||||
}
|
||||
|
||||
// --- engine surface ----------------------------------------------------------------
|
||||
task::DownloadHandle start(task::DownloadSpec spec, task::DownloadCallbacks cbs) {
|
||||
TaskId id{++next_id_};
|
||||
auto st = task::create_task(*this, id, std::move(spec), std::move(cbs));
|
||||
{
|
||||
std::lock_guard lk(reg_mu_);
|
||||
tasks_[id] = st;
|
||||
}
|
||||
return task::DownloadHandle(std::move(st));
|
||||
}
|
||||
|
||||
Config cfg_;
|
||||
net::HttpClient http_;
|
||||
net::Prober prober_;
|
||||
segment::SegmentBudget budget_;
|
||||
rate::RateLimiter limiter_;
|
||||
|
||||
std::atomic<std::uint64_t> next_id_{0};
|
||||
|
||||
std::mutex reg_mu_;
|
||||
std::unordered_map<TaskId, std::shared_ptr<task::DownloadTaskState>> tasks_;
|
||||
|
||||
struct Entry {
|
||||
task::TimerId id;
|
||||
std::function<void()> fn;
|
||||
};
|
||||
std::mutex timer_mu_;
|
||||
std::condition_variable timer_cv_;
|
||||
std::multimap<SteadyTime, Entry> timers_;
|
||||
std::atomic<std::uint64_t> next_timer_{0};
|
||||
std::jthread timer_;
|
||||
|
||||
void timer_loop(std::stop_token st) {
|
||||
std::unique_lock lk(timer_mu_);
|
||||
while (!st.stop_requested()) {
|
||||
if (timers_.empty()) {
|
||||
timer_cv_.wait_for(lk, std::chrono::seconds(1));
|
||||
continue;
|
||||
}
|
||||
auto next_at = timers_.begin()->first;
|
||||
if (timer_cv_.wait_until(lk, next_at, [&] {
|
||||
return st.stop_requested() ||
|
||||
(!timers_.empty() && timers_.begin()->first < next_at);
|
||||
})) {
|
||||
continue; // stop, or an earlier timer landed — re-evaluate
|
||||
}
|
||||
// fire everything due
|
||||
std::vector<std::function<void()>> due;
|
||||
auto now = std::chrono::steady_clock::now();
|
||||
for (auto it = timers_.begin(); it != timers_.end() && it->first <= now;)
|
||||
due.push_back(std::move(it->second.fn)), it = timers_.erase(it);
|
||||
lk.unlock();
|
||||
for (auto &fn : due)
|
||||
fn();
|
||||
lk.lock();
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
// --- Engine ---------------------------------------------------------------------------
|
||||
|
||||
Engine::Engine() : Engine(Config{}) {}
|
||||
Engine::Engine(Config cfg) : impl_(std::make_unique<Impl>(cfg)) {}
|
||||
Engine::~Engine() = default;
|
||||
|
||||
task::DownloadHandle Engine::start(task::DownloadSpec spec, task::DownloadCallbacks cbs) {
|
||||
return impl_->start(std::move(spec), std::move(cbs));
|
||||
}
|
||||
|
||||
segment::SegmentBudget &Engine::segment_budget() noexcept {
|
||||
return impl_->budget_;
|
||||
}
|
||||
rate::RateLimiter &Engine::rate_limiter() noexcept {
|
||||
return impl_->limiter_;
|
||||
}
|
||||
|
||||
void Engine::set_default_segments(std::uint32_t n) {
|
||||
impl_->cfg_.default_segments = n ? n : 1;
|
||||
}
|
||||
void Engine::set_default_buffer_bytes(std::uint64_t b) {
|
||||
impl_->cfg_.default_buffer_bytes = b;
|
||||
}
|
||||
void Engine::set_max_total_buffer_bytes(std::uint64_t b) {
|
||||
impl_->cfg_.max_total_buffer_bytes = b;
|
||||
}
|
||||
void Engine::set_probe_pool_size(std::uint32_t) { /* prober pool is fixed at construction in M1 */ }
|
||||
|
||||
void Engine::probe(net::ProbeRequest req, std::function<void(Result<net::ProbeResult>)> done) {
|
||||
impl_->prober_.probe(std::move(req), std::move(done));
|
||||
}
|
||||
|
||||
} // namespace vdm
|
||||
Reference in New Issue
Block a user