diff --git a/runtime/src/main.c b/runtime/src/main.c
index e439904..56cb4a4 100644
--- a/runtime/src/main.c
+++ b/runtime/src/main.c
@@ -23,7 +23,13 @@
#define WO_VERSION "0.0.0-dev"
#endif
-static wo_vm VM; /* 32K value stack: keep it off the C stack */
+/* The shard array: index 0 is the primary (runs the entry, owns the
+ * database); the arc's default is ONE VM PER CORE (WO_SHARDS overrides,
+ * =1 is the serial escape hatch). Static: each wo_vm carries its 32K
+ * value stack, kept off the C stack. */
+#define WO_MAX_SHARDS 64u
+static wo_vm SHARDS[WO_MAX_SHARDS];
+#define VM (SHARDS[0])
static wo_db DB; /* the per-shard engine (one shard until iteration 8) */
static wo_wal WAL;
@@ -159,6 +165,28 @@ int main(int argc, char **argv) {
wo_module_free(&mod);
return 2;
}
+ VM.shard_id = 0;
+ VM.is_primary = 1;
+ /* the arc's stage 2: all cores by default (the brave landing), one
+ * pinned worker vm per extra core; WO_SHARDS caps or forces it */
+ {
+ long cores = sysconf(_SC_NPROCESSORS_ONLN);
+ uint32_t nshards = cores > 0 ? (uint32_t)cores : 1;
+ const char *se = getenv("WO_SHARDS");
+ if (se && se[0]) {
+ unsigned long v = strtoul(se, NULL, 10);
+ if (v >= 1 && v <= WO_MAX_SHARDS) nshards = (uint32_t)v;
+ }
+ if (nshards > WO_MAX_SHARDS) nshards = WO_MAX_SHARDS;
+ wo_eng.shards = SHARDS;
+ if (wo_engine_start(&mod, heap_mb << 20, nshards) != 0) {
+ fprintf(stderr, "wovm: cannot start %u shards\n", nshards);
+ wo_engine_stop();
+ wo_vm_destroy(&VM);
+ wo_module_free(&mod);
+ return 2;
+ }
+ }
/* The database engine boots with the VM: every class IS a table.
* Durability is opt-in — WO_DATA=
opens /shard-0.wal,
* replays it before the entry runs (boot-before-listeners doctrine),
@@ -242,6 +270,7 @@ int main(int argc, char **argv) {
* heap is torn down, and after a trap too: the container outlives the
* unwind. */
if (argv_val) wo_drop_kind(&VM.rt, WO_K_MULTI, argv_val);
+ wo_engine_stop(); /* join + destroy the worker shards before the primary */
if (VM.rt.wal) wo_wal_close(&WAL);
wo_db_destroy(&DB);
gc_pump(&VM);
diff --git a/runtime/src/vm.c b/runtime/src/vm.c
index 1181a6d..33dfcc4 100644
--- a/runtime/src/vm.c
+++ b/runtime/src/vm.c
@@ -1,3 +1,4 @@
+#define _GNU_SOURCE /* pthread_setaffinity_np, CPU_SET */
#include "vm.h"
#include
@@ -11,8 +12,80 @@
#include "gc.h"
#include "park.h"
+#include
+#include
+#include
+#include
+
uint32_t wo_vm_depth(const wo_vm *vm) { return vm->cur->depth; }
+/* ---- the shard engine (arc stage 2, T5: threads exist and idle) -------- */
+
+wo_engine wo_eng = {0};
+
+static _Atomic int eng_shutdown = 0;
+
+/* A worker's whole life in T5: pinned, parked on its wake eventfd until
+ * shutdown. T6 gives it an inbox to adopt fibers from and the serve loop
+ * that runs them. */
+static void *shard_main(void *arg) {
+ wo_vm *vm = (wo_vm *)arg;
+ cpu_set_t set;
+ CPU_ZERO(&set);
+ CPU_SET((int)vm->shard_id, &set);
+ pthread_setaffinity_np(pthread_self(), sizeof set, &set);
+ while (!eng_shutdown) {
+ uint64_t v = 0;
+ ssize_t n = read(vm->wake_efd, &v, sizeof v); /* blocks until woken */
+ (void)n;
+ }
+ return NULL;
+}
+
+int wo_engine_start(const wo_module *mod, size_t heap_cap, uint32_t nshards) {
+ wo_eng.nshards = nshards;
+ if (nshards <= 1) return 0; /* the one-shard degenerate case: no threads */
+ pthread_t *ts = calloc(nshards - 1, sizeof(pthread_t));
+ if (!ts) return -1;
+ wo_eng.threads = ts;
+ for (uint32_t i = 1; i < nshards; i++) {
+ wo_vm *sv = &wo_eng.shards[i];
+ /* LAZY: a worker's full vm (64 MiB arena and all) is not paid for
+ * until its first fiber arrives (T6 adopts). T5 workers only need
+ * an identity and a wake fd — 20 idle shards must not cost 1.25 GiB
+ * (they did: the web-app gate flaked on exactly that). */
+ memset(sv, 0, sizeof *sv);
+ sv->mod = mod;
+ sv->shard_id = i;
+ sv->is_primary = 0;
+ sv->wake_efd = eventfd(0, 0);
+ if (sv->wake_efd < 0) return -1;
+ if (pthread_create(&ts[i - 1], NULL, shard_main, sv) != 0) return -1;
+ }
+ (void)heap_cap; /* consumed at lazy init (T6) */
+ return 0;
+}
+
+void wo_engine_stop(void) {
+ if (wo_eng.nshards <= 1) return;
+ eng_shutdown = 1;
+ pthread_t *ts = (pthread_t *)wo_eng.threads;
+ for (uint32_t i = 1; i < wo_eng.nshards; i++) {
+ uint64_t one = 1;
+ ssize_t n = write(wo_eng.shards[i].wake_efd, &one, sizeof one);
+ (void)n;
+ }
+ for (uint32_t i = 1; i < wo_eng.nshards; i++) pthread_join(ts[i - 1], NULL);
+ for (uint32_t i = 1; i < wo_eng.nshards; i++) {
+ close(wo_eng.shards[i].wake_efd);
+ if (wo_eng.shards[i].rt.arena.base) /* lazily init'ed only */
+ wo_vm_destroy(&wo_eng.shards[i]);
+ }
+ free(ts);
+ wo_eng.threads = NULL;
+ wo_eng.nshards = 1;
+}
+
int wo_vm_init(wo_vm *vm, const wo_module *mod, size_t heap_cap) {
memset(vm, 0, sizeof(*vm));
vm->mod = mod;
diff --git a/runtime/src/vm.h b/runtime/src/vm.h
index e7c0a7b..09d0f82 100644
--- a/runtime/src/vm.h
+++ b/runtime/src/vm.h
@@ -96,6 +96,11 @@ typedef struct wo_actor {
typedef struct wo_vm {
const wo_module *mod;
wo_rt rt;
+ /* the arc's stage 2: which shard this vm IS. Shard 0 is the primary
+ * (runs the entry, owns the database); workers run wo_vm_serve. */
+ uint32_t shard_id;
+ int is_primary;
+ int wake_efd; /* wakes an idle worker (inbox arrivals, shutdown) */
wo_fiber f0; /* fiber 0: main — embedded; spawned fibers are calloc'd */
wo_fiber *cur; /* the live fiber — every interpreter access goes here */
wo_fiber *qhead, *qtail; /* RUNNABLE fibers awaiting the interpreter */
@@ -117,6 +122,24 @@ int wo_vm_actor_spawn(wo_vm *vm, uint64_t instance, uint32_t method_idx,
uint64_t *out_addr, const char **msg);
int wo_vm_actor_send(wo_vm *vm, uint64_t addr, uint64_t msg_val, const char **msg);
+/* ---- the shard engine (arc stage 2) ------------------------------------
+ * One pinned thread per shard, each a full wo_vm (own arena, GC, I/O
+ * plane). Shard 0 is the caller's (main's); workers idle on their wake
+ * eventfd until fibers arrive (stage 2 T6) or shutdown. Count: WO_SHARDS
+ * or all cores (the arc's default). */
+typedef struct wo_engine {
+ wo_vm *shards; /* [nshards]; index 0 = primary */
+ void *threads; /* pthread_t[nshards-1], opaque here (libc-only header) */
+ uint32_t nshards;
+} wo_engine;
+
+extern wo_engine wo_eng; /* the process's one engine (vm.c) */
+
+/* Start shards 1..n-1 (0 is the caller's, already init'ed in shards[0]).
+ * 0 ok. Stop joins every worker and destroys their vms. */
+int wo_engine_start(const wo_module *mod, size_t heap_cap, uint32_t nshards);
+void wo_engine_stop(void);
+
/* Spawn a fiber that will run method_idx(args) — the runtime half the
* `spawn` expression lowers onto (stage 1 Task 3); Task 2's tests drive it
* directly. The fiber is RUNNABLE and queued; it runs when the scheduler