feat(runtime): shard engine — pinned worker threads, all cores default (arc T5)

- 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) <noreply@anthropic.com>
This commit is contained in:
shoney.arickathil 2026-08-20 11:46:28 +02:00
parent efbf2f36ca
commit fb74e41526
3 changed files with 126 additions and 1 deletions

View file

@ -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=<dir> opens <dir>/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);

View file

@ -1,3 +1,4 @@
#define _GNU_SOURCE /* pthread_setaffinity_np, CPU_SET */
#include "vm.h"
#include <stdarg.h>
@ -11,8 +12,80 @@
#include "gc.h"
#include "park.h"
#include <pthread.h>
#include <sched.h>
#include <sys/eventfd.h>
#include <unistd.h>
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;

View file

@ -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