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:
parent
7a614420f4
commit
5bd8813b9b
3 changed files with 126 additions and 1 deletions
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Reference in a new issue