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
8 changes: 4 additions & 4 deletions inc/cvc/core/async_lane.h
Original file line number Diff line number Diff line change
Expand Up @@ -42,15 +42,15 @@
* resolves the same channel to in (msg-recv ...). Every path posts exactly one value, so a waiter
* never hangs: a throwing on_error posts a fallback error string, and a stopped lane posts
* on_error("... is stopped") at once instead of dropping the task.
* Use '#' channels (e.g. "demo3#built.tracks") for anything a lane posts: a '#' channel is
* Use '#' channels (e.g. "app#built.tracks") for anything a lane posts: a '#' channel is
* the identity under chroot scoping and exempt from the section 12 channel policy, so the
* poster and a scoped receiver always agree on the key. A resident process can then wait on
* lane results forever:
*
* register_fn(env, "sim-launch", [&](auto) { // runs on the pump thread
* return se::value_t(cvc::launch_lane_task(sim_lane, sched, ictx.root_path, "demo3",
* return se::value_t(cvc::launch_lane_task(sim_lane, sched, ictx.root_path, "app",
* [&] { return run_ticks(); }, on_err));
* }); // -> "demo3#sim.<n>"
* }); // -> "app#sim.<n>"
* sched.execute("(while t (apply-sim (msg-recv (sim-launch))))", opts);
*
* The `work` rules of async_task.h apply unchanged: it runs OFF the scheduler thread, so marshal
Expand Down Expand Up @@ -214,7 +214,7 @@ using lane_task_on_error = pool_task_on_error;
/// Run `work` on `lane` and post its value -- or on_error(message) if it throws -- on `sched` to
/// the FIXED channel `channel`, resolved for chroot `root` with resolve_channel_key (the key a
/// process chrooted to `root` gets for (msg-recv channel); a '#' channel is used verbatim). For a
/// resident receiver such as (while t (apply-x (msg-recv "demo3#built.x"))). Callable from any
/// resident receiver such as (while t (apply-x (msg-recv "app#built.x"))). Callable from any
/// thread, including from inside another lane's job. Returns the resolved key.
std::string post_lane_task(async_lane &lane, state_exec::async_scheduler &sched,
const std::string &root, const std::string &channel, lane_task_work work,
Expand Down
53 changes: 29 additions & 24 deletions src/cvc/tests/async_lane_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -474,7 +474,7 @@ TEST_F(LaneTaskTest, ParkAndFutureResumeTheCallerWithTheLaneResult) {
std::mutex m;
std::vector<std::thread::id> work_threads;
async_lane lane("work");
proc_ctx &c = make_ctx("demo3");
proc_ctx &c = make_ctx("app");
const std::string root = c.ictx.root_path;
auto make_work = [&](std::int64_t n) {
return [&, n] {
Expand All @@ -487,12 +487,12 @@ TEST_F(LaneTaskTest, ParkAndFutureResumeTheCallerWithTheLaneResult) {
};
se::builtins::register_fn(
c.env, "lane-double", [&, root](std::span<const se::value_t> args) -> se::value_t {
return cvc::park_on_lane_task(lane, sched, root, "demo3",
return cvc::park_on_lane_task(lane, sched, root, "app",
make_work(std::get<std::int64_t>(args[0].v)), err_value);
});
se::builtins::register_fn(
c.env, "lane-double-async", [&, root](std::span<const se::value_t> args) -> se::value_t {
return cvc::future_lane_task(lane, sched, root, "demo3",
return cvc::future_lane_task(lane, sched, root, "app",
make_work(std::get<std::int64_t>(args[0].v)), err_value);
});
const int pid = spawn(c, "(+ (lane-double 21) (await (lane-double-async 10)))");
Expand All @@ -507,13 +507,13 @@ TEST_F(LaneTaskTest, ParkAndFutureResumeTheCallerWithTheLaneResult) {
EXPECT_NE(id, std::this_thread::get_id()) << "lane work must run off the pump thread";
}

TEST_F(LaneTaskTest, ProcessChrootedToDemo3ReceivesCrossThreadHashPosts) {
TEST_F(LaneTaskTest, ProcessChrootedToAppReceivesCrossThreadHashPosts) {
std::vector<std::int64_t> got; // written by apply-built, which runs on the pump (this thread)
std::vector<std::string> plain; // likewise, for the non-'#' channel
async_lane geom("geom");
async_lane sim("sim"); // destroyed first: its jobs submit to geom
proc_ctx &c = make_ctx("demo3");
proc_ctx &c2 = make_ctx("demo3");
proc_ctx &c = make_ctx("app");
proc_ctx &c2 = make_ctx("app");
se::builtins::register_fn(c.env, "apply-built", [&got](std::span<const se::value_t> args) {
got.push_back(std::get<std::int64_t>(args[0].v));
return se::value_t();
Expand All @@ -522,34 +522,39 @@ TEST_F(LaneTaskTest, ProcessChrootedToDemo3ReceivesCrossThreadHashPosts) {
plain.push_back(std::get<std::string>(args[0].v));
return se::value_t();
});
const int pid = spawn(c, "(while t (apply-built (msg-recv \"demo3#built.tracks\")))");
const int pid = spawn(c, "(while t (apply-built (msg-recv \"app#built.tracks\")))");
const int pid2 = spawn(c2, "(while t (apply-plain (msg-recv \"built\")))");
sched.sync_run(1000, 1.0);
ASSERT_EQ(status(pid), se::process_status::waiting);
ASSERT_EQ(status(pid2), se::process_status::waiting);

// The sim lane hands each generation straight to the geom lane (no main-thread hop), whose job
// posts to a FIXED '#' channel; a '#' channel survives the "demo3" chroot verbatim.
// posts to a FIXED '#' channel; a '#' channel survives the "app" chroot verbatim.
for (int gen = 1; gen <= 50; ++gen)
ASSERT_TRUE(sim.submit([this, &geom, gen] {
const std::string key = cvc::post_lane_task(
geom, sched, "demo3", "demo3#built.tracks",
geom, sched, "app", "app#built.tracks",
[gen] { return se::value_t(static_cast<std::int64_t>(gen)); }, err_value);
EXPECT_EQ(key, "demo3#built.tracks");
EXPECT_EQ(key, "app#built.tracks");
}));
// A non-'#' channel is scoped exactly as the chrooted (msg-recv "built") resolves it.
const std::string scoped = cvc::post_lane_task(
geom, sched, "demo3", "built", [] { return se::value_t("scoped"); }, err_value);
EXPECT_EQ(scoped, "demo3.channels.built");
geom, sched, "app", "built", [] { return se::value_t("scoped"); }, err_value);
EXPECT_EQ(scoped, "app.channels.built");

ASSERT_TRUE(pump_until([&] { return got.size() == 50 && plain.size() == 1; }))
<< "got " << got.size() << " of 50 '#' posts, " << plain.size() << " of 1 scoped post";
for (std::size_t i = 0; i < got.size(); ++i)
EXPECT_EQ(got[i], static_cast<std::int64_t>(i + 1)) << "lane results reordered end to end";
EXPECT_EQ(plain.front(), "scoped");
// Still resident: parked for the next generation.
EXPECT_EQ(status(pid), se::process_status::waiting);
EXPECT_EQ(status(pid2), se::process_status::waiting);
// Still resident: both re-park for the NEXT generation. Reaching got==50 only means the 50th
// value was applied INSIDE apply-built; the process re-enters (msg-recv …) a step later, so it is
// momentarily `ready`, not yet `waiting`. Pump until both have settled back to waiting rather
// than asserting it immediately — that race intermittently read `ready` (status 0) and flaked.
EXPECT_TRUE(pump_until([&] {
return status(pid) == se::process_status::waiting &&
status(pid2) == se::process_status::waiting;
})) << "processes did not re-park on msg-recv";
}

TEST_F(LaneTaskTest, EveryTaskPostsExactlyOneValueOnEveryErrorPath) {
Expand Down Expand Up @@ -593,11 +598,11 @@ TEST_F(LaneTaskTest, EveryTaskPostsExactlyOneValueOnEveryErrorPath) {
TEST_F(LaneTaskTest, LaunchReturnsAUniqueHashChannelPerTask) {
async_lane lane("sim");
const auto one = [] { return se::value_t(static_cast<std::int64_t>(1)); };
const std::string a = cvc::launch_lane_task(lane, sched, "demo3", "demo3", one, err_value);
const std::string b = cvc::launch_lane_task(lane, sched, "demo3", "demo3", one, err_value);
const std::string a = cvc::launch_lane_task(lane, sched, "app", "app", one, err_value);
const std::string b = cvc::launch_lane_task(lane, sched, "app", "app", one, err_value);
EXPECT_NE(a, b);
EXPECT_EQ(a.rfind("demo3#sim.", 0), 0u) << a;
EXPECT_EQ(b.rfind("demo3#sim.", 0), 0u) << b;
EXPECT_EQ(a.rfind("app#sim.", 0), 0u) << a;
EXPECT_EQ(b.rfind("app#sim.", 0), 0u) << b;
lane.stop();
sched.drain_ingress();
EXPECT_TRUE(take(a).has_value());
Expand All @@ -608,7 +613,7 @@ TEST_F(LaneTaskTest, LaunchReturnsAUniqueHashChannelPerTask) {
// and then yields a frame advances exactly one tick per pump, and never without one.
TEST_F(LaneTaskTest, PumpAdvancesAnAwaitFrameYielderOncePerFrame) {
int ticks = 0; // written by tick, which runs on the pump (this thread)
proc_ctx &c = make_ctx("demo3");
proc_ctx &c = make_ctx("app");
se::builtins::register_fn(c.env, "tick", [&ticks](std::span<const se::value_t>) {
++ticks;
return se::value_t();
Expand All @@ -632,12 +637,12 @@ TEST_F(LaneTaskTest, PumpAdvancesAnAwaitFrameYielderOncePerFrame) {
TEST_F(LaneTaskTest, SleeperIsWokenByLaneCompletionWhenNothingElseIsRunnable) {
std::atomic<int> batches{0};
async_lane lane("sim");
proc_ctx &c = make_ctx("demo3");
proc_ctx &c = make_ctx("app");
const std::string root = c.ictx.root_path;
se::builtins::register_fn(c.env, "sim-launch",
[&, root](std::span<const se::value_t>) -> se::value_t {
return se::value_t(cvc::launch_lane_task(
lane, sched, root, "demo3",
lane, sched, root, "app",
[&batches] {
++batches;
return se::value_t(0.001); // seconds until the next batch
Expand All @@ -661,12 +666,12 @@ TEST_F(LaneTaskTest, DeferredLaneDrivesTheSameResidentFromThePump) {
int batches = 0;
std::thread::id work_thread;
async_lane lane("sim", lane_mode::deferred);
proc_ctx &c = make_ctx("demo3");
proc_ctx &c = make_ctx("app");
const std::string root = c.ictx.root_path;
se::builtins::register_fn(c.env, "sim-launch",
[&, root](std::span<const se::value_t>) -> se::value_t {
return se::value_t(cvc::launch_lane_task(
lane, sched, root, "demo3",
lane, sched, root, "app",
[&] {
++batches;
work_thread = std::this_thread::get_id();
Expand Down
Loading