feat: cross-shard actors — placement, envelopes, home-routed frees, WO-E222 (arc T6)

- spawn placement: round-robin across shards (same-shard when the
  engine is absent/single); the actor's mailbox and delivery belong to
  its HOME thread — a spawn to another shard travels as an ADOPT
  envelope, a send as a SEND envelope (mutex-guarded inbox + eventfd
  wake; the spec's lock-free rings stay a disclosed deviation until
  9e measures the mutex)
- workers: first envelope triggers lazy full-vm init UNDER the inbox
  mutex (TSan caught the memset racing a concurrent push, twice — the
  second was inbox_push reading wake_efd outside the lock; both fixed,
  gate x8 + battery clean); serve loop = adopt -> run to drained ->
  wait on the plane (the wake eventfd is watched by io_uring POLL_ADD
  oneshot / epoll level-triggered on BOTH backends)
- ownership across heaps: every allocation stamps rt->shard_id into
  the header (the field reserved since iteration 2); a drop on the
  wrong shard routes home as a FREE envelope — the owner's arena stays
  single-threaded by construction; at teardown routed frees become
  no-ops (arenas die wholesale) which is what un-danced the freed-mutex
  ASan SEGV the first ordering had
- WO-E222: an actor's state or message type that is (or transitively
  contains) an inferred-traced class refuses at the spawn/send — with
  round-robin every actor is potentially remote; corpus-pinned
  (compile-fail/traced-send, inference-aware: Box contains ?Node)
- determinism narrowed per spec: oop-e2e pins WO_SHARDS=1 (exact
  outputs); the fibers gate grows multi-shard SET assertions + a TSan
  run (wovm-tsan target; setarch -R fallback for kernel 6.5+ ASLR)
- NEXT_RUNNABLE honors engine shutdown for parked workers (deadlock
  hole closed); io_wait's adopt-wake (rc 1) no longer reads as fatal
- battery: oop-e2e 93/0, fibers 10/0 x8 (+WO_IO=epoll), log-watcher
  7/0, employee 8/0, web-app 21/0, deps 8/0, runtime tests 16/16

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
shoney.arickathil 2026-08-20 12:06:13 +02:00
parent 5bd8813b9b
commit ec9d264355
14 changed files with 493 additions and 25 deletions

View file

