Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 21 additions & 5 deletions docs/roadmap/CVCGL-UI-DSL-ROADMAP.md
Original file line number Diff line number Diff line change
Expand Up @@ -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*

Expand Down
23 changes: 23 additions & 0 deletions inc/cvc/core/state_exec/intrinsics.h
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
#include <cvc/core/state_exec/types.h>
#include <functional>
#include <memory>
#include <optional>
#include <string>
#include <unordered_map>
#include <vector>
Expand Down Expand Up @@ -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__": "<reply-channel>"} (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 <future>) and (msg-recv <future>) 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<std::string> future_channel_of(const value_t &v);

} // namespace cvc::state_exec

#endif // CVC_STATE_EXEC_INTRINSICS_H
110 changes: 64 additions & 46 deletions src/cvc/ariadne/net_intrinsics.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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<const se::value_t> args) {
cvc::net::HttpRequest req;
if (!args.empty())
if (const std::string *url = std::get_if<std::string>(&args[0].v))
req.url = *url;
if (args.size() > 1)
if (const se::list_ptr *lp = std::get_if<se::list_ptr>(&args[1].v))
for (const se::value_t &e : **lp)
if (const std::string *h = std::get_if<std::string>(&e.v))
req.headers.push_back(*h);

static std::atomic<std::uint64_t> 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) {
Expand All @@ -64,52 +102,32 @@ void register_net_intrinsics(cvc::app &app) {
app.computePool();
app.exec_scheduler();

register_action_intrinsics(
[&app](std::shared_ptr<se::environment> env, se::intrinsics_context &ictx) {
// Capture the lane's chroot so the reply channel is scoped the same way the program's
// (msg-recv <chan>) 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<const se::value_t> args) -> se::value_t {
cvc::net::HttpRequest req;
if (!args.empty())
if (const std::string *url = std::get_if<std::string>(&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<se::list_ptr>(&args[1].v))
for (const se::value_t &e : **lp)
if (const std::string *h = std::get_if<std::string>(&e.v))
req.headers.push_back(*h);

static std::atomic<std::uint64_t> 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 <chan>)
});
});
register_action_intrinsics([&app](std::shared_ptr<se::environment> 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<const se::value_t> 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<const se::value_t> 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
Expand Down
96 changes: 69 additions & 27 deletions src/cvc/core/state_exec/intrinsics.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -405,8 +405,19 @@ value_t intrinsic_await(intrinsics_context *ctx, std::span<const value_t> 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<std::string> 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];
}
Expand Down Expand Up @@ -759,41 +770,28 @@ value_t intrinsic_msg_send(intrinsics_context *ctx, std::span<const value_t> arg
value_t intrinsic_msg_recv(intrinsics_context *ctx, std::span<const value_t> 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<std::string> 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<const value_t> args) {
Expand Down Expand Up @@ -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<std::string> future_channel_of(const value_t &v) {
const dict_ptr *dp = std::get_if<dict_ptr>(&v.v);
if (!dp || !*dp)
return std::nullopt;
const std::vector<std::pair<std::string, value_t>> &entries = **dp;
if (entries.size() != 1 || entries[0].first != "__future__")
return std::nullopt;
const std::string *chan = std::get_if<std::string>(&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<value_t> 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
Loading
Loading