The two items aimed at pointing GUI at a real veloxd instead of mockd.
rpc/event_hub — per-subscription fan-out shared by both transports.
subscribe() registers a connection with no interest; set_filter()
(session.subscribe, replaces not adds) turns on event kinds and an
optional per-task id filter; publish() delivers a pre-built
notification to every matching subscriber. session.subscribe on both
UdsServer and WsServer now does the real thing — registers/updates a
subscription, tears it down in close_conn.
sched/scheduler — the on_engine_state hook now actually publishes:
- transition() is the one place a task's row changes state; it reads
the store's own prior row for previousState (authoritative
regardless of engine/scheduler timing), writes the error columns,
and — when a hub is supplied — publishes event.task.state with
{taskId, state, previousState, summary, error}. Wired into every
transition: scheduler-driven (admission -> probing, resume ->
connecting, pause) and engine-reported (on_engine_state).
- progress_snapshot(): one row per task the engine is tracking
(EnginePort::progress(), a new interface method backed by
DownloadHandle::progress()), plus a store side-effect
(Tasks::update_progress) so download.list/get stay current between
state transitions. Returns rows; does NOT publish itself — batching
into one array message is the caller's job, per the schema's
x-maxRateHz: 4 and AGENT-DAEMON.md item 5 ("one message per task per
tick burns a core"). main.cpp's 250 ms timerfd is that caller: one
event.task.progress per tick, only when there's something to say.
dispatcher::on_download_add now publishes event.task.added (schema:
"summary is always present so a client can insert the row without a
follow-up download.get").
store/categories, store/queues — the two D3 handlers GUI's panels
call. category.list projects the six seeded built-ins; queue.list
derives taskIds from tasks.queue_id/queue_position (Queue's own schema
note: a queue's stored row never carries membership, download.update
/ queue.reorder do).
Verified live end to end against tools/testserver: a subscribed client
sees event.task.added on add, then the full event.task.state sequence
(queued -> probing -> connecting -> downloading -> assembling ->
verifying -> complete) with correct previousState at every step, and
real category.list / queue.list results.
Tests: event_hub (filter-by-kind, filter-by-task-id, replace-not-add,
unsubscribe), store_categories_queues, plus new sched_scheduler cases
for event.task.state publishing and progress_snapshot's store
side-effect. 38 daemon/cli tests green; sched_scheduler / event_hub /
ws_server / uds_roundtrip TSan-clean.
deferrals.md: D5 mostly closed (event.task.removed and the
still-unpublished events wait on their owning D3 handlers); D3 down to
the remaining download.* verbs, rules/settings/limiter/schedule,
queue mutation, category mutation, grabber, media.
Co-Authored-By: Claude Sonnet 5 <[email protected]>
Claude-Session: https://claude.ai/code/session_01Upd9WhG9oppieig5nRDLig
83 lines
2.8 KiB
C++
83 lines
2.8 KiB
C++
#pragma once
|
|
|
|
// The Unix-domain-socket RPC listener: $XDG_RUNTIME_DIR/velox/velox.sock, mode 0600,
|
|
// SO_PEERCRED same-UID check (docs/01 §2, AGENT-DAEMON.md build step 1). NDJSON framing.
|
|
// Non-blocking throughout; every fd runs through the shared EventLoop so one slow client
|
|
// never stalls another.
|
|
//
|
|
// session.hello and session.subscribe are handled here because they are connection- and
|
|
// transport-stateful (protocol-major check, sessionId, per-connection subscription set).
|
|
// Every other method is routed through the generated velox::proto::dispatch(), which does
|
|
// the envelope, the -32003 privileged-transport refusal and the typed param parse.
|
|
|
|
#include <cstdint>
|
|
#include <memory>
|
|
#include <string>
|
|
#include <optional>
|
|
#include <system_error>
|
|
#include <unordered_map>
|
|
#include <vector>
|
|
|
|
#include <nlohmann/json_fwd.hpp>
|
|
|
|
#include "rpc/event_hub.hpp"
|
|
#include "rpc/ndjson.hpp"
|
|
#include "velox_proto.hpp"
|
|
|
|
namespace velox::daemon::rpc {
|
|
|
|
class EventLoop;
|
|
|
|
class UdsServer {
|
|
public:
|
|
UdsServer(EventLoop& loop, velox::proto::Dispatcher& dispatcher, EventHub& hub,
|
|
std::string socket_path);
|
|
~UdsServer();
|
|
|
|
UdsServer(const UdsServer&) = delete;
|
|
UdsServer& operator=(const UdsServer&) = delete;
|
|
|
|
// Create the socket, bind, chmod 0600, listen, and register with the loop. A stale
|
|
// socket file left by a crashed daemon is removed first. Returns a non-ok error_code
|
|
// (and changes nothing) on any failure.
|
|
std::error_code start();
|
|
|
|
const std::string& socket_path() const noexcept { return path_; }
|
|
std::size_t connection_count() const noexcept { return conns_.size(); }
|
|
|
|
private:
|
|
struct Conn {
|
|
int fd;
|
|
FrameReader reader;
|
|
std::string outbuf;
|
|
std::size_t out_off = 0; // bytes of outbuf already written
|
|
bool close_after_flush = false;
|
|
bool hello_ok = false;
|
|
std::string session_id;
|
|
std::optional<EventHub::SubId> sub_id;
|
|
};
|
|
|
|
void on_listener_readable();
|
|
void on_conn_event(int fd, unsigned events);
|
|
void handle_line(Conn& c, const std::string& line);
|
|
|
|
// Returns true and fills `reply` if `method` is one this layer answers directly
|
|
// (session.hello / session.subscribe). Returns false to let dispatch() handle it.
|
|
bool handle_session_method(Conn& c, const std::string& method, const nlohmann::json& request,
|
|
nlohmann::json& reply);
|
|
|
|
void queue_reply(Conn& c, const nlohmann::json& reply);
|
|
void flush(Conn& c);
|
|
void close_conn(int fd);
|
|
|
|
EventLoop& loop_;
|
|
velox::proto::Dispatcher& dispatcher_;
|
|
EventHub& hub_;
|
|
std::string path_;
|
|
int listen_fd_ = -1;
|
|
bool bound_ = false; // path_ is ours to unlink on destruction
|
|
std::unordered_map<int, std::unique_ptr<Conn>> conns_;
|
|
};
|
|
|
|
} // namespace velox::daemon::rpc
|