From bc3ee0f03a0548910fd0c69a7836192447935dcc Mon Sep 17 00:00:00 2001 From: Joe Rivera Date: Tue, 29 Sep 2026 03:41:22 -0500 Subject: [PATCH] =?UTF-8?q?state=5Fexec/ariadne:=20async/await=20core=20?= =?UTF-8?q?=E2=80=94=20transparent=20(http-get)=20+=20proper=20await?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Deliver a transparent (http-get URL) verb and a value-carrying (await …) with ZERO evaluator edits: the suspend/resume-with-value mechanism already existed (msg-recv proves it), so this factors it into one audited primitive and layers the DSL surface on top. park_on_channel(sched, proc, pid, ch) (intrinsics.{h,cpp}) is the shared "suspend the running process, return a nil placeholder, resume with the value the scheduler delivers into the enclosing expression" primitive. It is msg-recv's body factored out (inbox -> pending -> receive_message -> nil), plus a top-level-call guard: a park needs an enclosing frame to receive the patched value, so a stepped top-level call (state.stack.size()==1) throws instead of leaving a done-and-waiting zombie (size 0 is a direct unit-test call, allowed). msg-recv becomes a thin caller; behavior is unchanged for every existing (nested) use. Futures: a future handle is a tagged dict {"__future__": ""} (no new value_t alternative, so no codec/serialization ripple). make_future/future_channel_of are the shared convention. (await x): a future handle parks value-carrying and resumes with the resolved value; an already-settled value keeps its one-frame cooperative yield (unchanged). (http-get-async) now returns a future, so (await (http-get-async url)) works; msg-recv accepts a future OR a channel string, so the two-step (msg-recv (http-get-async url)) still works verbatim. (http-get URL [HEADERS]) is the transparent one-shot (net_intrinsics.cpp): it kicks the fetch on the compute pool and SELF-PARKS via park_on_channel, reaching the scheduler through its captured &app (current_pid()/current_process() are valid around the evaluator step) -- so (get-attr (http-get url) "status") reads as one expression, no msg-recv. The blocking send stays off the scheduler thread and the worker always posts a dict (an error dict on a throw), so a waiter never hangs. http-get-async is kept as the low-level fan-out primitive. Tests (ariadne_runtime_test AriadneNetIntrinsics): transparent (http-get) yields the dict + the bytes body round-trips; (await (http-get-async …)) resolves the future. The msg-recv/await refactor is regression-clean: ariadne_runtime 104, state_exec intrinsics 154, async 71, integration 85, scheduler 108, multiprocess 5; 0/20 flaky on the async cases. Offline via the set_http_client fake. First slice of the async arc (§13.8). Next: async URI resolver for ari-generated requests, then render-pass park budget/safety. --- docs/roadmap/CVCGL-UI-DSL-ROADMAP.md | 26 ++++-- inc/cvc/core/state_exec/intrinsics.h | 23 ++++++ src/cvc/ariadne/net_intrinsics.cpp | 110 ++++++++++++++----------- src/cvc/core/state_exec/intrinsics.cpp | 96 +++++++++++++++------ src/cvc/tests/ariadne_runtime_test.cpp | 62 ++++++++++++++ 5 files changed, 239 insertions(+), 78 deletions(-) diff --git a/docs/roadmap/CVCGL-UI-DSL-ROADMAP.md b/docs/roadmap/CVCGL-UI-DSL-ROADMAP.md index e495698c..de747e8b 100644 --- a/docs/roadmap/CVCGL-UI-DSL-ROADMAP.md +++ b/docs/roadmap/CVCGL-UI-DSL-ROADMAP.md @@ -3261,11 +3261,27 @@ render/resolver thread**. wasm: `emscripten_fetch` is async-only (no sync on the > 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. +> the PR1 `set_http_client` fake. +> +> **LANDED (async/await core) — transparent `(http-get url)` + proper `await`, ZERO evaluator edits.** +> The suspend/resume-with-value mechanism already existed (`msg-recv` proves it), so the linchpin was +> *factoring* it, not extending the evaluator: `park_on_channel(sched, proc, pid, ch)` +> (`intrinsics.{h,cpp}`) is the one audited "suspend the running process, return a nil placeholder, +> resume with the delivered value in the enclosing expression" primitive `msg-recv`, `await`, and +> `http-get` all share (a host verb self-parks by reaching the scheduler through its captured `&app` — +> `current_pid()`/`receive_message()` — no `intrinsics_context` needed). A top-level-call guard rejects +> a park with no enclosing frame (else a done-and-waiting zombie). `(await x)` is now value-carrying: +> a **future handle** (a tagged dict `{"__future__": chan}`, no new `value_t` alternative) parks and +> resumes with the resolved value; a settled value keeps its one-frame yield. `(http-get-async)` now +> returns a future (so `(await (http-get-async url))` works), and `msg-recv` accepts a future *or* a +> channel string (so the two-step `(msg-recv (http-get-async url))` still works verbatim). +> **`(http-get url)`** is the transparent one-shot: it self-parks and yields the dict straight into the +> enclosing expression — `(get-attr (http-get url) "status")`, no `msg-recv`. Tests: AriadneNetIntrinsics +> gains the transparent-verb + await-future cases; the msg-recv/await refactor is regression-clean +> across the state_exec suites. **NOT YET:** chunked/streaming bodies; the dedicated libcurl-`multi` +> I/O thread (today each in-flight fetch blocks one compute-pool worker); making ALL ari-generated URI +> requests async (the resolver — separate PR); render-pass park budget (separate PR). 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* diff --git a/inc/cvc/core/state_exec/intrinsics.h b/inc/cvc/core/state_exec/intrinsics.h index 9d42cf5f..4fc32362 100644 --- a/inc/cvc/core/state_exec/intrinsics.h +++ b/inc/cvc/core/state_exec/intrinsics.h @@ -13,6 +13,7 @@ #include #include #include +#include #include #include #include @@ -106,6 +107,28 @@ void apply_chroot(intrinsics_context &ctx, cvc::state &tree_root, const std::str /// tree off-thread. std::string resolve_channel_key(const std::string &root_path, const std::string &channel); +// §13.8 async park primitive — suspend the CURRENTLY-RUNNING process until a value is delivered on +// scheduler channel `ch`, returning a nil placeholder that the scheduler's delivery +// (deliver_to_receivers) patches with the real value in the ENCLOSING expression. This is the +// msg-recv protocol, factored so `await` and host async verbs (e.g. an (http-get) that self-parks) +// share ONE audited suspend path — no evaluator change needed. `ch` must be an ALREADY-RESOLVED +// scheduler key (the caller applies channel policy/scoping; a '#'-runtime channel is exempt). The +// running process is named by `proc` + `pid` (from sched->current_process()/current_pid(), with the +// ctx->proc / ctx->pid fallback for unit-test contexts). Fast-paths a buffered inbox / pending so a +// value that beat the park is returned without suspending. THROWS if the process cannot be +// suspended, or if this is a top-level call with no enclosing frame to receive the delivered value +// (which would otherwise leave a done-and-waiting zombie). +value_t park_on_channel(scheduler_base *sched, process *proc, int pid, const std::string &ch); + +// §13.8 futures — a future handle is a tagged one-key dict {"__future__": ""} (a +// dict_ptr; NO new value_t alternative, so no codec/serialization ripple). An async producer such +// as (http-get-async) returns make_future(chan); (await ) and (msg-recv ) unwrap it +// with future_channel_of and park on the channel, while a plain value / non-future dict yields +// nullopt (treated as an already-settled value). make_future + future_channel_of are the single +// shared convention so producers and consumers never disagree on the tag. +value_t make_future(const std::string &channel); +std::optional future_channel_of(const value_t &v); + } // namespace cvc::state_exec #endif // CVC_STATE_EXEC_INTRINSICS_H diff --git a/src/cvc/ariadne/net_intrinsics.cpp b/src/cvc/ariadne/net_intrinsics.cpp index c4539c64..285b0521 100644 --- a/src/cvc/ariadne/net_intrinsics.cpp +++ b/src/cvc/ariadne/net_intrinsics.cpp @@ -54,6 +54,44 @@ se::value_t marshal_error(const std::string &url, const std::string &message) { return marshal_response(r); } +// Shared by both verbs: build the request from (url [, headers]), kick the BLOCKING cvc::net::send +// onto a compute-pool worker, and have the worker post the marshalled dict to a UNIQUE '#'-suffixed +// reply channel when it joins. Returns that channel (already resolved; '#' => policy-exempt + +// identity). The verb never blocks the scheduler thread — compute_async returns immediately — and +// the worker ALWAYS posts (an error dict on a throw), so a waiter never hangs. +std::string launch_http_fetch(cvc::app &app, const std::string &root, + std::span args) { + cvc::net::HttpRequest req; + if (!args.empty()) + if (const std::string *url = std::get_if(&args[0].v)) + req.url = *url; + 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 + 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 done; +} + } // namespace void register_net_intrinsics(cvc::app &app) { @@ -64,52 +102,32 @@ void register_net_intrinsics(cvc::app &app) { 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 ) - }); - }); + 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 an (msg-recv)/(await) + // resolves it (launch_http_fetch uses a UNIQUE '#'-suffixed channel per call — policy-exempt + + // identity — so two in-flight fetches never collide on the single recv_path a process has). + const std::string root = ictx.root_path; + + // (http-get-async URL [HEADERS]) — kicks the fetch and returns a FUTURE handle; the program + // awaits it with (await …) or (msg-recv …). The low-level primitive for fanning out N fetches. + se::builtins::register_fn(env, "http-get-async", + [&app, root](std::span args) -> se::value_t { + return se::make_future(launch_http_fetch(app, root, args)); + }); + + // (http-get URL [HEADERS]) — TRANSPARENT: kicks the fetch and SELF-PARKS the calling process, + // resuming with the response dict threaded straight into the enclosing expression (no channel, + // no msg-recv). Reaches the scheduler via the captured app; current_pid()/current_process() are + // valid here (set around the evaluator step). Equivalent to (await (http-get-async URL)). + se::builtins::register_fn(env, "http-get", + [&app, root](std::span args) -> se::value_t { + const std::string done = launch_http_fetch(app, root, args); + auto &sched = app.exec_scheduler(); + return se::park_on_channel(&sched, sched.current_process().get(), + sched.current_pid(), done); + }); + }); } } // namespace ariadne diff --git a/src/cvc/core/state_exec/intrinsics.cpp b/src/cvc/core/state_exec/intrinsics.cpp index 4f860926..a278f27d 100644 --- a/src/cvc/core/state_exec/intrinsics.cpp +++ b/src/cvc/core/state_exec/intrinsics.cpp @@ -405,8 +405,19 @@ value_t intrinsic_await(intrinsics_context *ctx, std::span args) int pid = ctx->sched->current_pid(); if (pid < 0) pid = ctx->pid; - // Yield the running process until the next frame boundary (a no-op on a scheduler with no - // frames, where await is identity), then resume with the already-evaluated argument value. + // A FUTURE handle {"__future__": chan} → value-carrying park (the §13.8 proper await): suspend + // until the producer delivers, then resume with the resolved value threaded into the enclosing + // expression, exactly as (msg-recv) does — via the shared park_on_channel primitive. + if (const std::optional fchan = future_channel_of(args[0])) { + process *proc = nullptr; + if (ctx->sched->current_process()) + proc = ctx->sched->current_process().get(); + else if (ctx->proc) + proc = ctx->proc.get(); + return park_on_channel(ctx->sched, proc, pid, *fchan); + } + // An already-settled value → the cooperative one-frame yield await has always been (a no-op on a + // scheduler with no frames, where await is identity), then resume with the value unchanged. ctx->sched->yield_frame(pid); return args[0]; } @@ -759,41 +770,28 @@ value_t intrinsic_msg_send(intrinsics_context *ctx, std::span arg value_t intrinsic_msg_recv(intrinsics_context *ctx, std::span args) { expect_exact(args, 1, "msg-recv"); require_sched(ctx, "msg-recv"); - const std::string &raw = as_string(args[0], "msg-recv"); - enforce_channel_policy(ctx, raw); // §12 runtime enforcement (strict → throw) - const std::string ch = resolve_channel(ctx, raw); // §12 chroot-scoped scheduler key + // Accept a future handle {"__future__": chan} (from an async producer like (http-get-async)) OR a + // raw channel string. A future's channel is an already-resolved '#'-channel (policy-exempt); a + // raw string goes through §12 policy + chroot-scoping as before. + std::string ch; + if (const std::optional fchan = future_channel_of(args[0])) { + ch = *fchan; + } else { + const std::string &raw = as_string(args[0], "msg-recv"); + enforce_channel_policy(ctx, raw); + ch = resolve_channel(ctx, raw); + } - // Resolve the actual PID of the currently executing process. int pid = ctx->sched->current_pid(); if (pid < 0) pid = ctx->pid; // fallback for unit-test contexts - - // Resolve the process pointer for inbox check. process *proc = nullptr; if (ctx->sched->current_process()) proc = ctx->sched->current_process().get(); else if (ctx->proc) proc = ctx->proc.get(); - // If the process already has a message in its inbox, return it immediately. - if (proc && !proc->inbox.empty()) { - auto msg = std::move(proc->inbox.front()); - proc->inbox.pop(); - return msg; - } - - // Check the scheduler's pending message queue for this channel. - auto pending = ctx->sched->pop_pending_message(ch); - if (pending) - return *pending; - - // Otherwise, suspend the process until a message arrives. - bool ok = ctx->sched->receive_message(pid, ch); - if (!ok) - throw std::runtime_error("msg-recv: cannot suspend process " + std::to_string(pid)); - // Return nil as a placeholder — the scheduler will overwrite - // state.result with the actual message when delivery occurs. - return nil_value; + return park_on_channel(ctx->sched, proc, pid, ch); } value_t intrinsic_msg_pending(intrinsics_context *ctx, std::span args) { @@ -893,4 +891,48 @@ std::string resolve_channel_key(const std::string &root_path, const std::string return root_path + cvc::state::SEPARATOR + "channels" + cvc::state::SEPARATOR + channel; } +value_t make_future(const std::string &channel) { + return make_dict({{"__future__", value_t(channel)}}); +} + +std::optional future_channel_of(const value_t &v) { + const dict_ptr *dp = std::get_if(&v.v); + if (!dp || !*dp) + return std::nullopt; + const std::vector> &entries = **dp; + if (entries.size() != 1 || entries[0].first != "__future__") + return std::nullopt; + const std::string *chan = std::get_if(&entries[0].second.v); + if (!chan) + return std::nullopt; + return *chan; +} + +value_t park_on_channel(scheduler_base *sched, process *proc, int pid, const std::string &ch) { + // A message that arrived before the park is returned without suspending: first the process inbox, + // then the scheduler's per-channel pending queue. + if (proc && !proc->inbox.empty()) { + value_t msg = std::move(proc->inbox.front()); + proc->inbox.pop(); + return msg; + } + if (const std::optional pending = sched->pop_pending_message(ch)) + return *pending; + + // A park needs an ENCLOSING frame to receive the value deliver_to_receivers patches in. During + // the native call the apply frame is on top of the stack; size 1 means it is the ONLY frame — a + // top-level call whose pop would empty the stack (done=true) and skip the patch, leaving a + // done-and-waiting zombie. (size 0 == a direct unit-test call with no evaluator frames — allowed; + // the test drives delivery itself.) + if (proc && proc->state.stack.size() == 1) + throw std::runtime_error("cannot suspend a top-level call (no enclosing expression to receive " + "the value); wrap it, e.g. bind the result or use (begin …)"); + + if (!sched->receive_message(pid, ch)) + throw std::runtime_error("cannot suspend process " + std::to_string(pid)); + // Nil placeholder — deliver_to_receivers overwrites the enclosing frame's results.back() with the + // delivered value when the message arrives on `ch`. + return nil_value; +} + } // namespace cvc::state_exec diff --git a/src/cvc/tests/ariadne_runtime_test.cpp b/src/cvc/tests/ariadne_runtime_test.cpp index 12848bee..579bb1e7 100644 --- a/src/cvc/tests/ariadne_runtime_test.cpp +++ b/src/cvc/tests/ariadne_runtime_test.cpp @@ -2328,3 +2328,65 @@ TEST(AriadneNetIntrinsics, HttpGetAsyncDoesNotBlockTheScheduler) { pump_until(rt, [&] { return !cvc::state::instance(app)("r.status").value().empty(); })); EXPECT_EQ(cvc::state::instance(app)("r.status").value(), "200"); } + +// PR-A: the TRANSPARENT verb — (http-get url) self-parks and yields the dict directly, no msg-recv. +TEST(AriadneNetIntrinsics, HttpGetTransparentReturnsDictNoMsgRecv) { + 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 (http-get \"http://ex/x\"))" + " (state-set \"r.status\" (get-attr r \"status\"))" + " (state-data-set \"r.body\" (get-attr r \"body\")))")})); + mb.button_click = true; + rt.render(); + rt.drain(); // (http-get …) self-parks the action; drain never blocks + mb.button_click = false; + + ASSERT_TRUE(pump_until(rt, [&] { + return !cvc::state::instance(app)("r.status").value().empty(); + })) << "the self-parked (http-get) never resumed"; + EXPECT_EQ(cvc::state::instance(app)("r.status").value(), "200"); + EXPECT_EQ(node_data_string(app, "r.body"), "hi"); // the dict flowed into `set r` transparently + EXPECT_EQ(fake->calls.load(), 1); // one transparent fetch, one transfer +} + +// PR-A: PROPER await — (await (http-get-async url)) resolves the future returned by the async verb. +TEST(AriadneNetIntrinsics, AwaitResolvesHttpGetAsyncFuture) { + 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")); + 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 (await (http-get-async \"http://ex/x\")) \"status\"))")})); + mb.button_click = true; + rt.render(); + rt.drain(); + mb.button_click = false; + + ASSERT_TRUE(pump_until(rt, [&] { + return !cvc::state::instance(app)("r.status").value().empty(); + })) << "(await ) never resolved"; + EXPECT_EQ(cvc::state::instance(app)("r.status").value(), "200"); + EXPECT_EQ(fake->calls.load(), 1); +}