daemon: wire the engine into veloxd — the vertical slice runs end to end
CORE stage 8 merged, so vdm::Engine is linkable. This closes D4a and narrows D4b: `velox add <url>` now actually downloads. - sched/engine_port_core.hpp — the real EnginePort: forwards to a live vdm::Engine, keeps the DownloadHandle per task for pause/resume/ cancel/provide_auth/decide/refresh_url, drives set_task_order / set_max_active_segments / set_host_segment_cap via engine.segment_budget(). CORE confirmed the admission model: DAEMON decides when to start(); the engine's own download_task calls register_task/set_want internally — DAEMON never touches per-task budget calls. EnginePort gains release(TaskId) so the port drops a handle when the task goes terminal. - rpc/event_loop — EventLoop::post(fn): thread-safe, runs fn on the loop thread next iteration. The marshaller for engine-thread callbacks. - main.cpp — constructs vdm::Engine + EnginePortCore + Scheduler (post_to_loop = loop.post). At startup: reconcile_after_restart() (ADR 0013 §5), reload_config(), tick(). A 1 s timerfd on the loop re-runs tick() (schedule windows, missed nudges); download.add nudges via dispatcher.set_on_mutation. End-to-end verified against tools/testserver: `velox add http://127.0.0.1:.../file/512K` -> task queued -> scheduler admits -> engine downloads 524288 bytes -> complete, file on disk. First byte-path all the way through the project. safepath-adversarial.md: re-verified per its own note — CORE landed O_NOFOLLOW on the target open (core/src/io/sparse_file.cpp), so the leaf-symlink TOCTOU is now closed; residual is down to one intermediate-dir gap (documented post-M1 chase). 36 daemon/cli tests green; scheduler + uds_roundtrip TSan-clean. Co-Authored-By: Claude Sonnet 5 <[email protected]> Claude-Session: https://claude.ai/code/session_01Upd9WhG9oppieig5nRDLig
This commit is contained in:
@@ -15,14 +15,20 @@
|
||||
#include <sys/un.h>
|
||||
#include <unistd.h>
|
||||
|
||||
#include <sys/timerfd.h>
|
||||
|
||||
#include "rpc/dispatcher.hpp"
|
||||
#include "rpc/event_loop.hpp"
|
||||
#include "rpc/pairing.hpp"
|
||||
#include "rpc/runtime_dir.hpp"
|
||||
#include "rpc/uds_server.hpp"
|
||||
#include "rpc/ws_server.hpp"
|
||||
#include "sched/engine_port_core.hpp"
|
||||
#include "sched/governor.hpp"
|
||||
#include "sched/scheduler.hpp"
|
||||
#include "store/migrations.hpp"
|
||||
#include "store/sqlite.hpp"
|
||||
#include "vdm/engine.hpp"
|
||||
#include "version.hpp"
|
||||
|
||||
namespace {
|
||||
@@ -101,7 +107,38 @@ int main() {
|
||||
return 1;
|
||||
}
|
||||
|
||||
// --- engine + scheduler ---------------------------------------------------------
|
||||
vdm::Engine engine;
|
||||
velox::daemon::sched::EnginePortCore engine_port(engine);
|
||||
velox::daemon::sched::Scheduler scheduler(
|
||||
*db, engine_port, velox::daemon::sched::Governor{},
|
||||
{/*local_now*/ {},
|
||||
/*post_to_loop*/ [&loop](std::function<void()> fn) { loop.post(std::move(fn)); }});
|
||||
|
||||
if (const auto ec = scheduler.reconcile_after_restart(); !ec)
|
||||
std::cerr << "veloxd: restart reconcile: " << ec.error().to_string() << "\n";
|
||||
(void)scheduler.reload_config();
|
||||
(void)scheduler.tick(); // admit anything already queued in the DB
|
||||
|
||||
velox::daemon::rpc::VeloxDispatcher dispatcher(*db);
|
||||
dispatcher.set_on_mutation([&loop, &scheduler] {
|
||||
loop.post([&scheduler] { (void)scheduler.tick(); });
|
||||
});
|
||||
|
||||
// A 1 s timer re-runs the scheduler so schedule windows opening/closing and any
|
||||
// missed nudge are picked up. Registered on the loop, no extra thread.
|
||||
const int tick_fd = ::timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK | TFD_CLOEXEC);
|
||||
if (tick_fd >= 0) {
|
||||
itimerspec spec{};
|
||||
spec.it_value.tv_sec = 1;
|
||||
spec.it_interval.tv_sec = 1;
|
||||
::timerfd_settime(tick_fd, 0, &spec, nullptr);
|
||||
loop.add_fd(tick_fd, velox::daemon::rpc::kRead, [&](int fd, unsigned) {
|
||||
std::uint64_t ticks = 0;
|
||||
[[maybe_unused]] ssize_t n = ::read(fd, &ticks, sizeof(ticks));
|
||||
(void)scheduler.tick();
|
||||
});
|
||||
}
|
||||
|
||||
velox::daemon::rpc::UdsServer uds(loop, dispatcher, rt.socket_path());
|
||||
if (const auto ec = uds.start()) {
|
||||
@@ -128,6 +165,10 @@ int main() {
|
||||
loop.run();
|
||||
std::cout << "veloxd: shutting down\n";
|
||||
|
||||
if (tick_fd >= 0) {
|
||||
loop.del_fd(tick_fd);
|
||||
::close(tick_fd);
|
||||
}
|
||||
g_loop = nullptr;
|
||||
::close(lock_fd);
|
||||
return 0;
|
||||
|
||||
@@ -151,7 +151,7 @@ VeloxDispatcher::on_download_add(const proto::DownloadSpec& spec) {
|
||||
row.created_at = velox::daemon::now_iso();
|
||||
row.start_mode = spec.startMode ? std::string(proto::to_string(*spec.startMode)) : "auto";
|
||||
// startMode 'manual' parks the task in `new`; anything else makes it eligible for the
|
||||
// scheduler (`queued`). The scheduler itself is not wired yet (deferrals.md D4).
|
||||
// scheduler (`queued`); on_mutation_ nudges it.
|
||||
row.state = row.start_mode == "manual" ? "new" : "queued";
|
||||
row.category_id = spec.categoryId;
|
||||
row.queue_id = spec.queueId;
|
||||
@@ -168,6 +168,8 @@ VeloxDispatcher::on_download_add(const proto::DownloadSpec& spec) {
|
||||
return std::unexpected(proto::HandlerError{proto::ErrorCode::InternalError,
|
||||
"download.add: " + ins.error().message});
|
||||
|
||||
if (on_mutation_) on_mutation_();
|
||||
|
||||
proto::DownloadAddResult r;
|
||||
r.taskId = row.task_id;
|
||||
if (auto st = proto::parse_TaskState(row.state)) r.state = *st;
|
||||
|
||||
@@ -12,6 +12,8 @@
|
||||
// "not implemented in this build" (-> -32603) until its handler and the scheduler land;
|
||||
// see daemon/docs/deferrals.md.
|
||||
|
||||
#include <functional>
|
||||
|
||||
#include "store/sqlite.hpp"
|
||||
#include "velox_proto.hpp"
|
||||
|
||||
@@ -21,6 +23,10 @@ class VeloxDispatcher final : public velox::proto::Dispatcher {
|
||||
public:
|
||||
explicit VeloxDispatcher(velox::daemon::store::Db& db) : db_(db) {}
|
||||
|
||||
// Called after a handler mutates task state (download.add for now). main.cpp wires it
|
||||
// to nudge the scheduler; unset in tests.
|
||||
void set_on_mutation(std::function<void()> fn) { on_mutation_ = std::move(fn); }
|
||||
|
||||
velox::proto::HandlerResult<velox::proto::CaptureRules>
|
||||
on_capture_getRules(const velox::proto::CaptureGetRulesParams&) override;
|
||||
velox::proto::HandlerResult<velox::proto::CaptureOfferResult>
|
||||
@@ -101,6 +107,7 @@ public:
|
||||
|
||||
private:
|
||||
velox::daemon::store::Db& db_;
|
||||
std::function<void()> on_mutation_;
|
||||
};
|
||||
|
||||
} // namespace velox::daemon::rpc
|
||||
|
||||
@@ -45,6 +45,23 @@ void EventLoop::stop() noexcept {
|
||||
wake();
|
||||
}
|
||||
|
||||
void EventLoop::post(std::function<void()> fn) {
|
||||
{
|
||||
std::lock_guard<std::mutex> lk(post_mu_);
|
||||
posts_.push_back(std::move(fn));
|
||||
}
|
||||
wake();
|
||||
}
|
||||
|
||||
void EventLoop::drain_posts() {
|
||||
std::vector<std::function<void()>> batch;
|
||||
{
|
||||
std::lock_guard<std::mutex> lk(post_mu_);
|
||||
batch.swap(posts_);
|
||||
}
|
||||
for (auto& fn : batch) fn();
|
||||
}
|
||||
|
||||
void EventLoop::drain_wakeup() noexcept {
|
||||
std::uint64_t sink = 0;
|
||||
while (::read(wake_fd_, &sink, sizeof(sink)) > 0) {
|
||||
@@ -87,6 +104,8 @@ void EventLoop::run() {
|
||||
if (p.revents != 0) fired.push_back(p.fd);
|
||||
}
|
||||
|
||||
drain_posts();
|
||||
|
||||
for (const int fd : fired) {
|
||||
const auto it = fds_.find(fd);
|
||||
if (it == fds_.end()) continue; // removed by an earlier callback this pass
|
||||
|
||||
@@ -12,7 +12,9 @@
|
||||
#include <atomic>
|
||||
#include <cstdint>
|
||||
#include <functional>
|
||||
#include <mutex>
|
||||
#include <unordered_map>
|
||||
#include <vector>
|
||||
|
||||
namespace velox::daemon::rpc {
|
||||
|
||||
@@ -52,6 +54,10 @@ public:
|
||||
// callback. Async-signal-safe.
|
||||
void wake() noexcept;
|
||||
|
||||
// Run `fn` on the loop thread at the next iteration. Thread-safe; the intended way to
|
||||
// marshal an engine-thread callback back onto the RPC loop.
|
||||
void post(std::function<void()> fn);
|
||||
|
||||
private:
|
||||
struct Entry {
|
||||
unsigned interest;
|
||||
@@ -59,11 +65,15 @@ private:
|
||||
};
|
||||
|
||||
void drain_wakeup() noexcept;
|
||||
void drain_posts();
|
||||
|
||||
int wake_fd_; // eventfd, always registered
|
||||
bool running_ = false;
|
||||
std::atomic<bool> stop_requested_ = false; // set from stop(), read by run()
|
||||
std::unordered_map<int, Entry> fds_;
|
||||
|
||||
std::mutex post_mu_;
|
||||
std::vector<std::function<void()>> posts_;
|
||||
};
|
||||
|
||||
} // namespace velox::daemon::rpc
|
||||
|
||||
@@ -37,6 +37,9 @@ public:
|
||||
virtual void decide(vdm::TaskId, vdm::task::Decision) = 0;
|
||||
virtual void refresh_url(vdm::TaskId, const std::string& url) = 0;
|
||||
|
||||
// The daemon is done with this task (it went terminal). Drop the handle. Idempotent.
|
||||
virtual void release(vdm::TaskId) = 0;
|
||||
|
||||
// ADR 0011 admission surface. Values are DAEMON's; enforcement is the engine's.
|
||||
virtual void set_task_order(const std::vector<vdm::TaskId>& order) = 0;
|
||||
virtual void set_max_active_segments(std::uint32_t n) = 0;
|
||||
|
||||
@@ -0,0 +1,68 @@
|
||||
#pragma once
|
||||
|
||||
// The real EnginePort: forwards to a live vdm::Engine and keeps the DownloadHandle per
|
||||
// task so pause/resume/cancel/… have something to call. All methods run on the RPC loop
|
||||
// thread (the Scheduler's thread); the handles map is only touched there.
|
||||
|
||||
#include <unordered_map>
|
||||
|
||||
#include "sched/engine_port.hpp"
|
||||
#include "vdm/engine.hpp"
|
||||
#include "vdm/task/download.hpp"
|
||||
|
||||
namespace velox::daemon::sched {
|
||||
|
||||
class EnginePortCore final : public EnginePort {
|
||||
public:
|
||||
explicit EnginePortCore(vdm::Engine& engine) : engine_(engine) {}
|
||||
|
||||
vdm::TaskId start(const vdm::task::DownloadSpec& spec,
|
||||
vdm::task::DownloadCallbacks callbacks) override {
|
||||
vdm::task::DownloadHandle h = engine_.start(spec, std::move(callbacks));
|
||||
const vdm::TaskId id = h.id();
|
||||
handles_.insert_or_assign(id, std::move(h));
|
||||
return id;
|
||||
}
|
||||
|
||||
void pause(vdm::TaskId id) override {
|
||||
if (auto* h = find(id)) h->pause();
|
||||
}
|
||||
void resume(vdm::TaskId id) override {
|
||||
if (auto* h = find(id)) h->resume();
|
||||
}
|
||||
void cancel(vdm::TaskId id, bool discard_partial) override {
|
||||
if (auto* h = find(id)) h->cancel(discard_partial);
|
||||
}
|
||||
void provide_auth(vdm::TaskId id, const std::string& u, const std::string& p,
|
||||
bool remember) override {
|
||||
if (auto* h = find(id)) h->provide_auth(u, p, remember);
|
||||
}
|
||||
void decide(vdm::TaskId id, vdm::task::Decision d) override {
|
||||
if (auto* h = find(id)) h->decide(d);
|
||||
}
|
||||
void refresh_url(vdm::TaskId id, const std::string& url) override {
|
||||
if (auto* h = find(id)) h->refresh_url(url);
|
||||
}
|
||||
void release(vdm::TaskId id) override { handles_.erase(id); }
|
||||
|
||||
void set_task_order(const std::vector<vdm::TaskId>& order) override {
|
||||
engine_.segment_budget().set_task_order(order);
|
||||
}
|
||||
void set_max_active_segments(std::uint32_t n) override {
|
||||
engine_.segment_budget().set_max_active_segments(n);
|
||||
}
|
||||
void set_host_segment_cap(const std::string& host, std::uint32_t cap) override {
|
||||
engine_.segment_budget().set_host_segment_cap(host, cap);
|
||||
}
|
||||
|
||||
private:
|
||||
vdm::task::DownloadHandle* find(vdm::TaskId id) {
|
||||
auto it = handles_.find(id);
|
||||
return it == handles_.end() ? nullptr : &it->second;
|
||||
}
|
||||
|
||||
vdm::Engine& engine_;
|
||||
std::unordered_map<vdm::TaskId, vdm::task::DownloadHandle> handles_;
|
||||
};
|
||||
|
||||
} // namespace velox::daemon::sched
|
||||
@@ -24,6 +24,7 @@ public:
|
||||
std::vector<vdm::TaskId> paused;
|
||||
std::vector<vdm::TaskId> resumed;
|
||||
std::vector<std::pair<vdm::TaskId, bool>> cancelled;
|
||||
std::vector<vdm::TaskId> released;
|
||||
std::vector<std::vector<vdm::TaskId>> orders;
|
||||
std::vector<std::uint32_t> max_active_segments;
|
||||
std::vector<std::pair<std::string, std::uint32_t>> host_caps;
|
||||
@@ -40,6 +41,7 @@ public:
|
||||
void provide_auth(vdm::TaskId, const std::string&, const std::string&, bool) override {}
|
||||
void decide(vdm::TaskId, vdm::task::Decision) override {}
|
||||
void refresh_url(vdm::TaskId, const std::string&) override {}
|
||||
void release(vdm::TaskId id) override { released.push_back(id); }
|
||||
void set_task_order(const std::vector<vdm::TaskId>& order) override { orders.push_back(order); }
|
||||
void set_max_active_segments(std::uint32_t n) override { max_active_segments.push_back(n); }
|
||||
void set_host_segment_cap(const std::string& h, std::uint32_t c) override {
|
||||
|
||||
@@ -307,7 +307,10 @@ void Scheduler::on_engine_state(const std::string& wire_id, std::string_view eng
|
||||
}
|
||||
|
||||
if (engine_state == "complete" || engine_state == "failed" || engine_state == "cancelled") {
|
||||
if (auto eid = engine_id_of(wire_id)) unmap_engine(*eid);
|
||||
if (auto eid = engine_id_of(wire_id)) {
|
||||
engine_.release(*eid);
|
||||
unmap_engine(*eid);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user