From e442c991303dad6b51d1017a5f18c648eed1272d Mon Sep 17 00:00:00 2001 From: sami Date: Wed, 9 Sep 2026 23:21:10 +0400 Subject: [PATCH 1/2] core: sign off on ADR 0011 with three amendments MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Accept the decision (each ceiling enforced once by its unit owner; DAEMON counts tasks, CORE counts segments; the single min() clamp; per-host caps split by unit). daemon/src/sched/ is unblocked. Answers to the five open questions: 1. Min-1 is implementable in the allocator without inversion: guarantee pass (zero-slot tasks, priority order) before growth pass; a released slot always re-enters allocation from the top, never handed back locally. 2. Probe pool size 4, outside the segment budget — confirmed. 3. set_max_active_segments is drain-not-kill; in-flight segments run to their boundary. 4. Priority = an ordered TaskId list pushed on change, not an integer, not per tick — tie-breaking is DAEMON policy. 5. on_budget_changed coalesced at 4 Hz, immediate on the tasks_starved zero-crossing. Amendments: (A1) add "yield" to §3.3 as the non-neutral slot-transfer op that satisfies min-1 when the budget is full — "steal" stays slot-neutral; (A2) "admission implies progress" is bounded-delay not immediate — bounded by an incumbent's next segment boundary, capped by the stall timeout, so §3.6's 2 s assertion window is too tight; (A3) add starved_tasks() / starved_since() and pin the segments_active / tasks_starved definitions (a connecting segment counts as a held slot and is not starvation). Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01HPPSGhiArbvQgwC2DNiURS --- core/docs/adr-0011-core-response.md | 177 ++++++++++++++++++++++++++++ 1 file changed, 177 insertions(+) create mode 100644 core/docs/adr-0011-core-response.md diff --git a/core/docs/adr-0011-core-response.md b/core/docs/adr-0011-core-response.md new file mode 100644 index 0000000..05678d5 --- /dev/null +++ b/core/docs/adr-0011-core-response.md @@ -0,0 +1,177 @@ +# CORE response to ADR 0011 (admission control & the segment budget) + +**Verdict: accept the decision, with three amendments and one caveat on §3.5.** None of +the amendments change DAEMON's task-unit model — `daemon/src/sched/` is unblocked. The +caveat tightens what "admission implies progress" can promise. + +Signing off on §1 (each ceiling enforced once, by unit owner), §2 (the single clamp), +§4 (per-host caps split by unit from one table), §6 (contract gap / daemon-local stopgap), +and the three rejected alternatives — the shared-semaphore rejection especially. + +--- + +## Answers to the five open questions + +### Q1 — Min-1 before seconds: implementable in the stealer without inversion? **Yes.** + +Slot allocation is a pure function of `(running tasks, their held-slot counts, their +effective per-task caps, DAEMON's priority order, total budget)`. It runs to completion on +every budget-changing edge (segment complete / start / pause / fail, task admit / pause / +resume, a steal, a cap change), in two ordered passes: + +1. **Guarantee pass.** Walk tasks holding **zero** slots *in DAEMON's priority order*; give + each one slot while the pool is non-empty. +2. **Growth pass.** While the pool is non-empty, round-robin over running tasks in priority + order, granting one slot to any task below its effective cap + (`min(spec.segments ?? maxSegmentsPerDownload, host_cap, resumable ? ∞ : 1)`), until the + pool empties or nobody wants more. + +The inversion DAEMON is worried about at slot release does not occur, because **a released +slot always goes to the pool and the allocator re-runs from pass 1** — it is never handed +back locally to the releasing task. If task A's segment finishes and both A and a +lower-priority C now hold zero, the guarantee pass processes the zero-slot set *in priority +order*, so A gets it back, not C. A task that drops to zero re-enters the guarantee queue +at its priority position, not the back. + +Liveness of min-1 depends on DAEMON honouring §2's clamp — never running more tasks than +`maxActiveSegments`. If it admits 3 running tasks against a budget of 2, one starves by +construction and no fairness rule fixes it. §2 already says this; calling it out because +it is load-bearing for Q1. + +### Q2 — Probe pool outside the segment budget, size 4: **confirmed.** + +CORE's probe path is a dedicated pool, independent of `maxActiveSegments`. A probe is a +`Request` with `method = HEAD` (or `range = {0,0}`) whose `on_head` returns `abort`; it +never allocates a transfer slot. Default pool size 4, exposed as +`set_probe_pool_size(uint32_t)`. Probe cancellation is immediate (a timed-out probe frees +its pool slot at once), so `capture.offer` answering `ignore`-first-probe-after works. +Probes still open real sockets — DAEMON should still bound probe *submission* on its side; +CORE bounds *concurrency*, not queue depth. + +### Q3 — `set_max_active_segments()` live-apply: **drain, never kill.** + +- **Raise:** new slots available immediately; allocator runs; starved then growing tasks + take them. +- **Lower:** in-flight segments run to their next boundary. No new segment starts while + `active > new_ceiling`. Completed segments' slots are withheld until + `active ≤ new_ceiling`. Nothing is aborted, no partial range is lost. +- **Edge:** if `new_ceiling < running_task_count`, CORE honours min-1 for the top + `new_ceiling` tasks in priority order; the remainder are held at zero slots and reported + in `tasks_starved`. CORE does **not** auto-pause them — that's policy. After a live + lower, DAEMON must reconcile its running set against the new clamp (pause the + lowest-priority excess). + +### Q4 — Priority shape: **an ordered list, pushed on change — not an integer, not per tick.** + +`set_task_order(std::span)`, called by DAEMON whenever the order changes +(admit, remove, reorder, priority edit, queue switch). CORE caches it and the allocator +walks it. An integer priority would force CORE to implement tie-breaking (FIFO by +admission time, queue precedence) — that's DAEMON policy and DAEMON already has the total +order. Not per tick: an unchanged order isn't re-sent. A running task absent from the list +sorts last (shouldn't happen; defensive). + +### Q5 — Budget-change callback coalescing: **4 Hz for counts, immediate on the edge.** + +`on_budget_changed` coalesced at ≤4 Hz to match `event.task.progress` — DAEMON isn't +making sub-250 ms admission decisions. **Exception:** fire immediately when +`tasks_starved` crosses 0→non-zero or non-zero→0, so DAEMON's invariant check and any UI +reaction see that transition without up to 250 ms of lag. + +--- + +## Amendments requested to the ADR + +### A1 — §3.3 wording: distinguish *steal* from *yield* + +"A steal is slot-neutral" is true for the steal the ADR means (a worker that finished its +range takes the tail of the largest remaining range — the same worker, the same slot). +Min-1 also needs a second, non-neutral operation: + +- **Yield** — the allocator marks an over-quota task to release one slot *at its next + segment boundary*. When that segment completes the slot goes to the pool → guarantee + pass → starved task. It is a slot transfer, not slot-neutral, and it is bounded by the + yielding segment's remaining bytes (never a mid-segment kill). + +Please add "yield" to §3.3 as the mechanism that satisfies §3.1 when the budget is full. +"Steal" stays exactly as written. + +### A2 — §3.5 / §3.6: "admission implies progress" is bounded-delay, not immediate + +When the budget is full of healthy incumbents, a newly admitted task's first slot +materialises only when some incumbent segment reaches a boundary (yield) — bounded by that +segment's remaining bytes, hard-capped by the stall timeout (`low_speed_secs`, default +30 s, after which a stalled segment fails and frees its slot). So the true bound is + +``` +time_to_first_slot ≤ min(incumbent's next boundary, low_speed_secs) + connect_timeout +``` + +not "connect timeout + per-host cap" alone. Two consequences: + +- §3.6's **2 s** assertion window is too tight — a legitimately full budget with a slow + incumbent tail can hold a new task at zero for longer than 2 s with nothing wrong. + Recommend the warning threshold be `low_speed_secs + connect_timeout_ms` (~45 s), or + configurable. +- CORE will expose `starved_since` (a monotonic timestamp) per starved task via the + diagnostics call below, so DAEMON can tell "briefly waiting for a boundary" from + "wedged" without guessing. + +Optional future tightening (not M1): a **preemptive split** — truncate an incumbent's +largest remaining range ahead of its current offset and hand the freed tail to the starved +task as a new segment. Zero bytes lost, first slot within one round-trip. CORE will add +this if the yield delay proves painful in the soak test; it doesn't change the API. + +### A3 — C1 API: add a starved-set accessor and pin two definitions + +``` +struct EngineBudget { uint32_t total; uint32_t active; uint32_t tasks_starved; }; +EngineBudget budget() const; +uint32_t segments_active(TaskId) const; // slots held, any state +std::vector starved_tasks() const; // diagnostics / velox ls --json +std::optional starved_since(TaskId) const; // per A2 +void on_budget_changed(std::function); +``` + +Definitions, so the projection to `TaskSummary.segments` (ADR 0010: effective count) is +unambiguous: + +- **`segments_active(id)`** = slots the task holds, counting a segment in `connecting` + (0 bytes yet) as well as `downloading`. This is what the user sees as "using N + connections." +- **`tasks_starved`** counts running tasks with `segments_active == 0`. A task with a + `connecting` segment is **not** starved — it is progressing. + +--- + +## C2–C6 confirmations + +- **C2:** drain-not-kill, per Q3. Confirmed as DAEMON assumed. +- **C3:** `set_host_segment_cap(std::string host, uint32_t)` — confirmed. CORE keeps the + `host → cap` map, applies it in the per-task effective cap and in the stealer (no Nth + connection to a host capped at N-1). DAEMON owns the table; CORE derives a task's host + from its URL + mirror set. +- **C4:** agreed it's PROTO's, same bundle as CORE's B2a / `buffer-sizing.md` asks + (`connection.maxActiveSegments`, `connection.maxTotalBufferBytes`). DAEMON's local-value + stopgap is fine. +- **C5:** signed off — see Q1 plus amendments A1/A2. Min-1 is buildable; "implies + progress" is bounded-delay; steal stays slot-neutral, yield is the transfer op. +- **C6:** confirmed — see Q2. + +--- + +## New API surface CORE will expose for `sched/` (summary) + +``` +void set_max_active_segments(uint32_t); // drain-not-kill +void set_host_segment_cap(std::string host, uint32_t); +void set_task_order(std::span); // pushed on change +void set_probe_pool_size(uint32_t); // default 4 +EngineBudget budget() const; +uint32_t segments_active(TaskId) const; +std::vector starved_tasks() const; +std::optional starved_since(TaskId) const; +void on_budget_changed(std::function); // 4 Hz + starved edge +``` + +All of it lands with stage 6 (segmenter/stealer) / stage 8 (download_task). None of it is +on the M1 critical path ahead of where DAEMON needs it; flag if the ordering is wrong. From bd1bc029f3001d62aa126bf0044aceb379410888 Mon Sep 17 00:00:00 2001 From: sami Date: Wed, 9 Sep 2026 23:27:38 +0400 Subject: [PATCH 2/2] =?UTF-8?q?core:=20net/http=5Fclient=20=E2=80=94=20lib?= =?UTF-8?q?curl=20multi=20wrapper=20(stage=202)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit One HttpClient owns a small pool of workers, each with its own CURLM; an easy handle lives on one worker for its life. Public start/pause/resume/ cancel enqueue a command + curl_multi_wakeup(); callbacks (on_head / on_data / on_finished) run on the worker thread and return a DataAction (proceed / pause / abort). Covers redirects (final-response head only), ranges (inclusive ByteRange -> CURLOPT_RANGE), proxy/SOCKS5, basic/digest auth, cookies, verbatim headers, stall detection, a coarse recv-rate cap, and a curl_share DNS/TLS cache across workers. CURLcode + HTTP status -> vdm::Error in net/curl_error. A probe is on_head returning abort: it finishes successfully (head_complete), not canceled. Tests drive tools/testserver: full GET, ranged 206, redirect chain, 404 -> not_found, connection refused -> connect_failed, HEAD probe + ranged 0-0 probe (no body), cancel mid-transfer, pause/resume completes. Skip cleanly if testserver isn't in the tree. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01HPPSGhiArbvQgwC2DNiURS --- core/CMakeLists.txt | 12 +- core/include/vdm/net/http_client.hpp | 125 ++++++ core/include/vdm/net/http_types.hpp | 156 +++++++ core/src/net/curl_error.cpp | 89 ++++ core/src/net/curl_error.hpp | 24 ++ core/src/net/http_client.cpp | 570 ++++++++++++++++++++++++++ core/tests/CMakeLists.txt | 14 + core/tests/net/http_client_test.cpp | 270 ++++++++++++ core/tests/net/testserver_fixture.hpp | 113 +++++ 9 files changed, 1369 insertions(+), 4 deletions(-) create mode 100644 core/include/vdm/net/http_client.hpp create mode 100644 core/include/vdm/net/http_types.hpp create mode 100644 core/src/net/curl_error.cpp create mode 100644 core/src/net/curl_error.hpp create mode 100644 core/src/net/http_client.cpp create mode 100644 core/tests/net/http_client_test.cpp create mode 100644 core/tests/net/testserver_fixture.hpp diff --git a/core/CMakeLists.txt b/core/CMakeLists.txt index affb361..299a4f2 100644 --- a/core/CMakeLists.txt +++ b/core/CMakeLists.txt @@ -5,16 +5,20 @@ # `add_subdirectory(core)` in the root CMakeLists.txt (see core/docs/pkg-requests-m1.md). find_package(Threads REQUIRED) +find_package(CURL 8.0 REQUIRED) add_library(veloxcore STATIC src/util/error.cpp src/util/log.cpp src/util/thread_pool.cpp + src/net/curl_error.cpp + src/net/http_client.cpp ) add_library(velox::core ALIAS veloxcore) -target_include_directories(veloxcore PUBLIC - ${CMAKE_CURRENT_SOURCE_DIR}/include +target_include_directories(veloxcore + PUBLIC ${CMAKE_CURRENT_SOURCE_DIR}/include + PRIVATE ${CMAKE_CURRENT_SOURCE_DIR}/src # net/*.cpp -> "net/curl_error.hpp" ) target_compile_features(veloxcore PUBLIC cxx_std_23) @@ -25,9 +29,9 @@ target_compile_options(veloxcore PRIVATE -Wall -Wextra -Wpedantic -Werror ) -target_link_libraries(veloxcore PUBLIC Threads::Threads) +target_link_libraries(veloxcore PUBLIC Threads::Threads CURL::libcurl) -# Later stages add: find_package(CURL 8.0) for net/, find_package(OpenSSL) for meta/. +# Later stages add: find_package(OpenSSL) for meta/ (streaming SHA-256 + resume CRC). if(VELOX_BUILD_TESTS) add_subdirectory(tests) diff --git a/core/include/vdm/net/http_client.hpp b/core/include/vdm/net/http_client.hpp new file mode 100644 index 0000000..f213f03 --- /dev/null +++ b/core/include/vdm/net/http_client.hpp @@ -0,0 +1,125 @@ +// vdm/net/http_client.hpp — a libcurl-multi wrapper for the download engine. +// +// One HttpClient owns a small pool of worker threads, each with its own curl_multi +// (curl handles are not thread-safe; a handle lives on exactly one worker for its life). +// The engine hands it a Request plus callbacks and gets back a Transfer handle it can +// pause / resume / cancel from any thread. +// +// Layering: this is the ONLY core header that pulls in libcurl, and only in its .cpp — +// nothing here exposes a curl type. No JSON / SQL / Qt / RPC (CLAUDE.md §3). +// +// Callbacks run ON THE WORKER THREAD, one transfer at a time for a given Transfer. +// They must not block (that stalls every other transfer on that worker) and must not +// call back into this Transfer's pause/resume/cancel re-entrantly — post that work +// elsewhere. on_data must not allocate on the hot path (AGENT-CORE); the ring buffer it +// writes into is preallocated by the caller (stage 4). +// +// This header compiles standalone. + +#ifndef VDM_NET_HTTP_CLIENT_HPP +#define VDM_NET_HTTP_CLIENT_HPP + +#include +#include +#include +#include +#include + +#include "vdm/net/http_types.hpp" +#include "vdm/util/bytes.hpp" +#include "vdm/util/result.hpp" + +namespace vdm::net { + +// What on_data tells the client to do with the transfer after this chunk. +enum class DataAction { + proceed, // keep receiving + pause, // stop receiving; resumes on Transfer::resume() + abort, // end the transfer now (finishes with Error::canceled) +}; + +struct TransferStats { + std::uint64_t bytes_received = 0; + long http_status = 0; + std::string effective_url; + // Timings in milliseconds (from CURLINFO_*_TIME_T), 0 if the phase didn't happen. + long namelookup_ms = 0; + long connect_ms = 0; + long appconnect_ms = 0; // TLS handshake done + long starttransfer_ms = 0; // first response byte + long total_ms = 0; +}; + +struct TransferCallbacks { + // Response headers are in. Called at most once per transfer (a followed redirect's + // intermediate headers are not delivered). If the caller only wanted headers (a + // probe), return DataAction::abort here. + std::function on_head; + + // A chunk of body bytes. The span is valid only for the duration of the call. + std::function on_data; + + // The transfer ended — success carries stats, failure carries the mapped error + // (ErrorInfo.http_status is set for HTTP-status failures). Always called exactly once, + // last. + std::function)> on_finished; +}; + +class HttpClient; + +// Lightweight handle to a running transfer. Copyable (shared state). All methods are +// safe to call from any thread; they post to the owning worker and return immediately. +// Dropping the last handle does NOT cancel — call cancel() for that. +class Transfer { + public: + Transfer() = default; + + [[nodiscard]] std::uint64_t id() const noexcept; + [[nodiscard]] bool valid() const noexcept { return static_cast(state_); } + + void pause(); // no-op if already paused / finished + void resume(); // no-op if not paused / finished + void cancel(); // idempotent; on_finished fires with Error::canceled + + private: + friend class HttpClient; + struct State; + explicit Transfer(std::shared_ptr s) : state_(std::move(s)) {} + std::shared_ptr state_; +}; + +class HttpClient { + public: + struct Options { + // Worker threads, each with its own curl_multi. New transfers are assigned + // round-robin. 0 => pick from hardware_concurrency (min 1, max 4). + unsigned workers = 0; + // CURLMOPT_MAX_TOTAL_CONNECTIONS per worker (0 = curl default). + long max_connections_per_worker = 0; + // Shared DNS + TLS-session cache across this client's workers (curl_share). + bool share_dns_and_tls = true; + }; + + HttpClient(); // default Options + explicit HttpClient(Options opts); + ~HttpClient(); + + HttpClient(const HttpClient &) = delete; + HttpClient &operator=(const HttpClient &) = delete; + + [[nodiscard]] unsigned worker_count() const noexcept; + + // Start a transfer. The Request is consumed (moved). Returns an invalid Transfer and + // never calls the callbacks only if the client is shutting down; every other failure + // (bad URL, DNS, ...) is delivered through on_finished. + Transfer start(Request req, TransferCallbacks cbs); + + private: + friend class Transfer; // Transfer posts pause/resume/cancel commands to Impl + struct Impl; + std::unique_ptr impl_; +}; + +} // namespace vdm::net + +#endif // VDM_NET_HTTP_CLIENT_HPP diff --git a/core/include/vdm/net/http_types.hpp b/core/include/vdm/net/http_types.hpp new file mode 100644 index 0000000..e476b77 --- /dev/null +++ b/core/include/vdm/net/http_types.hpp @@ -0,0 +1,156 @@ +// vdm/net/http_types.hpp — value types for HTTP requests and responses. +// +// No libcurl in this header: it is the vocabulary the rest of core/ speaks to the net +// layer. http_client.hpp is the only place curl leaks in, and only in its .cpp. +// +// This header compiles standalone. + +#ifndef VDM_NET_HTTP_TYPES_HPP +#define VDM_NET_HTTP_TYPES_HPP + +#include +#include +#include +#include +#include +#include + +namespace vdm::net { + +enum class Method { get, head }; + +[[nodiscard]] constexpr std::string_view method_name(Method m) noexcept { + return m == Method::head ? "HEAD" : "GET"; +} + +// A closed byte range [first, last]. `last == kUnbounded` means "to the end of the +// resource" (`Range: bytes=first-`). This mirrors HTTP Range semantics — inclusive on +// both ends — deliberately: see core/docs/proto-requests-m1.md B3. +struct ByteRange { + static constexpr std::uint64_t kUnbounded = ~std::uint64_t{0}; + + std::uint64_t first = 0; + std::uint64_t last = kUnbounded; + + [[nodiscard]] bool bounded() const noexcept { return last != kUnbounded; } + [[nodiscard]] std::optional length() const noexcept { + if (!bounded()) + return std::nullopt; + return last - first + 1; + } + // "bytes=100-199" or "bytes=100-" + [[nodiscard]] std::string to_header_value() const { + std::string v = "bytes="; + v += std::to_string(first); + v += '-'; + if (bounded()) + v += std::to_string(last); + return v; + } +}; + +struct HeaderField { + std::string name; + std::string value; +}; + +// Case-insensitive view over response headers. Not a multimap for perf — the header set +// on a download response is tiny; linear scan is fine and keeps this allocation-light. +class HeaderList { + public: + void add(std::string name, std::string value) { + fields_.push_back({std::move(name), std::move(value)}); + } + void clear() noexcept { fields_.clear(); } + + [[nodiscard]] const std::vector &fields() const noexcept { return fields_; } + [[nodiscard]] bool empty() const noexcept { return fields_.empty(); } + + // First value for `name` (ASCII case-insensitive), or nullopt. + [[nodiscard]] std::optional get(std::string_view name) const { + for (const auto &f : fields_) + if (iequals(f.name, name)) + return f.value; + return std::nullopt; + } + [[nodiscard]] bool has(std::string_view name) const { return get(name).has_value(); } + + static bool iequals(std::string_view a, std::string_view b) noexcept { + if (a.size() != b.size()) + return false; + for (std::size_t i = 0; i < a.size(); ++i) + if (lower(a[i]) != lower(b[i])) + return false; + return true; + } + + private: + static constexpr char lower(char c) noexcept { + return (c >= 'A' && c <= 'Z') ? char(c - 'A' + 'a') : c; + } + std::vector fields_; +}; + +enum class ProxyKind { none, http, socks5, socks5_hostname }; + +struct ProxyConfig { + ProxyKind kind = ProxyKind::none; + std::string host; // host[:port] + std::uint16_t port = 0; + std::string username; // empty = no proxy auth + std::string password; +}; + +enum class AuthScheme { none, basic, digest, any }; + +struct AuthConfig { + AuthScheme scheme = AuthScheme::none; + std::string username; + std::string password; +}; + +struct Cookie { + std::string name; + std::string value; +}; + +// One HTTP request the engine wants performed. Defaults are the well-behaved case. +struct Request { + std::string url; + Method method = Method::get; + std::vector headers; // verbatim; browser UA/Referer/cookies live here + std::optional range; + std::vector cookies; // merged into a Cookie: header + curl's jar + std::string user_agent; // convenience; also settable via headers + std::string referrer; + + ProxyConfig proxy; + AuthConfig auth; + + bool follow_redirects = true; + long max_redirects = 20; + bool accept_encoding = false; // OFF for downloads: a gzip'd body breaks Range math + + // Stall detection: abort if throughput stays under `low_speed_bytes_per_sec` for + // `low_speed_secs`. 0 disables. Distinct from an overall deadline (probe sets one). + long connect_timeout_ms = 15000; + long low_speed_bytes_per_sec = 1024; + long low_speed_secs = 30; + long overall_timeout_ms = 0; // 0 = none; probe uses a few seconds + + // Coarse download-rate ceiling handed to curl (CURLOPT_MAX_RECV_SPEED_LARGE). The + // precise limiter (stage 7) pauses/resumes on top of this. 0 = unlimited. + std::uint64_t max_recv_bytes_per_sec = 0; +}; + +// Delivered once, when response headers are in. +struct ResponseHead { + long status = 0; + std::string effective_url; // after redirects + HeaderList headers; + std::optional content_length; // from Content-Length, if present & sane +}; + +} // namespace vdm::net + +#endif // VDM_NET_HTTP_TYPES_HPP diff --git a/core/src/net/curl_error.cpp b/core/src/net/curl_error.cpp new file mode 100644 index 0000000..ebdac79 --- /dev/null +++ b/core/src/net/curl_error.cpp @@ -0,0 +1,89 @@ +// vdm/net/curl_error.cpp + +#include "net/curl_error.hpp" + +namespace vdm::net::detail { + +Error error_from_curl(CURLcode code, long http_status) noexcept { + // Transport-level failures first — these override any status. + switch (code) { + case CURLE_OK: + break; + case CURLE_COULDNT_RESOLVE_PROXY: + case CURLE_COULDNT_RESOLVE_HOST: + return Error::resolve_failed; + case CURLE_COULDNT_CONNECT: + case CURLE_INTERFACE_FAILED: + return Error::connect_failed; + case CURLE_OPERATION_TIMEDOUT: + return Error::timeout; + case CURLE_TOO_MANY_REDIRECTS: + return Error::too_many_redirects; + case CURLE_PEER_FAILED_VERIFICATION: + case CURLE_SSL_CONNECT_ERROR: + case CURLE_SSL_CERTPROBLEM: + case CURLE_SSL_CIPHER: + case CURLE_SSL_CACERT_BADFILE: + case CURLE_SSL_ISSUER_ERROR: + case CURLE_SSL_PINNEDPUBKEYNOTMATCH: + case CURLE_SSL_INVALIDCERTSTATUS: + return Error::tls_failed; + case CURLE_GOT_NOTHING: + case CURLE_RECV_ERROR: + case CURLE_SEND_ERROR: + case CURLE_PARTIAL_FILE: + case CURLE_HTTP2: + case CURLE_HTTP2_STREAM: + return Error::connection_reset; + case CURLE_WRITE_ERROR: + // Our write callback returned short — the sink (disk) failed or we're + // cancelling. The caller distinguishes; default to io_error. + return Error::io_error; + case CURLE_LOGIN_DENIED: + return Error::auth_required; + case CURLE_UNSUPPORTED_PROTOCOL: + case CURLE_URL_MALFORMAT: + return Error::unsupported_url_scheme; + case CURLE_ABORTED_BY_CALLBACK: + return Error::canceled; + default: + // Fall through to status-based classification, else generic. + break; + } + + // HTTP status classification (also reached on CURLE_OK). + if (http_status >= 400) { + switch (http_status) { + case 401: + case 407: + return Error::auth_required; + case 403: + return Error::forbidden; + case 404: + return Error::not_found; + case 410: + return Error::gone; + case 416: + return Error::range_not_satisfiable; + default: + return http_status >= 500 ? Error::http_server_error : Error::http_client_error; + } + } + + if (code != CURLE_OK) + return Error::internal; + return Error::ok; +} + +ErrorInfo make_error(CURLcode code, long http_status, const char *curl_msg) { + Error e = error_from_curl(code, http_status); + std::string ctx; + if (curl_msg && *curl_msg) + ctx = curl_msg; + else if (code != CURLE_OK) + ctx = curl_easy_strerror(code); + ErrorInfo info(e, std::move(ctx), static_cast(http_status)); + return info; +} + +} // namespace vdm::net::detail diff --git a/core/src/net/curl_error.hpp b/core/src/net/curl_error.hpp new file mode 100644 index 0000000..9366e6b --- /dev/null +++ b/core/src/net/curl_error.hpp @@ -0,0 +1,24 @@ +// vdm/net/curl_error.hpp — internal: CURLcode -> vdm::Error. Not a public header. + +#ifndef VDM_NET_CURL_ERROR_HPP +#define VDM_NET_CURL_ERROR_HPP + +#include + +#include + +#include "vdm/util/error.hpp" + +namespace vdm::net::detail { + +// Map a libcurl transfer result to the engine taxonomy. `http_status` (0 if none) lets +// the HTTP-status errors (403/404/416/...) be classified here too; pass it from +// CURLINFO_RESPONSE_CODE. `CURLE_OK` with a >= 400 status still yields an error. +[[nodiscard]] Error error_from_curl(CURLcode code, long http_status) noexcept; + +// A human-readable ErrorInfo, folding in curl's own message and the status. +[[nodiscard]] ErrorInfo make_error(CURLcode code, long http_status, const char *curl_msg = nullptr); + +} // namespace vdm::net::detail + +#endif // VDM_NET_CURL_ERROR_HPP diff --git a/core/src/net/http_client.cpp b/core/src/net/http_client.cpp new file mode 100644 index 0000000..25f9575 --- /dev/null +++ b/core/src/net/http_client.cpp @@ -0,0 +1,570 @@ +// vdm/net/http_client.cpp — libcurl multi implementation. +// +// Threading model: one Worker == one std::jthread + one CURLM. An easy handle is created, +// used, paused, and destroyed only on its Worker's thread. Public calls (start / pause / +// resume / cancel) just enqueue a Command and curl_multi_wakeup() the worker. + +#include "vdm/net/http_client.hpp" + +#include + +#include +#include +#include +#include +#include +#include +#include +#include + +#include "net/curl_error.hpp" +#include "vdm/util/log.hpp" + +namespace vdm::net { +namespace { + +struct CurlGlobal { + CurlGlobal() { + if (curl_global_init(CURL_GLOBAL_ALL) != CURLE_OK) + VDM_LOG_ERROR("net", "curl_global_init failed"); + } + ~CurlGlobal() { curl_global_cleanup(); } +}; +void ensure_curl_global() { + static CurlGlobal g; + (void)g; +} + +std::uint64_t next_id() { + static std::atomic counter{0}; + return ++counter; +} + +bool parse_header_line(std::string_view line, std::string &name, std::string &value) { + while (!line.empty() && (line.back() == '\r' || line.back() == '\n')) + line.remove_suffix(1); + if (line.empty()) + return false; + auto colon = line.find(':'); + if (colon == std::string_view::npos) + return false; + name.assign(line.substr(0, colon)); + auto v = line.substr(colon + 1); + while (!v.empty() && (v.front() == ' ' || v.front() == '\t')) + v.remove_prefix(1); + value.assign(v); + return true; +} + +// "HTTP/1.1 206 Partial Content" -> 206; 0 on parse failure. +long status_from_line(std::string_view line) { + auto sp = line.find(' '); + if (sp == std::string_view::npos) + return 0; + auto rest = line.substr(sp + 1); + long code = 0; + auto [p, ec] = std::from_chars(rest.data(), rest.data() + rest.size(), code); + (void)p; + return ec == std::errc{} ? code : 0; +} + +} // namespace + +// --- Transfer::State ----------------------------------------------------------------- +struct Transfer::State { + enum class Stop { none, head_complete, aborted }; + + std::uint64_t id = 0; + struct HttpClient::Impl *client = nullptr; + unsigned worker_index = 0; + + Request req; + TransferCallbacks cbs; + + // Worker-thread-owned. + CURL *easy = nullptr; + curl_slist *header_slist = nullptr; + std::string range_value; + std::string cookie_value; + + ResponseHead head; + long line_status = 0; // status from the most recent HTTP/ line + bool head_delivered = false; + + std::atomic pause_requested{false}; + std::atomic stop{Stop::none}; + bool curl_paused = false; + + std::uint64_t bytes_received = 0; + bool finished = false; +}; + +// --- Impl ------------------------------------------------------------------------- +struct HttpClient::Impl { + enum class CmdKind { add, pause, resume, cancel }; + struct Command { + CmdKind kind; + std::shared_ptr state; + }; + + struct Worker { + CURLM *multi = nullptr; + std::mutex mu; + std::deque queue; + std::vector> live; + std::jthread thread; + }; + + explicit Impl(Options o) : opts(o) { + ensure_curl_global(); + + unsigned n = opts.workers; + if (n == 0) { + unsigned hw = std::thread::hardware_concurrency(); + n = std::clamp(hw ? hw : 1, 1, 4); + } + + if (opts.share_dns_and_tls) { + share = curl_share_init(); + if (share) { + curl_share_setopt(share, CURLSHOPT_LOCKFUNC, &Impl::share_lock); + curl_share_setopt(share, CURLSHOPT_UNLOCKFUNC, &Impl::share_unlock); + curl_share_setopt(share, CURLSHOPT_USERDATA, this); + curl_share_setopt(share, CURLSHOPT_SHARE, CURL_LOCK_DATA_DNS); + curl_share_setopt(share, CURLSHOPT_SHARE, CURL_LOCK_DATA_SSL_SESSION); + } + } + + workers.reserve(n); + for (unsigned i = 0; i < n; ++i) { + auto w = std::make_unique(); + w->multi = curl_multi_init(); + if (opts.max_connections_per_worker > 0) + curl_multi_setopt(w->multi, CURLMOPT_MAX_TOTAL_CONNECTIONS, + opts.max_connections_per_worker); + Worker *raw = w.get(); + w->thread = std::jthread([this, raw](std::stop_token st) { run(*raw, st); }); + workers.push_back(std::move(w)); + } + } + + ~Impl() { + stopping.store(true); + for (auto &w : workers) { + w->thread.request_stop(); + if (w->multi) + curl_multi_wakeup(w->multi); + } + for (auto &w : workers) + if (w->thread.joinable()) + w->thread.join(); + for (auto &w : workers) + if (w->multi) + curl_multi_cleanup(w->multi); + if (share) + curl_share_cleanup(share); + } + + Options opts; + std::vector> workers; + CURLSH *share = nullptr; + std::mutex share_mu[CURL_LOCK_DATA_LAST]; + std::atomic rr{0}; + std::atomic stopping{false}; + + static void share_lock(CURL *, curl_lock_data data, curl_lock_access, void *userp) { + static_cast(userp)->share_mu[data].lock(); + } + static void share_unlock(CURL *, curl_lock_data data, void *userp) { + static_cast(userp)->share_mu[data].unlock(); + } + + void enqueue(unsigned wi, Command cmd) { + Worker &w = *workers[wi]; + { + std::lock_guard lk(w.mu); + w.queue.push_back(std::move(cmd)); + } + curl_multi_wakeup(w.multi); + } + + // ---- curl C callbacks ---- + static std::size_t header_cb(char *buf, std::size_t size, std::size_t n, void *userp) { + auto *st = static_cast(userp); + const std::size_t total = size * n; + std::string_view line(buf, total); + + if (line.starts_with("HTTP/")) { + st->line_status = status_from_line(line); + st->head.headers.clear(); // keep only the final response's headers + return total; + } + if (line == "\r\n" || line == "\n") { + const bool redirect = + st->req.follow_redirects && st->line_status >= 300 && st->line_status < 400; + if (!redirect) + deliver_head(st); + return total; + } + std::string name, value; + if (parse_header_line(line, name, value)) + st->head.headers.add(std::move(name), std::move(value)); + return total; + } + + static void deliver_head(Transfer::State *st) { + if (st->head_delivered) + return; + st->head_delivered = true; + + long code = 0; + curl_easy_getinfo(st->easy, CURLINFO_RESPONSE_CODE, &code); + st->head.status = code ? code : st->line_status; + char *eff = nullptr; + if (curl_easy_getinfo(st->easy, CURLINFO_EFFECTIVE_URL, &eff) == CURLE_OK && eff) + st->head.effective_url = eff; + curl_off_t clen = -1; + if (curl_easy_getinfo(st->easy, CURLINFO_CONTENT_LENGTH_DOWNLOAD_T, &clen) == CURLE_OK && + clen >= 0) + st->head.content_length = static_cast(clen); + + if (st->cbs.on_head) { + DataAction a = st->cbs.on_head(st->head); + if (a == DataAction::abort) + st->stop.store(Transfer::State::Stop::head_complete); + else if (a == DataAction::pause) + st->pause_requested.store(true); + } + } + + static std::size_t write_cb(char *ptr, std::size_t size, std::size_t n, void *userp) { + auto *st = static_cast(userp); + const std::size_t total = size * n; + + if (!st->head_delivered) + deliver_head(st); + if (st->stop.load() != Transfer::State::Stop::none) + return 0; // -> CURLE_WRITE_ERROR + if (st->pause_requested.load()) { + st->curl_paused = true; + return CURL_WRITEFUNC_PAUSE; + } + + if (total && st->cbs.on_data) { + ConstByteSpan span(reinterpret_cast(ptr), total); + DataAction a = st->cbs.on_data(span); + if (a == DataAction::abort) { + st->stop.store(Transfer::State::Stop::aborted); + return 0; + } + if (a == DataAction::pause) { + st->pause_requested.store(true); + st->curl_paused = true; + return CURL_WRITEFUNC_PAUSE; + } + } + st->bytes_received += total; + return total; + } + + // ---- worker thread ---- + void run(Worker &w, std::stop_token stok) { + while (!stok.stop_requested()) { + drain_commands(w); + int running = 0; + curl_multi_perform(w.multi, &running); + reap(w); + if (stok.stop_requested()) + break; + int numfds = 0; + curl_multi_poll(w.multi, nullptr, 0, 1000, &numfds); + } + shutdown_worker(w); + } + + void drain_commands(Worker &w) { + std::deque local; + { + std::lock_guard lk(w.mu); + local.swap(w.queue); + } + for (auto &cmd : local) { + auto &st = cmd.state; + switch (cmd.kind) { + case CmdKind::add: + attach(w, st); + break; + case CmdKind::pause: + st->pause_requested.store(true); + if (st->easy && !st->curl_paused) { + curl_easy_pause(st->easy, CURLPAUSE_RECV); + st->curl_paused = true; + } + break; + case CmdKind::resume: + st->pause_requested.store(false); + if (st->easy && st->curl_paused) { + st->curl_paused = false; + curl_easy_pause(st->easy, CURLPAUSE_CONT); + } + break; + case CmdKind::cancel: + st->stop.store(Transfer::State::Stop::aborted); + if (st->easy && st->curl_paused) { + st->curl_paused = false; + curl_easy_pause(st->easy, CURLPAUSE_CONT); // let write_cb return 0 + } + break; + } + } + } + + void attach(Worker &w, std::shared_ptr st) { + if (stopping.load()) { + complete(st, ErrorInfo(Error::canceled, "client shutting down")); + return; + } + CURL *e = curl_easy_init(); + if (!e) { + complete(st, ErrorInfo(Error::internal, "curl_easy_init")); + return; + } + st->easy = e; + const Request &r = st->req; + + curl_easy_setopt(e, CURLOPT_URL, r.url.c_str()); + curl_easy_setopt(e, CURLOPT_PRIVATE, st.get()); + curl_easy_setopt(e, CURLOPT_NOSIGNAL, 1L); + curl_easy_setopt(e, CURLOPT_NOPROGRESS, 1L); + curl_easy_setopt(e, CURLOPT_HEADERFUNCTION, &Impl::header_cb); + curl_easy_setopt(e, CURLOPT_HEADERDATA, st.get()); + curl_easy_setopt(e, CURLOPT_WRITEFUNCTION, &Impl::write_cb); + curl_easy_setopt(e, CURLOPT_WRITEDATA, st.get()); + curl_easy_setopt(e, CURLOPT_TCP_KEEPALIVE, 1L); + if (share) + curl_easy_setopt(e, CURLOPT_SHARE, share); + + if (r.method == Method::head) + curl_easy_setopt(e, CURLOPT_NOBODY, 1L); + + if (r.follow_redirects) { + curl_easy_setopt(e, CURLOPT_FOLLOWLOCATION, 1L); + curl_easy_setopt(e, CURLOPT_MAXREDIRS, r.max_redirects); + } + curl_easy_setopt(e, CURLOPT_ACCEPT_ENCODING, r.accept_encoding ? "" : nullptr); + + if (r.range) { + std::string v = r.range->to_header_value(); // "bytes=first-last" + std::string_view sv = v; + if (sv.starts_with("bytes=")) + sv.remove_prefix(6); + st->range_value.assign(sv); + curl_easy_setopt(e, CURLOPT_RANGE, st->range_value.c_str()); + } + + curl_easy_setopt(e, CURLOPT_CONNECTTIMEOUT_MS, static_cast(r.connect_timeout_ms)); + if (r.overall_timeout_ms > 0) + curl_easy_setopt(e, CURLOPT_TIMEOUT_MS, static_cast(r.overall_timeout_ms)); + if (r.low_speed_bytes_per_sec > 0 && r.low_speed_secs > 0) { + curl_easy_setopt(e, CURLOPT_LOW_SPEED_LIMIT, r.low_speed_bytes_per_sec); + curl_easy_setopt(e, CURLOPT_LOW_SPEED_TIME, r.low_speed_secs); + } + if (r.max_recv_bytes_per_sec > 0) + curl_easy_setopt(e, CURLOPT_MAX_RECV_SPEED_LARGE, + static_cast(r.max_recv_bytes_per_sec)); + + if (!r.user_agent.empty()) + curl_easy_setopt(e, CURLOPT_USERAGENT, r.user_agent.c_str()); + if (!r.referrer.empty()) + curl_easy_setopt(e, CURLOPT_REFERER, r.referrer.c_str()); + + if (r.proxy.kind != ProxyKind::none) { + curl_easy_setopt(e, CURLOPT_PROXY, r.proxy.host.c_str()); + if (r.proxy.port) + curl_easy_setopt(e, CURLOPT_PROXYPORT, long(r.proxy.port)); + long pt = CURLPROXY_HTTP; + if (r.proxy.kind == ProxyKind::socks5) + pt = CURLPROXY_SOCKS5; + else if (r.proxy.kind == ProxyKind::socks5_hostname) + pt = CURLPROXY_SOCKS5_HOSTNAME; + curl_easy_setopt(e, CURLOPT_PROXYTYPE, pt); + if (!r.proxy.username.empty()) { + std::string up = r.proxy.username + ":" + r.proxy.password; + curl_easy_setopt(e, CURLOPT_PROXYUSERPWD, up.c_str()); + } + } + + if (r.auth.scheme != AuthScheme::none) { + long m = CURLAUTH_ANY; + if (r.auth.scheme == AuthScheme::basic) + m = CURLAUTH_BASIC; + else if (r.auth.scheme == AuthScheme::digest) + m = CURLAUTH_DIGEST; + curl_easy_setopt(e, CURLOPT_HTTPAUTH, m); + std::string up = r.auth.username + ":" + r.auth.password; + curl_easy_setopt(e, CURLOPT_USERPWD, up.c_str()); + } + + if (!r.cookies.empty()) { + for (const auto &c : r.cookies) { + if (!st->cookie_value.empty()) + st->cookie_value += "; "; + st->cookie_value += c.name + "=" + c.value; + } + curl_easy_setopt(e, CURLOPT_COOKIE, st->cookie_value.c_str()); + } + + for (const auto &h : r.headers) { + std::string joined = h.name + ": " + h.value; + st->header_slist = curl_slist_append(st->header_slist, joined.c_str()); + } + if (st->header_slist) + curl_easy_setopt(e, CURLOPT_HTTPHEADER, st->header_slist); + + CURLMcode mc = curl_multi_add_handle(w.multi, e); + if (mc != CURLM_OK) { + curl_easy_cleanup(e); + st->easy = nullptr; + complete(st, ErrorInfo(Error::internal, curl_multi_strerror(mc))); + return; + } + w.live.push_back(std::move(st)); + } + + void reap(Worker &w) { + CURLMsg *msg = nullptr; + int inq = 0; + while ((msg = curl_multi_info_read(w.multi, &inq)) != nullptr) { + if (msg->msg != CURLMSG_DONE) + continue; + CURL *e = msg->easy_handle; + const CURLcode res = msg->data.result; + + Transfer::State *raw = nullptr; + curl_easy_getinfo(e, CURLINFO_PRIVATE, &raw); + long code = 0; + curl_easy_getinfo(e, CURLINFO_RESPONSE_CODE, &code); + + TransferStats stats; + stats.http_status = code; + gather_timings(e, stats); + + auto it = std::find_if(w.live.begin(), w.live.end(), + [raw](const auto &s) { return s.get() == raw; }); + std::shared_ptr st = (it != w.live.end()) ? *it : nullptr; + + curl_multi_remove_handle(w.multi, e); + curl_easy_cleanup(e); + if (st) { + st->easy = nullptr; + if (st->header_slist) { + curl_slist_free_all(st->header_slist); + st->header_slist = nullptr; + } + } + if (it != w.live.end()) + w.live.erase(it); + if (!st) + continue; + + using Stop = Transfer::State::Stop; + const Stop stop = st->stop.load(); + + if (stop == Stop::aborted) { + complete(st, ErrorInfo(Error::canceled)); + } else if (stop == Stop::head_complete) { + // A probe: on_head asked to stop. Headers were the goal -> success, even + // though a ranged GET body-stop surfaces as CURLE_WRITE_ERROR. + stats.bytes_received = st->bytes_received; + stats.effective_url = st->head.effective_url; + complete(st, std::move(stats)); + } else if (res == CURLE_OK && code < 400) { + stats.bytes_received = st->bytes_received; + stats.effective_url = st->head.effective_url; + complete(st, std::move(stats)); + } else { + complete(st, detail::make_error(res, code)); + } + } + } + + static void gather_timings(CURL *e, TransferStats &s) { + auto us_to_ms = [](curl_off_t us) { return us > 0 ? static_cast(us / 1000) : 0L; }; + curl_off_t t = 0; + if (curl_easy_getinfo(e, CURLINFO_NAMELOOKUP_TIME_T, &t) == CURLE_OK) + s.namelookup_ms = us_to_ms(t); + if (curl_easy_getinfo(e, CURLINFO_CONNECT_TIME_T, &t) == CURLE_OK) + s.connect_ms = us_to_ms(t); + if (curl_easy_getinfo(e, CURLINFO_APPCONNECT_TIME_T, &t) == CURLE_OK) + s.appconnect_ms = us_to_ms(t); + if (curl_easy_getinfo(e, CURLINFO_STARTTRANSFER_TIME_T, &t) == CURLE_OK) + s.starttransfer_ms = us_to_ms(t); + if (curl_easy_getinfo(e, CURLINFO_TOTAL_TIME_T, &t) == CURLE_OK) + s.total_ms = us_to_ms(t); + } + + void complete(const std::shared_ptr &st, Result r) { + if (st->finished) + return; + st->finished = true; + if (st->cbs.on_finished) + st->cbs.on_finished(std::move(r)); + } + + void shutdown_worker(Worker &w) { + for (auto &st : w.live) { + if (st->easy) { + curl_multi_remove_handle(w.multi, st->easy); + curl_easy_cleanup(st->easy); + st->easy = nullptr; + } + if (st->header_slist) { + curl_slist_free_all(st->header_slist); + st->header_slist = nullptr; + } + complete(st, ErrorInfo(Error::canceled, "client shutting down")); + } + w.live.clear(); + } +}; + +// --- Transfer ------------------------------------------------------------------- +std::uint64_t Transfer::id() const noexcept { + return state_ ? state_->id : 0; +} + +void Transfer::pause() { + if (state_ && state_->client) + state_->client->enqueue(state_->worker_index, {HttpClient::Impl::CmdKind::pause, state_}); +} +void Transfer::resume() { + if (state_ && state_->client) + state_->client->enqueue(state_->worker_index, {HttpClient::Impl::CmdKind::resume, state_}); +} +void Transfer::cancel() { + if (state_ && state_->client) + state_->client->enqueue(state_->worker_index, {HttpClient::Impl::CmdKind::cancel, state_}); +} + +// --- HttpClient --------------------------------------------------------------- +HttpClient::HttpClient() : HttpClient(Options{}) {} +HttpClient::HttpClient(Options opts) : impl_(std::make_unique(opts)) {} +HttpClient::~HttpClient() = default; + +unsigned HttpClient::worker_count() const noexcept { + return impl_ ? static_cast(impl_->workers.size()) : 0; +} + +Transfer HttpClient::start(Request req, TransferCallbacks cbs) { + auto st = std::make_shared(); + st->id = next_id(); + st->client = impl_.get(); + st->req = std::move(req); + st->cbs = std::move(cbs); + st->worker_index = impl_->workers.empty() ? 0 : impl_->rr.fetch_add(1) % impl_->workers.size(); + + impl_->enqueue(st->worker_index, {Impl::CmdKind::add, st}); + return Transfer(st); +} + +} // namespace vdm::net diff --git a/core/tests/CMakeLists.txt b/core/tests/CMakeLists.txt index 886f9fe..9f83bde 100644 --- a/core/tests/CMakeLists.txt +++ b/core/tests/CMakeLists.txt @@ -20,3 +20,17 @@ vdm_add_test(veloxcore_event_bus_test util/event_bus_test.cpp) vdm_add_test(veloxcore_thread_pool_test util/thread_pool_test.cpp) vdm_add_test(veloxcore_bytes_test util/bytes_test.cpp) vdm_add_test(veloxcore_log_test util/log_test.cpp) + +# net/ integration tests drive tools/testserver (lane PKG/QA). Skip cleanly if it isn't +# in the tree yet (lanes merge independently). +set(_testserver ${CMAKE_SOURCE_DIR}/tools/testserver/testserver.py) +vdm_add_test(veloxcore_http_client_test net/http_client_test.cpp) +target_include_directories(veloxcore_http_client_test PRIVATE ${CMAKE_CURRENT_SOURCE_DIR}/net) +if(EXISTS ${_testserver}) + target_compile_definitions(veloxcore_http_client_test + PRIVATE VDM_TESTSERVER_PY="${_testserver}") + set_tests_properties(veloxcore_http_client_test PROPERTIES TIMEOUT 120) +else() + message(STATUS "veloxcore: tools/testserver not present; http_client_test will skip " + "its server-backed cases.") +endif() diff --git a/core/tests/net/http_client_test.cpp b/core/tests/net/http_client_test.cpp new file mode 100644 index 0000000..bcaba26 --- /dev/null +++ b/core/tests/net/http_client_test.cpp @@ -0,0 +1,270 @@ +#include "vdm/net/http_client.hpp" + +#include +#include +#include +#include +#include +#include + +#include "testserver_fixture.hpp" +#include "vtest.hpp" + +using namespace vdm; +using namespace vdm::net; +using vdm::testing::TestServer; + +namespace { + +// Collects callback output from a transfer and lets the test thread wait for the end. +struct Recorder { + std::mutex mu; + ResponseHead head; + bool head_seen = false; + std::uint64_t bytes = 0; + std::promise> done; + std::future> done_fut = done.get_future(); + + DataAction want_on_head = DataAction::proceed; // set before start() + std::atomic want_on_data{DataAction::proceed}; + std::atomic data_calls{0}; + + TransferCallbacks callbacks() { + return TransferCallbacks{ + .on_head = + [this](const ResponseHead &h) { + std::lock_guard lk(mu); + head = h; + head_seen = true; + return want_on_head; + }, + .on_data = + [this](ConstByteSpan s) { + data_calls.fetch_add(1); + std::lock_guard lk(mu); + bytes += s.size(); + return want_on_data.load(); + }, + .on_finished = [this](Result r) { done.set_value(std::move(r)); }, + }; + } + + Result wait(std::chrono::seconds to = std::chrono::seconds(20)) { + if (done_fut.wait_for(to) != std::future_status::ready) + return Err{Error::timeout, "test wait timed out"}; + return done_fut.get(); + } +}; + +// on_data needs an atomic for the cancel/pause tests to flip it from the test thread. +struct AtomicAction { + std::atomic a{DataAction::proceed}; + void store(DataAction v) { a.store(v); } + operator DataAction() const { return a.load(); } + DataAction load() const { return a.load(); } +}; + +} // namespace + +VT_TEST(http_plain_full_get) { + TestServer srv; + VT_REQUIRE(srv.available()); + HttpClient client({.workers = 1}); + + Recorder rec; + Request req; + req.url = srv.url("/plain/file/64K"); + auto t = client.start(std::move(req), rec.callbacks()); + + auto r = rec.wait(); + VT_REQUIRE(r.has_value()); + VT_CHECK_EQ(r.value().http_status, 200); + VT_CHECK_EQ(r.value().bytes_received, 65536u); + VT_CHECK(rec.head_seen); + VT_CHECK_EQ(rec.head.status, 200); + VT_CHECK(rec.head.headers.has("Accept-Ranges")); + VT_CHECK_EQ(rec.bytes, 65536u); + VT_CHECK(t.id() != 0); +} + +VT_TEST(http_ranged_get_is_206) { + TestServer srv; + VT_REQUIRE(srv.available()); + HttpClient client({.workers = 1}); + + Recorder rec; + Request req; + req.url = srv.url("/plain/file/64K"); + req.range = ByteRange{1000, 1999}; + auto t = client.start(std::move(req), rec.callbacks()); + + auto r = rec.wait(); + VT_REQUIRE(r.has_value()); + VT_CHECK_EQ(rec.head.status, 206); + auto cr = rec.head.headers.get("Content-Range"); + VT_REQUIRE(cr.has_value()); + VT_CHECK_EQ(std::string(*cr), std::string("bytes 1000-1999/65536")); + VT_CHECK_EQ(rec.bytes, 1000u); + VT_CHECK_EQ(r.value().bytes_received, 1000u); +} + +VT_TEST(http_follows_redirect_chain) { + TestServer srv; + VT_REQUIRE(srv.available()); + HttpClient client({.workers = 1}); + + Recorder rec; + Request req; + req.url = srv.url("/redirect-chain/file/16K"); + auto t = client.start(std::move(req), rec.callbacks()); + + auto r = rec.wait(); + VT_REQUIRE(r.has_value()); + VT_CHECK_EQ(rec.head.status, 200); // on_head fires once, for the final response + VT_CHECK_EQ(rec.bytes, 16u * 1024u); + VT_CHECK(r.value().effective_url != srv.url("/redirect-chain/file/16K")); + VT_CHECK(r.value().effective_url.find("_r=done") != std::string::npos); +} + +VT_TEST(http_404_is_not_found_error) { + TestServer srv; + VT_REQUIRE(srv.available()); + HttpClient client({.workers = 1}); + + Recorder rec; + Request req; + req.url = srv.url("/plain/nope"); + auto t = client.start(std::move(req), rec.callbacks()); + + auto r = rec.wait(); + VT_REQUIRE(!r.has_value()); + VT_CHECK_EQ(r.error().code, Error::not_found); + VT_CHECK_EQ(r.error().http_status, 404); +} + +VT_TEST(http_connection_refused_is_connect_failed) { + HttpClient client({.workers = 1}); + Recorder rec; + Request req; + req.url = "http://127.0.0.1:1/nothing"; // nothing listens on :1 + req.connect_timeout_ms = 2000; + auto t = client.start(std::move(req), rec.callbacks()); + + auto r = rec.wait(); + VT_REQUIRE(!r.has_value()); + VT_CHECK_EQ(r.error().code, Error::connect_failed); +} + +VT_TEST(http_head_probe_stops_after_headers) { + TestServer srv; + VT_REQUIRE(srv.available()); + HttpClient client({.workers = 1}); + + Recorder rec; + rec.want_on_head = DataAction::abort; // probe: headers only + Request req; + req.url = srv.url("/plain/file/1M"); + req.method = Method::head; + auto t = client.start(std::move(req), rec.callbacks()); + + auto r = rec.wait(); + VT_REQUIRE(r.has_value()); // head_complete, NOT canceled + VT_CHECK(rec.head_seen); + VT_REQUIRE(rec.head.content_length.has_value()); + VT_CHECK_EQ(*rec.head.content_length, 1024u * 1024u); + VT_CHECK_EQ(rec.data_calls.load(), 0); +} + +VT_TEST(http_ranged_get_probe_stops_without_downloading_file) { + TestServer srv; + VT_REQUIRE(srv.available()); + HttpClient client({.workers = 1}); + + Recorder rec; + rec.want_on_head = DataAction::abort; + Request req; + req.url = srv.url("/plain/file/8M"); + req.range = ByteRange{0, 0}; // classic HEAD-refused fallback + auto t = client.start(std::move(req), rec.callbacks()); + + auto r = rec.wait(); + VT_REQUIRE(r.has_value()); + VT_CHECK(rec.bytes <= 1u); // at most the one probe byte, usually 0 +} + +VT_TEST(http_cancel_mid_transfer) { + TestServer srv; + VT_REQUIRE(srv.available()); + HttpClient client({.workers = 1}); + + AtomicAction data_action; + std::promise> done; + auto fut = done.get_future(); + std::atomic calls{0}; + + TransferCallbacks cbs{ + .on_head = [](const ResponseHead &) { return DataAction::proceed; }, + .on_data = + [&](ConstByteSpan) { + calls.fetch_add(1); + return data_action.load(); + }, + .on_finished = [&](Result r) { done.set_value(std::move(r)); }, + }; + + Request req; + req.url = srv.url("/throttled/file/4M"); // ~128 KiB/s => lots of chunks + auto t = client.start(std::move(req), std::move(cbs)); + + // wait for the transfer to actually start, then cancel + for (int i = 0; i < 200 && calls.load() == 0; ++i) + std::this_thread::sleep_for(std::chrono::milliseconds(10)); + VT_REQUIRE(calls.load() > 0); + t.cancel(); + + VT_REQUIRE(fut.wait_for(std::chrono::seconds(10)) == std::future_status::ready); + auto r = fut.get(); + VT_REQUIRE(!r.has_value()); + VT_CHECK_EQ(r.error().code, Error::canceled); +} + +VT_TEST(http_pause_then_resume_completes) { + TestServer srv; + VT_REQUIRE(srv.available()); + HttpClient client({.workers = 1}); + + std::atomic calls{0}; + std::atomic total{0}; + std::promise> done; + auto fut = done.get_future(); + + TransferCallbacks cbs{ + .on_head = [](const ResponseHead &) { return DataAction::proceed; }, + .on_data = + [&](ConstByteSpan s) { + calls.fetch_add(1); + total.fetch_add(s.size()); + return DataAction::proceed; + }, + .on_finished = [&](Result r) { done.set_value(std::move(r)); }, + }; + + Request req; + req.url = srv.url("/throttled/file/1M"); + auto t = client.start(std::move(req), std::move(cbs)); + + for (int i = 0; i < 200 && calls.load() == 0; ++i) + std::this_thread::sleep_for(std::chrono::milliseconds(10)); + VT_REQUIRE(calls.load() > 0); + + t.pause(); + int calls_at_pause = calls.load(); + std::this_thread::sleep_for(std::chrono::milliseconds(400)); + VT_CHECK(calls.load() - calls_at_pause <= 1); // at most one in-flight chunk slips through + t.resume(); + + VT_REQUIRE(fut.wait_for(std::chrono::seconds(20)) == std::future_status::ready); + auto r = fut.get(); + VT_REQUIRE(r.has_value()); + VT_CHECK_EQ(total.load(), 1024u * 1024u); +} diff --git a/core/tests/net/testserver_fixture.hpp b/core/tests/net/testserver_fixture.hpp new file mode 100644 index 0000000..28a4bc7 --- /dev/null +++ b/core/tests/net/testserver_fixture.hpp @@ -0,0 +1,113 @@ +// testserver_fixture.hpp — spawn tools/testserver for a test, tear it down after. +// +// Linux-only (fork/exec/pipe/kill). The path to testserver.py is injected by CMake as +// VDM_TESTSERVER_PY; if it's empty or missing the fixture reports unavailable() and the +// test should skip. + +#ifndef VDM_TESTS_NET_TESTSERVER_FIXTURE_HPP +#define VDM_TESTS_NET_TESTSERVER_FIXTURE_HPP + +#include +#include +#include +#include + +#include +#include +#include +#include +#include +#include +#include + +#ifndef VDM_TESTSERVER_PY +#define VDM_TESTSERVER_PY "" +#endif + +namespace vdm::testing { + +class TestServer { + public: + TestServer() { + const char *script = VDM_TESTSERVER_PY; + if (!script || !*script || ::access(script, R_OK) != 0) + return; + + int pipefd[2]; + if (::pipe(pipefd) != 0) + return; + + pid_ = ::fork(); + if (pid_ < 0) { + ::close(pipefd[0]); + ::close(pipefd[1]); + return; + } + if (pid_ == 0) { + ::dup2(pipefd[1], STDOUT_FILENO); + ::close(pipefd[0]); + ::close(pipefd[1]); + int devnull = ::open("/dev/null", O_WRONLY); + if (devnull >= 0) + ::dup2(devnull, STDERR_FILENO); + ::execlp("python3", "python3", script, "--port", "0", "--seed", "9", "--loris-seconds", + "1", "--throttle-bps", "131072", static_cast(nullptr)); + ::_exit(127); + } + ::close(pipefd[1]); + + // Read the port line the server prints to stdout. + std::string line; + char c = 0; + auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds(10); + while (std::chrono::steady_clock::now() < deadline) { + ssize_t r = ::read(pipefd[0], &c, 1); + if (r == 1) { + if (c == '\n') + break; + line += c; + } else if (r == 0) { + break; + } else if (errno != EINTR) { + break; + } + } + ::close(pipefd[0]); + if (!line.empty()) + port_ = std::atoi(line.c_str()); + + // Give the listener a moment to accept. + std::this_thread::sleep_for(std::chrono::milliseconds(150)); + } + + ~TestServer() { + if (pid_ > 0) { + ::kill(pid_, SIGTERM); + int status = 0; + for (int i = 0; i < 50; ++i) { + if (::waitpid(pid_, &status, WNOHANG) == pid_) + return; + std::this_thread::sleep_for(std::chrono::milliseconds(20)); + } + ::kill(pid_, SIGKILL); + ::waitpid(pid_, &status, 0); + } + } + + TestServer(const TestServer &) = delete; + TestServer &operator=(const TestServer &) = delete; + + [[nodiscard]] bool available() const { return port_ > 0; } + [[nodiscard]] int port() const { return port_; } + [[nodiscard]] std::string url(const std::string &path) const { + return "http://127.0.0.1:" + std::to_string(port_) + path; + } + + private: + pid_t pid_ = -1; + int port_ = 0; +}; + +} // namespace vdm::testing + +#endif // VDM_TESTS_NET_TESTSERVER_FIXTURE_HPP