From 5bd8813b9b759c9c8cfa0c229cd0def387c6fde3 Mon Sep 17 00:00:00 2001 From: "shoney.arickathil" Date: Thu, 20 Aug 2026 11:46:28 +0200 Subject: [PATCH] =?UTF-8?q?feat(runtime):=20shard=20engine=20=E2=80=94=20p?= =?UTF-8?q?inned=20worker=20threads,=20all=20cores=20default=20(arc=20T5)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - wo_engine + wo_engine_start/stop: one pinned pthread per extra core (pthread_setaffinity_np); shard 0 is the primary (entry + database); workers idle on a wake eventfd until fibers arrive (T6) or shutdown - default = all cores (the arc's brave landing), WO_SHARDS=1..64 overrides; N=1 spawns no threads — byte-identical to stage 1 - worker vms are LAZY: identity + wake fd only until their first fiber arrives — 20 idle shards must not cost 1.25 GiB of eager arenas (they did: the web-app gate flaked on exactly that before the fix; 3 consecutive green runs after) - engine stops (join + destroy) before the primary's teardown - battery at the 20-core default: oop-e2e 92/0, fibers 8/0, log-watcher 7/0 (+WO_SHARDS=1 identical), web-app 21/0 x3, employee 8/0, deps 8/0, runtime tests 16/16 files green Co-Authored-By: Claude Opus 5 (1M context) --- runtime/src/main.c | 31 +++++++++++++++++++- runtime/src/vm.c | 73 ++++++++++++++++++++++++++++++++++++++++++++++ runtime/src/vm.h | 23 +++++++++++++++ 3 files changed, 126 insertions(+), 1 deletion(-) 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