diff --git a/docs/roadmap/CVCGL-UI-DSL-ROADMAP.md b/docs/roadmap/CVCGL-UI-DSL-ROADMAP.md index d7a93efc..e495698c 100644 --- a/docs/roadmap/CVCGL-UI-DSL-ROADMAP.md +++ b/docs/roadmap/CVCGL-UI-DSL-ROADMAP.md @@ -3249,6 +3249,24 @@ render/resolver thread**. wasm: `emscripten_fetch` is async-only (no sync on the `fetch_sync()` is unavailable there — the async path is identical across backends. Returns `resource{kind::bytes|local_path}` to slot into §13.4's dispatch-by-kind. +> **LANDED (first slice) — the async `(http-get-async)` state_exec intrinsic.** +> `ariadne/net_intrinsics.cpp` + `register_net_intrinsics(cvc::app&)` (inc/cvc/ariadne/net_intrinsics.h) +> bind a program-lane verb `(http-get-async URL [HEADERS])` that runs the blocking `cvc::net::send` +> on `app.computePool()` (OFF the scheduler thread) and posts the response to a UNIQUE `#`-suffixed +> reply channel it RETURNS; a program awaits with the existing `(msg-recv )`, resuming with a +> dict `{ ok status body(bytes) url headers(list) error }`. This is the nav_compute pattern +> (host-intrinsic → `compute_async` → `exec_scheduler().post_message` → parked `msg-recv` woken by +> `drain_ingress`, the delivered value threaded into the enclosing expression). The verb ALWAYS posts +> a dict (an error dict on failure) so a parked recv never hangs; the body is `bytes` (the §13.9 +> precursor bridge). Registered host-level (needs `cvc::net` + the app), not a core builtin. Tests in +> `ariadne_runtime_test` (AriadneNetIntrinsics: full-dict await, error-path resume, and a +> BlockingHttpClient proof that the scheduler is not blocked while the fetch is in flight), offline via +> the PR1 `set_http_client` fake. **NOT YET:** a *transparent* `(http-get url)` that returns the body +> directly (needs a new evaluator park-token hook — a host `native_fn` gets no `intrinsics_context` at +> call time, so it can't self-park); chunked/streaming bodies; and the dedicated libcurl-`multi` I/O +> thread (today each in-flight fetch blocks one compute-pool worker). NATIVE only until the wasm +> worker/`-sASYNCIFY` fetch path is confirmed. + ### 13.9 A cvc::app-wide HTTP(s) cache in the state tree (TTL + conditional GET) — *planned* Repeated `import:`/`load:`/`source:` of the same `http(s)://` URL must not re-fetch every time. The diff --git a/inc/cvc/ariadne/net_intrinsics.h b/inc/cvc/ariadne/net_intrinsics.h new file mode 100644 index 00000000..77a02a34 --- /dev/null +++ b/inc/cvc/ariadne/net_intrinsics.h @@ -0,0 +1,44 @@ +#ifndef CVC_ARIADNE_NET_INTRINSICS_H +#define CVC_ARIADNE_NET_INTRINSICS_H + +// Ariadne — async HTTP host intrinsics for the state_exec program lanes (roadmap §13.8). Registers +// an `(http-get-async URL [HEADERS])` verb that fetches over cvc::net OFF the scheduler thread and, +// on completion, posts the response to a UNIQUE reply channel the verb RETURNS. A program awaits it +// with the existing `(msg-recv )`, which resumes with a dict: +// +// (get-attr (msg-recv (http-get-async "https://host/x")) "body") +// ; => dict { ok status body(bytes) url headers(list) error } +// +// WHY two-step (verb + msg-recv) and not a transparent `(http-get url)` that returns the body: a +// host verb is a native_fn(std::span) with NO intrinsics_context at call time, so it +// cannot read the calling process's pid or self-park; only a core intrinsic (like msg-recv) can. So +// the verb SUBMITS and returns the channel, and msg-recv (a core intrinsic) does the parking + +// value-threading — the exact nav_compute pattern. A transparent single-verb await needs a new +// evaluator hook (a future item). +// +// THREADING: the fetch runs on app.computePool() (a background worker), so the scheduler/host +// thread is never blocked; the worker touches only the one thread-safe seam, +// exec_scheduler().post_message. The verb ALWAYS posts a dict (an error dict on a transport/arg +// failure), so a parked (msg-recv) never hangs forever. +// +// LANES + LIMITS: use it only in a program `on:`/`on:tick`/`on_key`/`on_pointer` lane — NOT in +// `init:` (no per-frame pump, so a park never resumes) and NOT in the reactive read lane +// (visible_when/computed, default-deny). NATIVE ONLY for now: the background-worker model assumes +// native threads (wasm needs the §13.8 async-fetch path). A no-op without a compiled cvc::net +// backend or without state_exec. + +namespace cvc { +class app; +namespace ariadne { + +// Register the async net intrinsics against `app` (used for the compute pool + the scheduler +// ingress the completion posts to). Appends to the process-global action-intrinsic providers, so +// call it once at setup, before load_*; tear down with cvc::ariadne::clear_action_intrinsics() +// before `app` dies (the provider captures `app` by reference). A no-op if this build lacks +// state_exec or an HTTP backend. +void register_net_intrinsics(cvc::app &app); + +} // namespace ariadne +} // namespace cvc + +#endif // CVC_ARIADNE_NET_INTRINSICS_H diff --git a/src/cvc/CMakeLists.txt b/src/cvc/CMakeLists.txt index c72b0697..3d28bd62 100644 --- a/src/cvc/CMakeLists.txt +++ b/src/cvc/CMakeLists.txt @@ -1118,7 +1118,7 @@ list(APPEND INCLUDE_FILES ../../inc/cvc/net/http_client.h) # the retained Widget tree, the Runtime + reconcile boundary, and direct # cvc::state binding, walking a pluggable Backend. It touches NO VTK/ImGui/GL — # the ImGui-over-VTK backend lives in cvcGL (src/cvcGL/ariadne). Roadmap §16.1. -list(APPEND SOURCE_FILES ariadne/ariadne.cpp ariadne/loader.cpp ariadne/uri.cpp ariadne/uri_state.cpp ariadne/uri_http.cpp ariadne/uri_http_cache.cpp ariadne/state_io.cpp ariadne/ftxui_backend.cpp) +list(APPEND SOURCE_FILES ariadne/ariadne.cpp ariadne/loader.cpp ariadne/uri.cpp ariadne/uri_state.cpp ariadne/uri_http.cpp ariadne/uri_http_cache.cpp ariadne/net_intrinsics.cpp ariadne/state_io.cpp ariadne/ftxui_backend.cpp) list(APPEND INCLUDE_FILES ../../inc/cvc/ariadne/widget.h ../../inc/cvc/ariadne/backend.h @@ -1129,6 +1129,7 @@ list(APPEND INCLUDE_FILES ../../inc/cvc/ariadne/uri_state.h ../../inc/cvc/ariadne/uri_http.h ../../inc/cvc/ariadne/uri_http_cache.h + ../../inc/cvc/ariadne/net_intrinsics.h ../../inc/cvc/ariadne/state_io.h ../../inc/cvc/ariadne/scene.h ../../inc/cvc/ariadne/value.h diff --git a/src/cvc/ariadne/net_intrinsics.cpp b/src/cvc/ariadne/net_intrinsics.cpp new file mode 100644 index 00000000..c4539c64 --- /dev/null +++ b/src/cvc/ariadne/net_intrinsics.cpp @@ -0,0 +1,128 @@ +// Ariadne — async HTTP host intrinsics for the state_exec program lanes (roadmap §13.8). See +// net_intrinsics.h. The whole implementation is gated on CVC_STATE_EXEC (no program lanes without +// it) and on a compiled cvc::net backend (cvc::net::have_http_backend()). + +#include // register_action_intrinsics, have_state_exec +#include +#include +#include + +#ifdef CVC_STATE_EXEC + +#include +#include +#include // exec_scheduler().post_message +#include // register_fn — bind the host verb into the lanes +#include // resolve_channel_key — scope the reply channel +#include // value_t, make_dict/make_list/make_bytes +#include +#include +#include +#include + +namespace cvc { +namespace ariadne { + +namespace se = cvc::state_exec; + +namespace { + +// Marshal a cvc::net::HttpResponse into the DSL reply dict. EVERY key is populated on every path +// (get-attr throws on a missing key), so an error response is a well-formed dict too. The body is +// `bytes` (opaque octets, the sanctioned home for an HTTP body — the state_exec bytes track). +se::value_t marshal_response(const cvc::net::HttpResponse &r) { + std::vector headers; + headers.reserve(r.headers.size()); + for (const std::string &h : r.headers) + headers.push_back(se::value_t(h)); + return se::make_dict({ + {"ok", se::value_t(r.ok)}, + {"status", se::value_t(static_cast(r.status))}, + {"body", se::make_bytes(r.body)}, + {"url", se::value_t(r.canonical_url)}, + {"headers", se::make_list(std::move(headers))}, + {"error", se::value_t(r.error)}, + }); +} + +se::value_t marshal_error(const std::string &url, const std::string &message) { + cvc::net::HttpResponse r; + r.ok = false; + r.status = 0; + r.canonical_url = url; + r.error = message; + return marshal_response(r); +} + +} // namespace + +void register_net_intrinsics(cvc::app &app) { + if (!have_state_exec() || !cvc::net::have_http_backend()) + return; + // Warm the lazy per-app singletons on THIS thread before any background worker touches them, so a + // worker never races their first construction (the nav_compute discipline). + app.computePool(); + app.exec_scheduler(); + + register_action_intrinsics( + [&app](std::shared_ptr env, se::intrinsics_context &ictx) { + // Capture the lane's chroot so the reply channel is scoped the same way the program's + // (msg-recv ) will resolve it. The verb generates a UNIQUE '#'-suffixed channel per + // call (returned verbatim by resolve_channel_key + exempt from channel-policy), so two + // in-flight fetches never collide on the single recv_path a process has. + const std::string root = ictx.root_path; + se::builtins::register_fn( + env, "http-get-async", [&app, root](std::span args) -> se::value_t { + cvc::net::HttpRequest req; + if (!args.empty()) + if (const std::string *url = std::get_if(&args[0].v)) + req.url = *url; + // Optional second arg: a list of verbatim "Name: value" header strings. + if (args.size() > 1) + if (const se::list_ptr *lp = std::get_if(&args[1].v)) + for (const se::value_t &e : **lp) + if (const std::string *h = std::get_if(&e.v)) + req.headers.push_back(*h); + + static std::atomic seq{0}; + const std::string chan = + "http.reply#" + std::to_string(seq.fetch_add(1, std::memory_order_relaxed)); + const std::string done = se::resolve_channel_key(root, chan); // '#' => == chan + + // Run the BLOCKING fetch on a compute-pool worker; post the marshalled dict when it + // joins. compute_async returns immediately, so the verb never blocks the scheduler + // thread. The worker ALWAYS posts (an error dict on a throw), so the parked + // (msg-recv) always resumes. + const std::string url = req.url; + app.compute_async( + 1, [](int) {}, + [&app, req, done, url] { + se::value_t payload; + try { + payload = marshal_response(cvc::net::send(req)); + } catch (const std::exception &e) { + payload = marshal_error(url, std::string("http-get: ") + e.what()); + } catch (...) { + payload = marshal_error(url, "http-get: unknown error"); + } + app.exec_scheduler().post_message(done, payload); + }); + return se::value_t(chan); // the program awaits with (msg-recv ) + }); + }); +} + +} // namespace ariadne +} // namespace cvc + +#else // !CVC_STATE_EXEC + +namespace cvc { +namespace ariadne { +void register_net_intrinsics(cvc::app & /*app*/) { + // No program lanes without state_exec — nothing to register. +} +} // namespace ariadne +} // namespace cvc + +#endif // CVC_STATE_EXEC diff --git a/src/cvc/tests/ariadne_runtime_test.cpp b/src/cvc/tests/ariadne_runtime_test.cpp index 5a4e943b..12848bee 100644 --- a/src/cvc/tests/ariadne_runtime_test.cpp +++ b/src/cvc/tests/ariadne_runtime_test.cpp @@ -5,9 +5,13 @@ #include #include // item 2: the nav-step worker's arrived-count accumulator +#include +#include +#include #include #include -#include // §12 end-to-end load: mount through the real Runtime +#include // §12 end-to-end load: mount through the real Runtime +#include // §13.8 async (http-get-async) intrinsic test #include #include #include @@ -15,9 +19,12 @@ #include // host-intrinsic seam test: register_fn #include // value_t #include // item 2: full type for app.computePool().parallel_for +#include // §13.8: swap the transport for a fake in the http-get test #include #include #include +#include +#include #include #include #include @@ -2133,3 +2140,191 @@ TEST(AriadneChannelPolicy, ExemptsHashChannel) { rt.set_channel_policy({"nav.done"}, {}, true, false); EXPECT_TRUE(msg_delivered(app, rt, mb, "sys#evt", "sys#evt")); // '#'-runtime channel is exempt } + +// --------------------------------------------------------------------------- +// §13.8 async (http-get-async) intrinsic — offline via the cvc::net fake transport. +// --------------------------------------------------------------------------- +namespace { + +// A canned transport (PR1's set_http_client seam): returns one fixed response, records the call. +class CannedHttpClient : public cvc::net::HttpClient { +public: + explicit CannedHttpClient(cvc::net::HttpResponse r) : resp_(std::move(r)) {} + std::atomic calls{0}; + std::string last_url; + cvc::net::HttpResponse send(const cvc::net::HttpRequest &req) override { + last_url = req.url; + ++calls; + return resp_; + } + +private: + cvc::net::HttpResponse resp_; +}; + +// A transport that blocks in send() until released — proves the verb does NOT block the scheduler. +class BlockingHttpClient : public cvc::net::HttpClient { +public: + explicit BlockingHttpClient(cvc::net::HttpResponse r) : resp_(std::move(r)) {} + cvc::net::HttpResponse send(const cvc::net::HttpRequest &) override { + std::unique_lock lk(m_); + cv_.wait(lk, [&] { return release_; }); + return resp_; + } + void release() { + { + std::lock_guard lk(m_); + release_ = true; + } + cv_.notify_all(); + } + +private: + std::mutex m_; + std::condition_variable cv_; + bool release_ = false; + cvc::net::HttpResponse resp_; +}; + +// Restore process-global state so one test never leaks into another. +struct NetIntrinsicsGuard { + ~NetIntrinsicsGuard() { + clear_action_intrinsics(); + cvc::net::set_http_client(nullptr); + } +}; + +cvc::net::HttpResponse http_ok(long status, std::string body, std::vector headers = {}, + std::string url = "http://ex/x") { + cvc::net::HttpResponse r; + r.ok = true; + r.status = status; + r.body = std::move(body); + r.headers = std::move(headers); + r.canonical_url = std::move(url); + return r; +} + +// Pump render+drain until `pred` holds or the budget elapses (compute_async posts from a background +// worker, so a bounded pump stands in for the nav test's hand-rolled join()). +template bool pump_until(Runtime &rt, Pred pred, int budget_ms = 5000) { + using namespace std::chrono; + const auto deadline = steady_clock::now() + milliseconds(budget_ms); + while (steady_clock::now() < deadline) { + rt.render(); + rt.drain(); + if (pred()) + return true; + std::this_thread::sleep_for(milliseconds(2)); + } + return pred(); +} + +std::string node_data_string(cvc::app &app, const char *path) { + const boost::any d = cvc::state::instance(app)(path).data(); + const std::string *s = boost::any_cast(&d); + return s ? *s : std::string(); +} + +} // namespace + +TEST(AriadneNetIntrinsics, HttpGetAsyncAwaitsResponseDict) { + if (!have_state_exec()) + GTEST_SKIP() << "libcvc built without state_exec"; + if (!cvc::net::have_http_backend()) + GTEST_SKIP() << "libcvc built without an HTTP backend"; + cvc::app app; + NetIntrinsicsGuard guard; + auto *fake = new CannedHttpClient(http_ok(200, "hi", {"Content-Type: text/plain"})); + cvc::net::set_http_client(std::unique_ptr(fake)); + register_net_intrinsics(app); + + Runtime rt(app, ""); + MockBackend mb; + rt.set_backend(&mb); + rt.set_root(group({button("Go", "(begin" + " (set r (msg-recv (http-get-async \"http://ex/x\")))" + " (state-set \"r.status\" (get-attr r \"status\"))" + " (state-set \"r.ok\" (get-attr r \"ok\"))" + " (state-set \"r.error\" (get-attr r \"error\"))" + " (state-data-set \"r.body\" (get-attr r \"body\")))")})); + mb.button_click = true; + rt.render(); + rt.drain(); // http-get-async launches the worker; the action parks on msg-recv (drain never + // blocks) + mb.button_click = false; + + ASSERT_TRUE(pump_until(rt, [&] { + return !cvc::state::instance(app)("r.status").value().empty(); + })) << "the parked action never resumed"; + EXPECT_EQ(cvc::state::instance(app)("r.status").value(), "200"); + EXPECT_EQ(cvc::state::instance(app)("r.ok").value(), "true"); + EXPECT_EQ(cvc::state::instance(app)("r.error").value(), ""); // empty on success + EXPECT_EQ(node_data_string(app, "r.body"), "hi"); // the bytes body round-trips byte-exact + EXPECT_EQ(fake->calls.load(), 1); + EXPECT_EQ(fake->last_url, "http://ex/x"); +} + +TEST(AriadneNetIntrinsics, HttpGetAsyncErrorPathResumesWithErrorDict) { + if (!have_state_exec()) + GTEST_SKIP() << "libcvc built without state_exec"; + if (!cvc::net::have_http_backend()) + GTEST_SKIP() << "libcvc built without an HTTP backend"; + cvc::app app; + NetIntrinsicsGuard guard; + cvc::net::HttpResponse err; + err.ok = false; + err.status = 0; + err.error = "boom"; + cvc::net::set_http_client(std::make_unique(err)); + register_net_intrinsics(app); + + Runtime rt(app, ""); + MockBackend mb; + rt.set_backend(&mb); + rt.set_root(group({button("Go", "(begin" + " (set r (msg-recv (http-get-async \"http://ex/down\")))" + " (state-set \"r.ok\" (get-attr r \"ok\"))" + " (state-set \"r.error\" (get-attr r \"error\")))")})); + mb.button_click = true; + rt.render(); + rt.drain(); + mb.button_click = false; + + ASSERT_TRUE(pump_until(rt, [&] { return !cvc::state::instance(app)("r.error").value().empty(); })) + << "the action must resume even on a transport error"; + EXPECT_EQ(cvc::state::instance(app)("r.ok").value(), "false"); + EXPECT_EQ(cvc::state::instance(app)("r.error").value(), "boom"); +} + +TEST(AriadneNetIntrinsics, HttpGetAsyncDoesNotBlockTheScheduler) { + if (!have_state_exec()) + GTEST_SKIP() << "libcvc built without state_exec"; + if (!cvc::net::have_http_backend()) + GTEST_SKIP() << "libcvc built without an HTTP backend"; + cvc::app app; + NetIntrinsicsGuard guard; + auto *fake = new BlockingHttpClient(http_ok(200, "later")); + cvc::net::set_http_client(std::unique_ptr(fake)); + register_net_intrinsics(app); + + Runtime rt(app, ""); + MockBackend mb; + rt.set_backend(&mb); + rt.set_root(group( + {button("Go", "(state-set \"r.status\" " + "(get-attr (msg-recv (http-get-async \"http://ex/slow\")) \"status\"))")})); + mb.button_click = true; + rt.render(); + rt.drain(); // the verb returns immediately + the action parks; the fetch is still blocked in + // send() + mb.button_click = false; + // drain() returned though the transport is blocked → the scheduler thread was not blocked. + EXPECT_TRUE(cvc::state::instance(app)("r.status").value().empty()) + << "the result must not be delivered while the fetch is in flight"; + + fake->release(); // let the worker's send() complete + ASSERT_TRUE( + pump_until(rt, [&] { return !cvc::state::instance(app)("r.status").value().empty(); })); + EXPECT_EQ(cvc::state::instance(app)("r.status").value(), "200"); +}