The Scheduler that D4 was waiting on. Built against CORE's engine
HEADERS (now in main); the real EnginePort and the veloxd wiring wait
for lane/core's stage-8 bodies to reach main (deferrals.md D4a/D4b) —
core/src/task/ is still .gitkeep there, so linking vdm::Engine now
would be an unresolved symbol.
- sched/engine_port — the abstract seam: start/pause/resume/cancel/
provide_auth/decide/refresh_url + the ADR 0011 admission config
(set_task_order / set_max_active_segments / set_host_segment_cap).
Keeps the Scheduler testable without a live engine and the daemon
unbound from the concrete vdm::Engine.
- sched/fake_engine_port — a recording impl for tests.
- sched/scheduler:
* owns the wire-UUID <-> vdm::TaskId map.
* tick(): snapshot queues (schedule window evaluated with an
injectable clock) + non-terminal tasks -> governor.evaluate ->
apply. to_start builds a vdm::task::DownloadSpec from the row and
calls EnginePort::start; to_resume -> resume(); to_pause ->
pause() + writes the pause_reason; priority_order -> set_task_order
over the mapped engine ids. `new` tasks are parked (startMode
manual) and skipped.
* on_engine_state(wire_id, state, err): projects an engine
transition onto the store row (state, pause_reason='auto' when an
error rides a paused transition per ADR 0013 §2, flattened error
columns) so the next tick sees ground truth. This is also the hook
event.task.state will fire from (D5).
* reconcile_after_restart(): CORE-owned states -> queued, paused
keeps its reason (ADR 0013 §5).
* reload_config(): reads connection.maxConcurrentDownloads /
maxActiveSegments + a daemon-local host-cap map, pushes caps to
the engine, updates the governor.
* Deps: injectable local-now clock and a post_to_loop marshaller
(engine callbacks arrive on engine threads; default runs inline
for tests).
Test veloxd.sched_scheduler (ASan+UBSan and TSan clean): admission +
ordering, a slot freeing on completion, queue-stop -> pause
(queue_stopped) then queue-restart -> resume (not a fresh start),
engine auto-pause -> pause_reason 'auto' + never auto-resumed,
reconcile_after_restart, reload_config caps push. 35 daemon/cli tests
green.
Co-Authored-By: Claude Sonnet 5 <[email protected]>
Claude-Session: https://claude.ai/code/session_01Upd9WhG9oppieig5nRDLig
315 lines
13 KiB
C++
315 lines
13 KiB
C++
#include "sched/scheduler.hpp"
|
|
|
|
#include <algorithm>
|
|
#include <cctype>
|
|
|
|
#include <nlohmann/json.hpp>
|
|
|
|
#include "sched/schedule_window.hpp"
|
|
#include "store/settings.hpp"
|
|
#include "store/tasks.hpp"
|
|
|
|
namespace velox::daemon::sched {
|
|
|
|
namespace proto = velox::proto;
|
|
|
|
namespace {
|
|
|
|
std::tm local_now_default() {
|
|
const std::time_t t = std::time(nullptr);
|
|
std::tm tm{};
|
|
::localtime_r(&t, &tm);
|
|
return tm;
|
|
}
|
|
|
|
// host[:port] out of a URL, lowercased. "" when it cannot be parsed (opts the task out of
|
|
// the per-host cap rather than lumping unrelated tasks under "").
|
|
std::string host_of(std::string_view url) {
|
|
auto scheme = url.find("://");
|
|
std::string_view rest = scheme == std::string_view::npos ? url : url.substr(scheme + 3);
|
|
const auto at = rest.find('@');
|
|
if (at != std::string_view::npos) rest = rest.substr(at + 1);
|
|
const auto end = rest.find_first_of("/:?#");
|
|
std::string h(end == std::string_view::npos ? rest : rest.substr(0, end));
|
|
std::transform(h.begin(), h.end(), h.begin(),
|
|
[](unsigned char c) { return static_cast<char>(std::tolower(c)); });
|
|
return h;
|
|
}
|
|
|
|
// ISO-8601 timestamp -> a monotonic integer for FIFO tiebreaking ("2026-09-10T15:57:50Z"
|
|
// -> 20260910155750). Lexical order of the digits is chronological.
|
|
std::int64_t rank_of(std::string_view iso) {
|
|
std::string digits;
|
|
for (char c : iso)
|
|
if (std::isdigit(static_cast<unsigned char>(c))) digits.push_back(c);
|
|
digits.resize(std::min<std::size_t>(digits.size(), 17)); // fits in int64
|
|
return digits.empty() ? 0 : std::stoll(digits);
|
|
}
|
|
|
|
RunState run_state_of(const std::string& s) {
|
|
if (s == "queued") return RunState::Queued;
|
|
if (s == "paused") return RunState::Paused;
|
|
if (s == "complete" || s == "failed" || s == "cancelled") return RunState::Terminal;
|
|
return RunState::Running; // probing / connecting / downloading / retry_wait / assembling / verifying
|
|
}
|
|
|
|
std::optional<PauseReason> pause_reason_of(const std::optional<std::string>& s) {
|
|
if (!s) return std::nullopt;
|
|
if (*s == "user") return PauseReason::User;
|
|
if (*s == "schedule") return PauseReason::Schedule;
|
|
if (*s == "queue_stopped") return PauseReason::QueueStopped;
|
|
if (*s == "admission_reconcile") return PauseReason::AdmissionReconcile;
|
|
if (*s == "auto") return PauseReason::Auto;
|
|
return std::nullopt;
|
|
}
|
|
|
|
const char* pause_reason_str(PauseReason r) {
|
|
switch (r) {
|
|
case PauseReason::User: return "user";
|
|
case PauseReason::Schedule: return "schedule";
|
|
case PauseReason::QueueStopped: return "queue_stopped";
|
|
case PauseReason::AdmissionReconcile: return "admission_reconcile";
|
|
case PauseReason::Auto: return "auto";
|
|
}
|
|
return "user";
|
|
}
|
|
|
|
const char* engine_state_name(vdm::task::EngineState s) {
|
|
using E = vdm::task::EngineState;
|
|
switch (s) {
|
|
case E::probing: return "probing";
|
|
case E::connecting: return "connecting";
|
|
case E::downloading: return "downloading";
|
|
case E::paused: return "paused";
|
|
case E::retry_wait: return "retry_wait";
|
|
case E::assembling: return "assembling";
|
|
case E::verifying: return "verifying";
|
|
case E::complete: return "complete";
|
|
case E::failed: return "failed";
|
|
case E::cancelled: return "cancelled";
|
|
}
|
|
return "connecting";
|
|
}
|
|
|
|
// The non-terminal states the scheduler cares about — feeds the tasks.list filter.
|
|
std::vector<proto::TaskState> non_terminal_states() {
|
|
return {proto::TaskState::New, proto::TaskState::Probing,
|
|
proto::TaskState::Queued, proto::TaskState::Connecting,
|
|
proto::TaskState::Downloading, proto::TaskState::Paused,
|
|
proto::TaskState::RetryWait, proto::TaskState::Assembling,
|
|
proto::TaskState::Verifying};
|
|
}
|
|
|
|
} // namespace
|
|
|
|
Scheduler::Scheduler(store::Db& db, EnginePort& engine, Governor governor, Deps deps)
|
|
: db_(db), engine_(engine), governor_(std::move(governor)), deps_(std::move(deps)) {
|
|
if (!deps_.local_now) deps_.local_now = local_now_default;
|
|
if (!deps_.post_to_loop) deps_.post_to_loop = [](std::function<void()> f) { f(); };
|
|
}
|
|
|
|
void Scheduler::map(const std::string& wire_id, vdm::TaskId engine_id) {
|
|
to_engine_[wire_id] = engine_id;
|
|
to_wire_[engine_id] = wire_id;
|
|
}
|
|
void Scheduler::unmap_engine(vdm::TaskId engine_id) {
|
|
if (auto it = to_wire_.find(engine_id); it != to_wire_.end()) {
|
|
to_engine_.erase(it->second);
|
|
to_wire_.erase(it);
|
|
}
|
|
}
|
|
std::optional<std::string> Scheduler::wire_id_of(vdm::TaskId id) const {
|
|
auto it = to_wire_.find(id);
|
|
return it == to_wire_.end() ? std::nullopt : std::optional<std::string>(it->second);
|
|
}
|
|
std::optional<vdm::TaskId> Scheduler::engine_id_of(const std::string& wire_id) const {
|
|
auto it = to_engine_.find(wire_id);
|
|
return it == to_engine_.end() ? std::nullopt : std::optional<vdm::TaskId>(it->second);
|
|
}
|
|
|
|
store::DbResult<void> Scheduler::reconcile_after_restart() {
|
|
// Any CORE-owned state (probing..verifying) becomes `queued`; the engine knows nothing
|
|
// across a restart and will re-probe / re-resume from the .veloxpart.meta sidecar when
|
|
// the scheduler starts it again (ADR 0013 §5). `paused` and `new` are left alone.
|
|
return db_.exec(
|
|
"UPDATE tasks SET state = 'queued', pause_reason = NULL "
|
|
"WHERE state IN ('probing','connecting','downloading','retry_wait','assembling','verifying')");
|
|
}
|
|
|
|
store::DbResult<void> Scheduler::reload_config() {
|
|
store::Settings settings(db_);
|
|
GovernorConfig cfg;
|
|
cfg.max_concurrent_downloads = settings.get_int("connection.maxConcurrentDownloads");
|
|
cfg.max_active_segments = settings.get_int("connection.maxActiveSegments");
|
|
|
|
// Per-host caps: a JSON object {host: n} under a daemon-local key. Absent => none.
|
|
if (auto raw = settings.get_raw("saveTo.hostSegmentCaps"); raw && *raw) {
|
|
auto j = nlohmann::json::parse(**raw, nullptr, false);
|
|
if (j.is_object()) {
|
|
for (const auto& [host, n] : j.items()) {
|
|
if (n.is_number_integer()) {
|
|
const auto cap = n.get<std::int64_t>();
|
|
cfg.host_caps[host] = cap;
|
|
engine_.set_host_segment_cap(host, static_cast<std::uint32_t>(std::max<std::int64_t>(cap, 0)));
|
|
}
|
|
}
|
|
}
|
|
}
|
|
governor_.set_config(cfg);
|
|
engine_.set_max_active_segments(
|
|
static_cast<std::uint32_t>(std::max<std::int64_t>(cfg.max_active_segments, 1)));
|
|
return {};
|
|
}
|
|
|
|
store::DbResult<void> Scheduler::tick() {
|
|
// --- snapshot: queues -----------------------------------------------------------
|
|
std::vector<QueueView> queues;
|
|
{
|
|
auto st = db_.prepare("SELECT queue_id, state, max_concurrent, schedule FROM queues");
|
|
if (!st) return std::unexpected(st.error());
|
|
const std::tm now = deps_.local_now();
|
|
for (;;) {
|
|
auto row = st->step();
|
|
if (!row) return std::unexpected(row.error());
|
|
if (!*row) break;
|
|
QueueView q;
|
|
q.queue_id = st->column_text(0);
|
|
q.running = st->column_text(1) == "running";
|
|
q.max_concurrent = st->column_int(2);
|
|
if (!st->column_is_null(3)) {
|
|
auto j = nlohmann::json::parse(st->column_text(3), nullptr, false);
|
|
proto::Schedule sched;
|
|
if (auto p = proto::parse<proto::Schedule>(j, "schedule")) sched = *p;
|
|
q.window_open = window_open(sched, now);
|
|
}
|
|
queues.push_back(std::move(q));
|
|
}
|
|
}
|
|
|
|
// --- snapshot: tasks ----------------------------------------------------------
|
|
store::Tasks tasks(db_);
|
|
proto::TaskFilter filter;
|
|
filter.states = non_terminal_states();
|
|
auto page = tasks.list(filter, std::nullopt, 0, 100000);
|
|
if (!page) return std::unexpected(page.error());
|
|
|
|
std::vector<TaskView> views;
|
|
views.reserve(page->rows.size());
|
|
for (const auto& r : page->rows) {
|
|
if (r.state == "new") continue; // parked until the user starts it
|
|
TaskView v;
|
|
v.task_id = r.task_id;
|
|
v.run_state = run_state_of(r.state);
|
|
v.pause_reason = pause_reason_of(r.pause_reason);
|
|
v.queue_id = r.queue_id;
|
|
v.queue_position = r.queue_position.value_or(0);
|
|
v.host = host_of(r.url);
|
|
v.admit_rank = rank_of(r.created_at);
|
|
views.push_back(std::move(v));
|
|
}
|
|
|
|
const Decision d = governor_.evaluate(views, queues);
|
|
|
|
// --- apply ------------------------------------------------------------------
|
|
for (const auto& wire_id : d.to_start) {
|
|
auto got = tasks.get(wire_id);
|
|
if (!got || !got->has_value()) continue;
|
|
const store::TaskRow& row = **got;
|
|
|
|
vdm::task::DownloadSpec spec;
|
|
spec.url = row.url;
|
|
spec.save_path = row.save_dir + "/" + row.filename;
|
|
if (row.req_segments) spec.segments = static_cast<std::uint32_t>(*row.req_segments);
|
|
if (row.req_buffer_bytes)
|
|
spec.buffer_bytes = static_cast<std::uint64_t>(*row.req_buffer_bytes);
|
|
if (row.checksum_algo && row.checksum_value) {
|
|
vdm::task::Checksum ck;
|
|
ck.hex = *row.checksum_value;
|
|
if (*row.checksum_algo == "md5") ck.algo = vdm::task::Checksum::Algo::md5;
|
|
else if (*row.checksum_algo == "sha1") ck.algo = vdm::task::Checksum::Algo::sha1;
|
|
else if (*row.checksum_algo == "sha512") ck.algo = vdm::task::Checksum::Algo::sha512;
|
|
else ck.algo = vdm::task::Checksum::Algo::sha256;
|
|
spec.checksum = ck;
|
|
}
|
|
spec.allow_resume = true; // resume from a sidecar if one is beside save_path
|
|
// headers / cookies / referrer / user_agent are not persisted yet (a URL-only
|
|
// `velox add` has none); the capture path will fill them when it lands.
|
|
|
|
vdm::task::DownloadCallbacks cbs;
|
|
const std::string id_copy = wire_id;
|
|
cbs.on_state = [this, id_copy](vdm::task::EngineState, vdm::task::EngineState to,
|
|
const std::optional<vdm::ErrorInfo>& err) {
|
|
std::optional<TaskErrorFields> ef;
|
|
if (err) {
|
|
ef = TaskErrorFields{};
|
|
ef->code = std::string(vdm::error_name(err->code)); // matches TaskErrorCode
|
|
ef->message = err->context;
|
|
if (err->http_status != 0) ef->http_status = err->http_status;
|
|
ef->retryable = err->retryable;
|
|
}
|
|
const std::string to_name = engine_state_name(to);
|
|
deps_.post_to_loop(
|
|
[this, id_copy, to_name, ef]() { on_engine_state(id_copy, to_name, ef); });
|
|
};
|
|
|
|
const vdm::TaskId engine_id = engine_.start(spec, std::move(cbs));
|
|
map(wire_id, engine_id);
|
|
(void)tasks.set_state(wire_id, "probing", std::nullopt);
|
|
}
|
|
|
|
for (const auto& wire_id : d.to_resume) {
|
|
if (auto eid = engine_id_of(wire_id)) {
|
|
engine_.resume(*eid);
|
|
(void)tasks.set_state(wire_id, "connecting", std::nullopt);
|
|
}
|
|
}
|
|
|
|
for (const auto& wire_id : d.to_pause) {
|
|
const auto reason = d.pause_reasons.count(wire_id)
|
|
? pause_reason_str(d.pause_reasons.at(wire_id))
|
|
: "user";
|
|
if (auto eid = engine_id_of(wire_id)) engine_.pause(*eid);
|
|
(void)tasks.set_state(wire_id, "paused", std::string(reason));
|
|
}
|
|
|
|
std::vector<vdm::TaskId> order;
|
|
order.reserve(d.priority_order.size());
|
|
for (const auto& wire_id : d.priority_order)
|
|
if (auto eid = engine_id_of(wire_id)) order.push_back(*eid);
|
|
engine_.set_task_order(order);
|
|
|
|
return {};
|
|
}
|
|
|
|
void Scheduler::on_engine_state(const std::string& wire_id, std::string_view engine_state,
|
|
const std::optional<TaskErrorFields>& err) {
|
|
store::Tasks tasks(db_);
|
|
// pause_reason: an engine-initiated pause carries an error => 'auto' (ADR 0013 §2);
|
|
// otherwise set_state clears the column.
|
|
std::optional<std::string> reason;
|
|
if (engine_state == "paused" && err) reason = "auto";
|
|
(void)tasks.set_state(wire_id, engine_state, reason);
|
|
|
|
if (err) {
|
|
auto st = db_.prepare(
|
|
"UPDATE tasks SET error_code=?2, error_message=?3, error_http_status=?4, "
|
|
"error_retryable=?5 WHERE task_id=?1");
|
|
if (st) {
|
|
(void)st->bind(1, std::string_view(wire_id));
|
|
(void)st->bind(2, std::string_view(err->code));
|
|
(void)st->bind(3, std::string_view(err->message));
|
|
if (err->http_status) (void)st->bind(4, *err->http_status);
|
|
else (void)st->bind_null(4);
|
|
if (err->retryable) (void)st->bind(5, static_cast<std::int64_t>(*err->retryable));
|
|
else (void)st->bind_null(5);
|
|
(void)st->step();
|
|
}
|
|
}
|
|
|
|
if (engine_state == "complete" || engine_state == "failed" || engine_state == "cancelled") {
|
|
if (auto eid = engine_id_of(wire_id)) unmap_engine(*eid);
|
|
}
|
|
}
|
|
|
|
} // namespace velox::daemon::sched
|