@ -394,6 +394,7 @@ let module_not_imported_code = Diag.types_prefix ^ "10"
let nullable_used_without_check_code = Diag.types_prefix ^ "11" let nullable_used_without_check_code = Diag.types_prefix ^ "11"
let nullable_assign_mismatch_code = Diag.types_prefix ^ "12" let nullable_assign_mismatch_code = Diag.types_prefix ^ "12"
let spawn_no_receive_code = Diag.types_prefix ^ "21" (* WO-E221: spawn target lacks fn receive(msg: M); E219/E220 are taken on the language-surface-strictness branch *) let spawn_no_receive_code = Diag.types_prefix ^ "21" (* WO-E221: spawn target lacks fn receive(msg: M); E219/E220 are taken on the language-surface-strictness branch *)
let traced_send_code = Diag.types_prefix ^ "22" (* WO-E222: traced(-containing) type in an actor message or actor state — aliased graphs cannot cross heap boundaries *)
let missing_nil_check_code = Diag.types_prefix ^ "13" let missing_nil_check_code = Diag.types_prefix ^ "13"
(* haxe-parity Task 1 (modules). module_not_imported_code (WO-E210, (* haxe-parity Task 1 (modules). module_not_imported_code (WO-E210,
@ -852,6 +853,43 @@ let req_label = function
| ReqContainer -> "a `multi` or `map`" | ReqContainer -> "a `multi` or `map`"
| ReqAny -> "any" (* matches_req is always true here -- never rendered *) | ReqAny -> "any" (* matches_req is always true here -- never rendered *)
(* arc T6 (WO-E222): does this type name a traced class, or a class/union
* that transitively CONTAINS one? An actor's state and messages may cross
* heap boundaries (placement is round-robin — every spawn/send may cross),
* and aliased graphs cannot: their lifetime is one shard's collector's.
* Fixpoint over the class graph, memoized per query via a visited set. *)
let contains_traced (syms : symbols) (root : string) : bool =
let rec go (seen : StringSet.t) (name : string) : bool =
if StringSet.mem name seen then false
else if is_gc_class syms name then true
else
let seen = StringSet.add name seen in
let field_hits fields =
List.exists
(fun (_, ft, _, _) ->
match unwrap_nullable (typ_of_field_ty ft) with
| TScalar n | TMulti (TScalar n) | TMap (_, TScalar n) -> go seen n
| _ -> false)
fields
in
match StringMap.find_opt name syms.classes with
| Some cls -> field_hits cls.fields
| None -> (
match StringMap.find_opt name syms.unions with
| Some u ->
List.exists
(fun (v : variant_info) ->
List.exists
(fun (_, ft) ->
match unwrap_nullable (typ_of_field_ty ft) with
| TScalar n | TMulti (TScalar n) | TMap (_, TScalar n) -> go seen n
| _ -> false)
v.vi_fields)
u.u_variants
| None -> false)
in
go StringSet.empty root
let rec typ_label (t : typ) : string = let rec typ_label (t : typ) : string =
match t with match t with
| TScalar name -> name | TScalar name -> name
@ -1627,6 +1665,15 @@ let typecheck_program ~file ~(module_of : string -> string)
()); ());
{ typ = TScalar "Int"; is_nil = false } { typ = TScalar "Int"; is_nil = false }
in in
let e222 (pos : pos) (what : string) (tname : string) : unit =
Diag.Collector.add collector
(Diag.error ~code:traced_send_code ~file ~line:pos.line ~col:pos.col
~message:
(Printf.sprintf
"%s type `%s` is traced (or contains a traced class) — aliased graphs cannot cross shard heaps; spawn placement makes every actor potentially remote"
what tname)
())
in
(match StringMap.find_opt cn syms.classes with (match StringMap.find_opt cn syms.classes with
| Some cls -> ( | Some cls -> (
match List.find_opt (fun (m : method_info) -> m.name = "receive") cls.methods with match List.find_opt (fun (m : method_info) -> m.name = "receive") cls.methods with
@ -1635,6 +1682,8 @@ let typecheck_program ~file ~(module_of : string -> string)
| TScalar mname | TScalar mname
when StringMap.mem mname syms.classes when StringMap.mem mname syms.classes
|| StringMap.mem mname syms.unions -> || StringMap.mem mname syms.unions ->
if contains_traced syms cn then e222 e.pos "actor state" cn;
if contains_traced syms mname then e222 e.pos "message" mname;
{ typ = TActor mname; is_nil = false } { typ = TActor mname; is_nil = false }
| TScalar mname -> bad (Printf.sprintf "receive's message type `%s` is not a declared class, record, or union" mname) | TScalar mname -> bad (Printf.sprintf "receive's message type `%s` is not a declared class, record, or union" mname)
| _ -> bad "receive's parameter must be a plain class, record, or union type") | _ -> bad "receive's parameter must be a plain class, record, or union type")

View file

@ -58,6 +58,12 @@ build/wovm_asan: src/main.c $(VMSRC) $(VMHDR) | build
wovm-asan: build/wovm_asan wovm-asan: build/wovm_asan
# TSan flavor (arc stage 2): the cross-shard proofs run under this
build/wovm_tsan: src/main.c $(VMSRC) $(VMHDR) | build
$(CC) $(CFLAGS) -fsanitize=thread -g -Isrc -I../database/src -o $@ src/main.c $(VMSRC) $(LDFLAGS)
wovm-tsan: build/wovm_tsan
# fixture generator for the CLI smoke test # fixture generator for the CLI smoke test
build/mkwob: test/mkwob.c test/wob_build.c $(VMHDR) | build build/mkwob: test/mkwob.c test/wob_build.c $(VMHDR) | build
$(CC) $(CFLAGS) -Isrc -Itest -o $@ test/mkwob.c test/wob_build.c $(CC) $(CFLAGS) -Isrc -Itest -o $@ test/mkwob.c test/wob_build.c
@ -74,4 +80,4 @@ clean:
rm -f wo-rt bench/bench wovm rm -f wo-rt bench/bench wovm
rm -rf build rm -rf build
.PHONY: bench run clean test test-iso wovm-asan .PHONY: bench run clean test test-iso wovm-asan wovm-tsan

View file

@ -8,6 +8,7 @@ wo_multi *wo_multi_new(wo_rt *rt, uint8_t elem_kind) {
if (!m) return NULL; if (!m) return NULL;
memset(m, 0, sizeof(*m)); memset(m, 0, sizeof(*m));
m->h.class_id = WO_CLS_MULTI; m->h.class_id = WO_CLS_MULTI;
m->h.shard_id = rt->shard_id;
m->elem_kind = elem_kind; m->elem_kind = elem_kind;
return m; return m;
} }
@ -35,6 +36,7 @@ wo_map *wo_map_new(wo_rt *rt, uint8_t key_kind, uint8_t val_kind) {
if (!m) return NULL; if (!m) return NULL;
memset(m, 0, sizeof(*m)); memset(m, 0, sizeof(*m));
m->h.class_id = WO_CLS_MAP; m->h.class_id = WO_CLS_MAP;
m->h.shard_id = rt->shard_id;
m->key_kind = key_kind; m->key_kind = key_kind;
m->val_kind = val_kind; m->val_kind = val_kind;
return m; return m;

View file

@ -68,7 +68,16 @@ static void class_free(wo_rt *rt, wo_hdr *o) {
wo_arena_free(&rt->arena, o, wo_obj_size(c)); wo_arena_free(&rt->arena, o, wo_obj_size(c));
} }
/* arc T6: a drop on the wrong shard routes home — the owner's arena is
* single-threaded by doctrine, so the free travels as an envelope.
* (vm.c owns the engine; this hook keeps gc.c engine-blind.) */
void wo_route_free(wo_hdr *h);
void wo_drop_obj(wo_rt *rt, wo_hdr *o) { void wo_drop_obj(wo_rt *rt, wo_hdr *o) {
if (o && o->shard_id != rt->shard_id && !(o->flags & WO_F_CONST)) {
wo_route_free(o);
return;
}
if (!o) return; if (!o) return;
switch (o->class_id) { switch (o->class_id) {
case WO_CLS_STR: case WO_CLS_STR:
@ -104,9 +113,15 @@ void wo_drop_kind(wo_rt *rt, uint8_t kind, uint64_t v) {
* edge's death means nothing. */ * edge's death means nothing. */
if (rt->gc_phase == WO_GC_MARK) wo_gc_shade(rt, (wo_hdr *)(uintptr_t)v); if (rt->gc_phase == WO_GC_MARK) wo_gc_shade(rt, (wo_hdr *)(uintptr_t)v);
return; return;
case WO_K_TEXT: case WO_K_TEXT: {
wo_str_free(rt, (wo_str *)(uintptr_t)v); wo_str *sp = (wo_str *)(uintptr_t)v;
if (sp->h.shard_id != rt->shard_id && !(sp->h.flags & WO_F_CONST)) {
wo_route_free(&sp->h); /* home arena frees it (arc T6) */
return;
}
wo_str_free(rt, sp);
return; return;
}
default: default:
return; /* loader guarantees kinds; defensive no-op */ return; /* loader guarantees kinds; defensive no-op */
} }

View file

@ -10,6 +10,8 @@
#include <sys/mman.h> #include <sys/mman.h>
#include <sys/stat.h> #include <sys/stat.h>
#include <unistd.h> #include <unistd.h>
#include <pthread.h>
#include <sys/eventfd.h>
#include "cont.h" #include "cont.h"
#include "gc.h" #include "gc.h"
@ -167,6 +169,17 @@ int main(int argc, char **argv) {
} }
VM.shard_id = 0; VM.shard_id = 0;
VM.is_primary = 1; VM.is_primary = 1;
VM.rt.shard_id = 0;
VM.wake_efd = eventfd(0, EFD_NONBLOCK);
VM.in_mu = calloc(1, sizeof(pthread_mutex_t));
if (!VM.in_mu || VM.wake_efd < 0
|| pthread_mutex_init((pthread_mutex_t *)VM.in_mu, NULL) != 0) {
fprintf(stderr, "wovm: cannot set up the primary shard\n");
wo_vm_destroy(&VM);
wo_module_free(&mod);
return 2;
}
wo_tls_set(&VM);
/* the arc's stage 2: all cores by default (the brave landing), one /* 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 */ * pinned worker vm per extra core; WO_SHARDS caps or forces it */
{ {

View file

@ -148,6 +148,7 @@ wo_hdr *wo_obj_new(wo_rt *rt, uint32_t class_id) {
if (!o) return NULL; if (!o) return NULL;
memset(o, 0, sz); /* color: WHITE by construction (all-zero) */ memset(o, 0, sz); /* color: WHITE by construction (all-zero) */
o->class_id = class_id; o->class_id = class_id;
o->shard_id = rt->shard_id;
if (c->flags & WO_CLASSF_GC) { if (c->flags & WO_CLASSF_GC) {
o->flags = WO_F_GC; o->flags = WO_F_GC;
/* born black while a cycle runs: live-at-birth for that cycle */ /* born black while a cycle runs: live-at-birth for that cycle */
@ -162,6 +163,7 @@ wo_str *wo_str_alloc(wo_rt *rt, uint32_t len) {
if (!s) return NULL; if (!s) return NULL;
memset(&s->h, 0, sizeof(s->h)); memset(&s->h, 0, sizeof(s->h));
s->h.class_id = WO_CLS_STR; s->h.class_id = WO_CLS_STR;
s->h.shard_id = rt->shard_id;
s->len = len; s->len = len;
return s; return s;
} }

View file

@ -36,6 +36,8 @@ enum { WO_GC_IDLE = 0, WO_GC_MARK = 1, WO_GC_SWEEP = 2 };
* state (gc.c), and the output stream builtin print writes to (tests point * state (gc.c), and the output stream builtin print writes to (tests point
* it at a temp file to capture output). */ * it at a temp file to capture output). */
typedef struct wo_rt { typedef struct wo_rt {
uint16_t shard_id; /* arc T6: stamped into every allocation's header;
a drop whose header disagrees routes home */
wo_arena arena; wo_arena arena;
const wo_classdesc *classes; const wo_classdesc *classes;
uint32_t class_cnt; uint32_t class_cnt;

View file

@ -7,6 +7,7 @@
#include "park.h" #include "park.h"
#include <errno.h> #include <errno.h>
#include <stdio.h>
#include <poll.h> #include <poll.h>
#include <stdlib.h> #include <stdlib.h>
#include <string.h> #include <string.h>
@ -231,10 +232,31 @@ int wo_io_arm(wo_vm *vm, wo_fiber *fb) {
return 0; return 0;
} }
/* user_data sentinel for the wake-eventfd's own readiness (fibers are
* heap pointers, never 1) */
#define EFD_SENTINEL 1ull
static void efd_drain(wo_vm *vm) {
uint64_t v = 0;
ssize_t n = read(vm->wake_efd, &v, sizeof v);
(void)n;
}
int wo_io_wait(wo_vm *vm) { int wo_io_wait(wo_vm *vm) {
for (;;) { for (;;) {
if (wo_sys_stop_pending()) return WO_IO_STOP; if (wo_sys_stop_pending()) return WO_IO_STOP;
if (vm->io_kind == 0) { if (vm->io_kind == 0) {
/* keep the wake eventfd armed (oneshot POLL_ADD, re-armed
* after each firing) so inbox pushes interrupt the wait */
if (vm->wake_efd >= 0 && !vm->efd_armed) {
struct io_uring_sqe sqe;
memset(&sqe, 0, sizeof sqe);
sqe.opcode = IORING_OP_POLL_ADD;
sqe.fd = vm->wake_efd;
sqe.poll32_events = POLLIN;
sqe.user_data = EFD_SENTINEL;
if (uring_submit(vm, &sqe) == 0) vm->efd_armed = 1;
}
rings r = ring_ptrs(vm); rings r = ring_ptrs(vm);
uint32_t head = *r.cq_head; uint32_t head = *r.cq_head;
uint32_t tail = __atomic_load_n(r.cq_tail, __ATOMIC_ACQUIRE); uint32_t tail = __atomic_load_n(r.cq_tail, __ATOMIC_ACQUIRE);
@ -242,24 +264,44 @@ int wo_io_wait(wo_vm *vm) {
long rc = syscall(SYS_io_uring_enter, vm->io_fd, 0u, 1u, long rc = syscall(SYS_io_uring_enter, vm->io_fd, 0u, 1u,
IORING_ENTER_GETEVENTS, NULL, 0); IORING_ENTER_GETEVENTS, NULL, 0);
if (rc < 0 && errno == EINTR) continue; /* stop checked on loop */ if (rc < 0 && errno == EINTR) continue; /* stop checked on loop */
if (rc < 0) return -1; if (rc < 0) {
fprintf(stderr, "DBG uring_enter shard=%u errno=%d\n", vm->shard_id, errno);
return -1;
}
tail = __atomic_load_n(r.cq_tail, __ATOMIC_ACQUIRE); tail = __atomic_load_n(r.cq_tail, __ATOMIC_ACQUIRE);
} }
int woke = 0; int woke = 0;
while (head != tail) { while (head != tail) {
struct io_uring_cqe *cqe = &r.cqes[head & *r.cq_mask]; struct io_uring_cqe *cqe = &r.cqes[head & *r.cq_mask];
wo_fiber *fb = (wo_fiber *)(uintptr_t)cqe->user_data; if (cqe->user_data == EFD_SENTINEL) {
if (fb && fb->state == WO_FIB_PARKED) { vm->efd_armed = 0;
wake(vm, fb); efd_drain(vm);
woke = 1; woke = 2; /* inbox wake: the caller adopts */
} else {
wo_fiber *fb = (wo_fiber *)(uintptr_t)cqe->user_data;
if (fb && fb->state == WO_FIB_PARKED) {
wake(vm, fb);
if (!woke) woke = 1;
}
} }
head++; head++;
} }
__atomic_store_n(r.cq_head, head, __ATOMIC_RELEASE); __atomic_store_n(r.cq_head, head, __ATOMIC_RELEASE);
if (woke == 2) return 1; /* adopt-needed */
if (woke) return 0; if (woke) return 0;
continue; continue;
} }
/* epoll: timeout from the nearest sleep deadline */ /* epoll: the wake eventfd is registered once, level-triggered
* (data.ptr NULL = the sentinel) */
if (vm->wake_efd >= 0 && !vm->efd_armed) {
struct epoll_event ev;
memset(&ev, 0, sizeof ev);
ev.events = EPOLLIN;
ev.data.ptr = NULL;
if (epoll_ctl(vm->io_fd, EPOLL_CTL_ADD, vm->wake_efd, &ev) == 0
|| errno == EEXIST)
vm->efd_armed = 1;
}
int timeout = -1; int timeout = -1;
int64_t now = now_ms(); int64_t now = now_ms();
for (wo_fiber *fb = vm->parked; fb; fb = fb->pnext) for (wo_fiber *fb = vm->parked; fb; fb = fb->pnext)
@ -271,14 +313,20 @@ int wo_io_wait(wo_vm *vm) {
struct epoll_event evs[16]; struct epoll_event evs[16];
int n = epoll_wait(vm->io_fd, evs, 16, timeout); int n = epoll_wait(vm->io_fd, evs, 16, timeout);
if (n < 0 && errno == EINTR) continue; if (n < 0 && errno == EINTR) continue;
if (n < 0) return -1; if (n < 0) {
fprintf(stderr, "DBG epoll_wait shard=%u errno=%d\n", vm->shard_id, errno);
return -1;
}
int woke = 0; int woke = 0;
for (int i = 0; i < n; i++) { for (int i = 0; i < n; i++) {
wo_fiber *fb = (wo_fiber *)evs[i].data.ptr; wo_fiber *fb = (wo_fiber *)evs[i].data.ptr;
if (fb && fb->state == WO_FIB_PARKED) { if (!fb) { /* the wake eventfd: adopt-needed */
efd_drain(vm);
woke = 2;
} else if (fb->state == WO_FIB_PARKED) {
epoll_ctl(vm->io_fd, EPOLL_CTL_DEL, fb->park_fd, NULL); epoll_ctl(vm->io_fd, EPOLL_CTL_DEL, fb->park_fd, NULL);
wake(vm, fb); wake(vm, fb);
woke = 1; if (!woke) woke = 1;
} }
} }
now = now_ms(); now = now_ms();
@ -291,6 +339,7 @@ int wo_io_wait(wo_vm *vm) {
} }
fb = nx; fb = nx;
} }
if (woke == 2) return 1; /* adopt-needed */
if (woke) return 0; if (woke) return 0;
} }
} }

View file

@ -13,7 +13,9 @@
#include "park.h" #include "park.h"
#include <pthread.h> #include <pthread.h>
#include <poll.h>
#include <sched.h> #include <sched.h>
#include <errno.h>
#include <sys/eventfd.h> #include <sys/eventfd.h>
#include <unistd.h> #include <unistd.h>
@ -24,26 +26,157 @@ uint32_t wo_vm_depth(const wo_vm *vm) { return vm->cur->depth; }
wo_engine wo_eng = {0}; wo_engine wo_eng = {0};
static _Atomic int eng_shutdown = 0; static _Atomic int eng_shutdown = 0;
static _Atomic int eng_teardown = 0; /* set once threads are joined: routed
frees become no-ops (every arena dies wholesale) and envelopes are
discarded, so teardown order cannot dangle a mutex */
static _Atomic uint32_t eng_rr = 0; /* round-robin spawn cursor */
static _Thread_local wo_vm *tls_vm = NULL;
wo_vm *wo_tls_vm(void) { return tls_vm; }
void wo_tls_set(wo_vm *vm) { tls_vm = vm; }
/* push an envelope into a shard's inbox and wake it (any thread) */
static void inbox_push(wo_vm *to, wo_envelope *e) {
pthread_mutex_t *mu = (pthread_mutex_t *)to->in_mu;
pthread_mutex_lock(mu);
e->next = NULL;
if (to->in_tail) to->in_tail->next = e;
else to->in_head = e;
to->in_tail = e;
/* capture the wake fd UNDER the lock: worker_late_init rewrites the
* whole vm under this same mutex (TSan caught the unlocked read) */
int efd = to->wake_efd;
pthread_mutex_unlock(mu);
if (efd >= 0) {
uint64_t one = 1;
ssize_t n = write(efd, &one, sizeof one);
(void)n;
}
}
static void fib_enqueue(wo_vm *vm, wo_fiber *fb);
static int actor_push(wo_actor *a, uint64_t m);
static int actor_activate(wo_vm *vm, wo_actor *a);
static void fib_reap_all(wo_vm *vm);
/* the owning thread drains its inbox: adopt actors, deliver sends,
* execute home-routed frees. Returns how many envelopes were handled. */
static int wo_vm_adopt(wo_vm *vm) {
pthread_mutex_lock((pthread_mutex_t *)vm->in_mu);
wo_envelope *e = vm->in_head;
vm->in_head = vm->in_tail = NULL;
pthread_mutex_unlock((pthread_mutex_t *)vm->in_mu);
int n = 0;
while (e) {
wo_envelope *nx = e->next;
switch (e->kind) {
case 1: /* adopt a freshly spawned actor: link it, nothing runs yet */
e->actor->next_all = vm->actors;
vm->actors = e->actor;
break;
case 0: /* a cross-shard send: mailbox + activation on the HOME thread */
if (actor_push(e->actor, e->payload) == 0 && !e->actor->active)
(void)actor_activate(vm, e->actor);
break;
case 2: /* a home-routed free: this arena owns the object */
wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)e->payload);
break;
}
free(e);
n++;
e = nx;
}
return n;
}
/* route a drop to the object's home shard (gc.c calls through this when
* the header's shard id is not the current thread's) */
void wo_route_free(wo_hdr *h) {
if (eng_teardown) return; /* arenas are torn down wholesale */
wo_vm *to = &wo_eng.shards[h->shard_id];
wo_envelope *e = calloc(1, sizeof *e);
if (!e) return; /* OOM on the free path: leak rather than crash */
e->kind = 2;
e->payload = (uint64_t)(uintptr_t)h;
inbox_push(to, e);
}
/* A worker's whole life in T5: pinned, parked on its wake eventfd until /* 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 * shutdown. T6 gives it an inbox to adopt fibers from and the serve loop
* that runs them. */ * that runs them. */
static size_t eng_heap_cap = 0;
/* lazily give a worker its full vm (arena, GC, I/O plane) — paid on the
* first envelope, not at boot (20 idle shards must stay ~free) */
static int worker_late_init(wo_vm *vm) {
if (vm->rt.arena.base) return 0;
/* under the inbox mutex: wo_vm_init memsets the whole vm, and a
* concurrent inbox_push would race the in_head/in_tail wipe (TSan
* caught exactly this). The mutex OBJECT is malloc'd and stable;
* pushers block on it while the fields are rebuilt. */
void *mu = vm->in_mu;
pthread_mutex_lock((pthread_mutex_t *)mu);
const wo_module *mod = vm->mod;
uint32_t id = vm->shard_id;
int efd = vm->wake_efd;
wo_envelope *h = vm->in_head, *t = vm->in_tail;
int rc = wo_vm_init(vm, mod, eng_heap_cap);
if (rc == 0) {
vm->shard_id = id;
vm->rt.shard_id = (uint16_t)id;
vm->is_primary = 0;
vm->wake_efd = efd;
vm->in_mu = mu;
vm->in_head = h;
vm->in_tail = t;
tls_vm = vm;
}
pthread_mutex_unlock((pthread_mutex_t *)mu);
return rc;
}
int wo_vm_serve(wo_vm *vm); /* vm_run's worker flavor, defined below it */
/* one blocking wait for the FIRST envelope (the vm — and its I/O plane —
* does not exist yet); after late init the plane's own wait watches the
* eventfd and this poll never runs again */
static void worker_first_wait(wo_vm *vm) {
struct pollfd p = { .fd = vm->wake_efd, .events = POLLIN };
while (!eng_shutdown) {
int n = poll(&p, 1, -1);
if (n > 0 || (n < 0 && errno != EINTR)) return;
if (wo_sys_stop_pending()) return;
}
}
static void *shard_main(void *arg) { static void *shard_main(void *arg) {
wo_vm *vm = (wo_vm *)arg; wo_vm *vm = (wo_vm *)arg;
tls_vm = vm;
cpu_set_t set; cpu_set_t set;
CPU_ZERO(&set); CPU_ZERO(&set);
CPU_SET((int)vm->shard_id, &set); CPU_SET((int)(vm->shard_id % 64u), &set);
pthread_setaffinity_np(pthread_self(), sizeof set, &set); pthread_setaffinity_np(pthread_self(), sizeof set, &set);
worker_first_wait(vm);
if (eng_shutdown || worker_late_init(vm) != 0) return NULL;
while (!eng_shutdown) { while (!eng_shutdown) {
uint64_t v = 0; (void)wo_vm_adopt(vm);
ssize_t n = read(vm->wake_efd, &v, sizeof v); /* blocks until woken */ if (vm->qhead) {
(void)n; int rc = wo_vm_serve(vm); /* runs until drained (2) or stop */
if (rc == 1) break; /* stop: everything reaped inside */
} else {
int rc = wo_io_wait(vm); /* parked fibers AND the wake eventfd */
if (rc == WO_IO_STOP) {
fib_reap_all(vm);
break;
}
}
} }
return NULL; return NULL;
} }
int wo_engine_start(const wo_module *mod, size_t heap_cap, uint32_t nshards) { int wo_engine_start(const wo_module *mod, size_t heap_cap, uint32_t nshards) {
wo_eng.nshards = nshards; wo_eng.nshards = nshards;
eng_heap_cap = heap_cap;
if (nshards <= 1) return 0; /* the one-shard degenerate case: no threads */ if (nshards <= 1) return 0; /* the one-shard degenerate case: no threads */
pthread_t *ts = calloc(nshards - 1, sizeof(pthread_t)); pthread_t *ts = calloc(nshards - 1, sizeof(pthread_t));
if (!ts) return -1; if (!ts) return -1;
@ -58,8 +191,11 @@ int wo_engine_start(const wo_module *mod, size_t heap_cap, uint32_t nshards) {
sv->mod = mod; sv->mod = mod;
sv->shard_id = i; sv->shard_id = i;
sv->is_primary = 0; sv->is_primary = 0;
sv->wake_efd = eventfd(0, 0); sv->wake_efd = eventfd(0, EFD_NONBLOCK);
if (sv->wake_efd < 0) return -1; if (sv->wake_efd < 0) return -1;
sv->in_mu = calloc(1, sizeof(pthread_mutex_t));
if (!sv->in_mu || pthread_mutex_init((pthread_mutex_t *)sv->in_mu, NULL) != 0)
return -1;
if (pthread_create(&ts[i - 1], NULL, shard_main, sv) != 0) return -1; if (pthread_create(&ts[i - 1], NULL, shard_main, sv) != 0) return -1;
} }
(void)heap_cap; /* consumed at lazy init (T6) */ (void)heap_cap; /* consumed at lazy init (T6) */
@ -76,10 +212,57 @@ void wo_engine_stop(void) {
(void)n; (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++) pthread_join(ts[i - 1], NULL);
/* single-threaded from here. Every arena dies wholesale, so routed
* frees and queued payloads need no per-object drops — DISCARD the
* envelopes (freeing the malloc'd nodes/actors) and let the arenas
* take their contents with them. The flag also turns any route_free
* raised by the destroys below into a no-op, so no teardown ordering
* can lock a freed mutex (the ASan SEGV this replaces). */
eng_teardown = 1;
for (uint32_t i = 0; i < wo_eng.nshards; i++) {
wo_vm *sv = &wo_eng.shards[i];
if (!sv->in_mu) continue;
wo_envelope *e = sv->in_head;
sv->in_head = sv->in_tail = NULL;
while (e) {
wo_envelope *nx = e->next;
if (e->kind == 1 && e->actor) {
free(e->actor->msgs);
free(e->actor);
}
free(e);
e = nx;
}
}
for (uint32_t i = 1; i < wo_eng.nshards; i++) { for (uint32_t i = 1; i < wo_eng.nshards; i++) {
close(wo_eng.shards[i].wake_efd); close(wo_eng.shards[i].wake_efd);
if (wo_eng.shards[i].rt.arena.base) /* lazily init'ed only */ if (wo_eng.shards[i].rt.arena.base) /* lazily init'ed only */
wo_vm_destroy(&wo_eng.shards[i]); wo_vm_destroy(&wo_eng.shards[i]);
free(wo_eng.shards[i].in_mu);
wo_eng.shards[i].in_mu = NULL;
}
/* the primary's inbox: same discard (main destroys its vm right after) */
{
wo_vm *pv = &wo_eng.shards[0];
if (pv->in_mu) {
wo_envelope *e = pv->in_head;
pv->in_head = pv->in_tail = NULL;
while (e) {
wo_envelope *nx = e->next;
if (e->kind == 1 && e->actor) {
free(e->actor->msgs);
free(e->actor);
}
free(e);
e = nx;
}
free(pv->in_mu);
pv->in_mu = NULL;
}
if (pv->wake_efd >= 0) {
close(pv->wake_efd);
pv->wake_efd = -1;
}
} }
free(ts); free(ts);
wo_eng.threads = NULL; wo_eng.threads = NULL;
@ -90,6 +273,7 @@ int wo_vm_init(wo_vm *vm, const wo_module *mod, size_t heap_cap) {
memset(vm, 0, sizeof(*vm)); memset(vm, 0, sizeof(*vm));
vm->mod = mod; vm->mod = mod;
vm->cur = &vm->f0; /* fiber 0: main — the one-fiber degenerate case */ vm->cur = &vm->f0; /* fiber 0: main — the one-fiber degenerate case */
vm->wake_efd = -1; /* engines/main wire a real one; tests run without */
vm->budget0 = 4000; /* reductions per slice, the BEAM-ish default */ vm->budget0 = 4000; /* reductions per slice, the BEAM-ish default */
{ {
const char *e = getenv("WO_REDUCTIONS"); const char *e = getenv("WO_REDUCTIONS");
@ -253,8 +437,26 @@ int wo_vm_actor_spawn(wo_vm *vm, uint64_t instance, uint32_t method_idx,
} }
a->instance = instance; a->instance = instance;
a->method = method_idx; a->method = method_idx;
a->next_all = vm->actors; /* placement (arc T6): round-robin across shards; same-shard when the
vm->actors = a; * engine is absent (tests) or single. The actor's list membership
* belongs to its HOME thread — an adopt envelope carries it there. */
uint32_t n = wo_eng.nshards ? wo_eng.nshards : 1;
uint32_t home = n > 1 ? (eng_rr++ % n) : vm->shard_id;
a->home = home;
if (home == vm->shard_id) {
a->next_all = vm->actors;
vm->actors = a;
} else {
wo_envelope *e = calloc(1, sizeof *e);
if (!e) {
free(a);
*msg = "out of memory";
return WO_T_OOM;
}
e->kind = 1;
e->actor = a;
inbox_push(&wo_eng.shards[home], e);
}
*out_addr = (uint64_t)(uintptr_t)a; *out_addr = (uint64_t)(uintptr_t)a;
return 0; return 0;
} }
@ -269,6 +471,21 @@ int wo_vm_actor_send(wo_vm *vm, uint64_t addr, uint64_t msg_val, const char **ms
*msg = "send: nil message"; *msg = "send: nil message";
return WO_T_BOUNDS; return WO_T_BOUNDS;
} }
if (a->home != vm->shard_id) {
/* cross-shard: the HOME thread owns the mailbox — send travels as
* an inbox envelope, ownership moves with it (the mutex is the
* happens-before edge TSan sees) */
wo_envelope *e = calloc(1, sizeof *e);
if (!e) {
*msg = "out of memory";
return WO_T_OOM;
}
e->kind = 0;
e->actor = a;
e->payload = msg_val;
inbox_push(&wo_eng.shards[a->home], e);
return 0;
}
if (actor_push(a, msg_val) != 0) { if (actor_push(a, msg_val) != 0) {
*msg = "out of memory"; *msg = "out of memory";
return WO_T_OOM; return WO_T_OOM;
@ -541,15 +758,25 @@ static int vm_run(wo_vm *vm, uint64_t *ret, wo_err *err) {
* cur/queued/parked). */ * cur/queued/parked). */
#define NEXT_RUNNABLE() \ #define NEXT_RUNNABLE() \
do { \ do { \
if (vm->in_mu) (void)wo_vm_adopt(vm); \
vm->cur = fib_dequeue(vm); \ vm->cur = fib_dequeue(vm); \
while (!vm->cur) { \ while (!vm->cur) { \
if (!vm->is_primary && !vm->parked) { \
vm->cur = &vm->f0; /* parked-safe sentinel */ \
return 2; /* worker drained: back to the serve loop */ \
} \
if (!vm->is_primary && eng_shutdown) { \
fib_reap_all(vm); \
vm->cur = &vm->f0; \
return 1; /* engine stopping: die clean */ \
} \
int iorc_ = wo_io_wait(vm); \ int iorc_ = wo_io_wait(vm); \
if (iorc_ == WO_IO_STOP) { \ if (iorc_ == WO_IO_STOP) { \
fib_reap_all(vm); \ fib_reap_all(vm); \
vm->cur = &vm->f0; \ vm->cur = &vm->f0; \
return 1; \ return 1; \
} \ } \
if (iorc_ != 0) { \ if (iorc_ < 0) { \
fib_reap_all(vm); \ fib_reap_all(vm); \
vm->cur = &vm->f0; \ vm->cur = &vm->f0; \
if (err) { \ if (err) { \
@ -558,6 +785,7 @@ static int vm_run(wo_vm *vm, uint64_t *ret, wo_err *err) {
} \ } \
return -1; \ return -1; \
} \ } \
if (vm->in_mu) (void)wo_vm_adopt(vm); \
vm->cur = fib_dequeue(vm); \ vm->cur = fib_dequeue(vm); \
} \ } \
vm->budget = vm->budget0; \ vm->budget = vm->budget0; \
@ -1057,6 +1285,19 @@ dispatch:
#undef DROP_CATCHES #undef DROP_CATCHES
} }
/* the worker flavor of wo_vm_call: no entry frame — run whatever the run
* queue holds (adopted fibers, actor deliveries) until drained (rc 2),
* stopped (1), or a fatal error (-1). */
int wo_vm_serve(wo_vm *vm) {
tls_vm = vm;
if (!vm->qhead) return 2;
vm->cur = fib_dequeue(vm);
vm->budget = vm->budget0;
uint64_t ret = 0;
wo_err err;
return vm_run(vm, &ret, &err);
}
int wo_vm_call(wo_vm *vm, uint32_t method_idx, const uint64_t *args, int wo_vm_call(wo_vm *vm, uint32_t method_idx, const uint64_t *args,
uint32_t argc, uint64_t *ret, wo_err *err) { uint32_t argc, uint64_t *ret, wo_err *err) {
if (err) memset(err, 0, sizeof(*err)); if (err) memset(err, 0, sizeof(*err));

View file

@ -87,6 +87,7 @@ typedef struct wo_fiber {
typedef struct wo_actor { typedef struct wo_actor {
uint64_t instance; /* the moved-in state object (runtime-owned) */ uint64_t instance; /* the moved-in state object (runtime-owned) */
uint32_t method; /* receive's method index (self + msg = 2 args) */ uint32_t method; /* receive's method index (self + msg = 2 args) */
uint32_t home; /* the shard whose thread owns mailbox + delivery */
uint64_t *msgs; /* FIFO ring, growable */ uint64_t *msgs; /* FIFO ring, growable */
uint32_t mhead, mlen, mcap; uint32_t mhead, mlen, mcap;
wo_fiber *active; /* the delivery fiber, NULL when idle */ wo_fiber *active; /* the delivery fiber, NULL when idle */
@ -100,7 +101,16 @@ typedef struct wo_vm {
* (runs the entry, owns the database); workers run wo_vm_serve. */ * (runs the entry, owns the database); workers run wo_vm_serve. */
uint32_t shard_id; uint32_t shard_id;
int is_primary; int is_primary;
int wake_efd; /* wakes an idle worker (inbox arrivals, shutdown) */ int wake_efd; /* wakes this shard's I/O wait (inbox arrivals, shutdown) */
/* the cross-shard inbox (arc T6): OTHER shards push envelopes here
* under in_mu and write wake_efd; only the OWNING thread pops. A
* mutex-guarded list, not the spec's lock-free ring — disclosed
* deviation, rings arrive when 9e measures the mutex. */
void *in_mu; /* pthread_mutex_t*, opaque here */
struct wo_envelope *in_head, *in_tail;
/* home-routed frees: objects owned by THIS shard's arena, dropped on
* another shard, come back here to die (header shard_id routes) */
struct wo_envelope *free_head;
wo_fiber f0; /* fiber 0: main — embedded; spawned fibers are calloc'd */ 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 *cur; /* the live fiber — every interpreter access goes here */
wo_fiber *qhead, *qtail; /* RUNNABLE fibers awaiting the interpreter */ wo_fiber *qhead, *qtail; /* RUNNABLE fibers awaiting the interpreter */
@ -112,6 +122,7 @@ typedef struct wo_vm {
wo_fiber *parked; /* fibers waiting on the plane */ wo_fiber *parked; /* fibers waiting on the plane */
uint32_t nparked; uint32_t nparked;
int io_kind; /* 0 = uring, 1 = epoll */ int io_kind; /* 0 = uring, 1 = epoll */
int efd_armed; /* wake_efd registered on the plane (uring oneshot) */
int io_fd; /* ring fd or epoll fd */ int io_fd; /* ring fd or epoll fd */
void *io_sq, *io_cq, *io_sqes; /* uring mmaps (NULL under epoll) */ void *io_sq, *io_cq, *io_sqes; /* uring mmaps (NULL under epoll) */
size_t io_sq_len, io_cq_len, io_sqes_len; size_t io_sq_len, io_cq_len, io_sqes_len;
@ -135,6 +146,19 @@ typedef struct wo_engine {
extern wo_engine wo_eng; /* the process's one engine (vm.c) */ extern wo_engine wo_eng; /* the process's one engine (vm.c) */
/* inbox envelope kinds (arc T6) */
typedef struct wo_envelope {
struct wo_envelope *next;
int kind; /* 0 = SEND (actor, payload), 1 = SPAWN-ADOPT (actor), 2 = FREE (payload = wo_hdr*) */
struct wo_actor *actor;
uint64_t payload;
} wo_envelope;
/* the shard whose thread we are on (thread-local; obj.c stamps and gc.c
* routes with it). NULL only before main's vm exists. */
wo_vm *wo_tls_vm(void);
void wo_tls_set(wo_vm *vm);
/* Start shards 1..n-1 (0 is the caller's, already init'ed in shards[0]). /* 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. */ * 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); int wo_engine_start(const wo_module *mod, size_t heap_cap, uint32_t nshards);

View file

@ -52,9 +52,43 @@ check_run() { # name [env pairs...]
fi fi
} }
check_run "auto" # EXACT checks pin one shard (deterministic by construction); the
check_run "uring" WO_IO=uring # multi-shard runs assert output SETS — the arc's honest narrowing
check_run "epoll" WO_IO=epoll check_run "auto" WO_SHARDS=1
check_run "uring" WO_SHARDS=1 WO_IO=uring
check_run "epoll" WO_SHARDS=1 WO_IO=epoll
# multi-shard (default = all cores): cross-shard placement + envelopes.
# Order across shards is scheduling; the SET of lines is the contract.
mout="$("$DIR/target/fibers" 2>&1)"
mc=$(printf '%s
' "$mout" | grep -c "count +")
mt=$(printf '%s
' "$mout" | grep -c "main tick")
if [[ "$mc" == "3" && "$mt" == "8" ]] && printf '%s' "$mout" | grep -q "count +3 = 6" && printf '%s' "$mout" | grep -q "sleeper: up" && printf '%s' "$mout" | grep -q "part2 done"; then
ok "multi-shard: cross-shard actor delivered the full set"
else
bad "multi-shard" "$(printf '%s' "$mout" | tr '
' '|')"
fi
# TSan: the cross-shard path race-checked (multi-shard, both parts)
make -C "$ROOT/runtime" wovm-tsan -s >/dev/null 2>&1
if "$WOC" build "$DIR" -o "$DIR/target/fibers_tsan" --runtime "$ROOT/runtime/build/wovm_tsan" >/dev/null 2>&1; then
tout="$("$DIR/target/fibers_tsan" 2>&1)"
if printf '%s' "$tout" | grep -q "unexpected memory mapping"; then
# kernel 6.5+ high-entropy ASLR vs TSan: the standard workaround
tout="$(setarch "$(uname -m)" -R "$DIR/target/fibers_tsan" 2>&1)"
fi
if printf '%s' "$tout" | grep -q "part2 done" && ! printf '%s' "$tout" | grep -qi "ThreadSanitizer"; then
ok "TSan multi-shard run clean"
else
bad "tsan" "$(printf '%s' "$tout" | grep -i -m2 "SUMMARY\|WARNING" | tr '
' '|')"
fi
else
bad "tsan" "build failed"
fi
# ASan flavor: rebuild the binary against the ASan runtime and repeat once # ASan flavor: rebuild the binary against the ASan runtime and repeat once
make -C "$ROOT/runtime" wovm-asan -s >/dev/null 2>&1 make -C "$ROOT/runtime" wovm-asan -s >/dev/null 2>&1

View file

@ -24,7 +24,13 @@
# rejection, a crash, a hang — fails and names the fixture. One line # rejection, a crash, a hang — fails and names the fixture. One line
# per fixture, a final tally, nonzero exit if anything failed. # per fixture, a final tally, nonzero exit if anything failed.
set -uo pipefail # no -e: a failing fixture is handled explicitly, one at a time set -uo pipefail
# The corpus asserts EXACT outputs — deterministic single-shard semantics.
# Multi-shard scheduling is nondeterministic by nature (the arc's spec
# narrows determinism to output SETS there); the multi-shard/TSan proofs
# live in the fibers gate, not here.
export WO_SHARDS=1 # no -e: a failing fixture is handled explicitly, one at a time
shopt -s nullglob shopt -s nullglob
ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"

View file

@ -0,0 +1 @@
WO-E222

View file

@ -0,0 +1,24 @@
-- WO-E222: an actor's state or message must not be (or contain) a traced
-- class — aliased graphs cannot cross shard heaps. Node is traced by
-- inference (self-referential); Box contains it.
class Node {
next: ?Node
v: Int
}
class Box {
head: ?Node
}
class Keeper {
pad: Int
fn receive(msg: Box) {
if msg.head == nil { print("empty"); } else { print("full"); }
}
}
fn main() -> Int {
let a = spawn Keeper { pad: 0 };
return 0;
}