From 897c8442f0c370558e08f6dadd1c2fce3a2c3aba Mon Sep 17 00:00:00 2001 From: "shoney.arickathil" Date: Thu, 20 Aug 2026 09:58:04 +0200 Subject: [PATCH] =?UTF-8?q?feat(runtime):=20the=20per-shard=20I/O=20plane?= =?UTF-8?q?=20=E2=80=94=20io=5Furing-first=20fiber=20parking=20(arc=20T4)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - park.c/park.h: one event loop per shard. io_uring PRIMARY (raw io_uring_setup/io_uring_enter, uapi structs mirrored, 5.4-floor ops: POLL_ADD for fd readiness, TIMEOUT for sleeps, user_data = the fiber); epoll+deadline-scan FALLBACK behind the startup probe; WO_IO= uring|epoll forces either so CI proves both on one kernel - park protocol: a blocking builtin fills cur->park_* and returns WO_SYS_PARKED; resume either RE-EXECUTES it (fd readiness: accept/ read/write retry) or continues PAST it (sleep: result preset, park_done=1 — re-executing would restart the full duration) - sysio: listener + accepted fds nonblocking (accept4 SOCK_NONBLOCK); accept/read park on EAGAIN; write parks on EAGAIN with its partial progress carried across the retry in park_wr_at; sleep parks on a deadline — with ONE fiber the plane's wait IS the blocking call, program mode is the degenerate case, not a special one - scheduler: NEXT_RUNNABLE waits on the plane when the queue empties; a stop interrupting the wait reaps EVERY fiber (queued and parked) and returns the clean-stop status; parked fibers are GC roots and fib_reap_all drains them - proof: full battery green on the uring path (oop-e2e 92/0, log-watcher 7/0 incl. the mcp accept/read/write loop, employee 8/0, web-app 21/0, deps 8/0), WO_IO=epoll battery green (log-watcher 7/0, web-app 21/0), WO_IO=uring forced green, LW_SOAK=8 10/0 (fd + RSS flatness holds over parked I/O) Co-Authored-By: Claude Opus 5 (1M context) --- runtime/src/builtin.h | 6 + runtime/src/park.c | 296 ++++++++++++++++++++++++++++++++++++++++++ runtime/src/park.h | 34 +++++ runtime/src/sysio.c | 57 ++++++-- runtime/src/vm.c | 68 +++++++++- runtime/src/vm.h | 22 ++++ 6 files changed, 465 insertions(+), 18 deletions(-) create mode 100644 runtime/src/park.c create mode 100644 runtime/src/park.h diff --git a/runtime/src/builtin.h b/runtime/src/builtin.h index 88c40ab..6c2bf70 100644 --- a/runtime/src/builtin.h +++ b/runtime/src/builtin.h @@ -14,6 +14,12 @@ * `kill -9` ends it. Negative so it cannot collide with a WO_T_* code, and * deliberately NOT catchable: `try` must not be able to swallow a stop. */ #define WO_SYS_STOPPED (-2) +/* arc T4: the builtin parked the calling fiber against the shard's I/O + * plane (fb->park_* filled); the interpreter schedules another fiber. */ +#define WO_SYS_PARKED (-3) + +/* the stop flag, readable by the I/O plane's wait loop (park.c) */ +int wo_sys_stop_pending(void); int wo_builtin(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg); diff --git a/runtime/src/park.c b/runtime/src/park.c new file mode 100644 index 0000000..ca6a7fe --- /dev/null +++ b/runtime/src/park.c @@ -0,0 +1,296 @@ +/* park.c — the per-shard I/O plane (arc stage 1 Task 4). See park.h. + * + * The uring structs below mirror include/uapi/linux/io_uring.h exactly + * (the fields this file touches; trailing space is padded to the kernel's + * sizes). The layout is part of the kernel ABI — stable by contract. */ +#define _GNU_SOURCE /* syscall(), struct layouts */ +#include "park.h" + +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include "builtin.h" + +/* ---- raw io_uring ABI (uapi mirror, the subset used) ------------------ */ + +struct io_sqring_offsets { + uint32_t head, tail, ring_mask, ring_entries, flags, dropped, array; + uint32_t resv1; + uint64_t user_addr; +}; +struct io_cqring_offsets { + uint32_t head, tail, ring_mask, ring_entries, overflow, cqes, flags; + uint32_t resv1; + uint64_t user_addr; +}; +struct io_uring_params { + uint32_t sq_entries, cq_entries, flags, sq_thread_cpu, sq_thread_idle; + uint32_t features, wq_fd, resv[3]; + struct io_sqring_offsets sq_off; + struct io_cqring_offsets cq_off; +}; +struct io_uring_sqe { + uint8_t opcode, flags; + uint16_t ioprio; + int32_t fd; + uint64_t off, addr; + uint32_t len; + union { + uint32_t poll32_events; /* POLL_ADD (little-endian u32 of poll mask) */ + uint32_t timeout_flags; /* TIMEOUT */ + uint32_t rw_flags; + }; + uint64_t user_data; + uint64_t pad2[3]; +}; +struct io_uring_cqe { + uint64_t user_data; + int32_t res; + uint32_t flags; +}; + +#define IORING_OP_POLL_ADD 6 +#define IORING_OP_TIMEOUT 11 +#define IORING_ENTER_GETEVENTS 1u +#define IORING_OFF_SQ_RING 0ULL +#define IORING_OFF_CQ_RING 0x8000000ULL +#define IORING_OFF_SQES 0x10000000ULL + +/* the sq/cq ring pointers, resolved once from the params offsets */ +typedef struct { + uint32_t *sq_head, *sq_tail, *sq_mask, *sq_array; + uint32_t *cq_head, *cq_tail, *cq_mask; + struct io_uring_cqe *cqes; + struct io_uring_sqe *sqes; +} rings; + +static struct io_uring_params g_params; /* offsets survive init */ + +static rings ring_ptrs(const wo_vm *vm) { + rings r; + uint8_t *sq = (uint8_t *)vm->io_sq, *cq = (uint8_t *)vm->io_cq; + r.sq_head = (uint32_t *)(sq + g_params.sq_off.head); + r.sq_tail = (uint32_t *)(sq + g_params.sq_off.tail); + r.sq_mask = (uint32_t *)(sq + g_params.sq_off.ring_mask); + r.sq_array = (uint32_t *)(sq + g_params.sq_off.array); + r.cq_head = (uint32_t *)(cq + g_params.cq_off.head); + r.cq_tail = (uint32_t *)(cq + g_params.cq_off.tail); + r.cq_mask = (uint32_t *)(cq + g_params.cq_off.ring_mask); + r.cqes = (struct io_uring_cqe *)(cq + g_params.cq_off.cqes); + r.sqes = (struct io_uring_sqe *)vm->io_sqes; + return r; +} + +static int uring_init(wo_vm *vm) { + memset(&g_params, 0, sizeof g_params); + long fd = syscall(SYS_io_uring_setup, 64u, &g_params); + if (fd < 0) return -1; + size_t sq_len = g_params.sq_off.array + g_params.sq_entries * sizeof(uint32_t); + size_t cq_len = g_params.cq_off.cqes + g_params.cq_entries * sizeof(struct io_uring_cqe); + size_t sqes_len = g_params.sq_entries * sizeof(struct io_uring_sqe); + void *sq = mmap(NULL, sq_len, PROT_READ | PROT_WRITE, MAP_SHARED | MAP_POPULATE, (int)fd, + IORING_OFF_SQ_RING); + void *cq = mmap(NULL, cq_len, PROT_READ | PROT_WRITE, MAP_SHARED | MAP_POPULATE, (int)fd, + IORING_OFF_CQ_RING); + void *sqes = mmap(NULL, sqes_len, PROT_READ | PROT_WRITE, MAP_SHARED | MAP_POPULATE, (int)fd, + IORING_OFF_SQES); + if (sq == MAP_FAILED || cq == MAP_FAILED || sqes == MAP_FAILED) { + if (sq != MAP_FAILED) munmap(sq, sq_len); + if (cq != MAP_FAILED) munmap(cq, cq_len); + if (sqes != MAP_FAILED) munmap(sqes, sqes_len); + close((int)fd); + return -1; + } + vm->io_kind = 0; + vm->io_fd = (int)fd; + vm->io_sq = sq; + vm->io_cq = cq; + vm->io_sqes = sqes; + vm->io_sq_len = sq_len; + vm->io_cq_len = cq_len; + vm->io_sqes_len = sqes_len; + return 0; +} + +static int uring_submit(wo_vm *vm, const struct io_uring_sqe *sqe) { + rings r = ring_ptrs(vm); + uint32_t tail = *r.sq_tail; + uint32_t idx = tail & *r.sq_mask; + r.sqes[idx] = *sqe; + r.sq_array[idx] = idx; + __atomic_store_n(r.sq_tail, tail + 1, __ATOMIC_RELEASE); + long rc = syscall(SYS_io_uring_enter, vm->io_fd, 1u, 0u, 0u, NULL, 0); + return rc < 0 ? -1 : 0; +} + +/* ---- backend-neutral helpers ------------------------------------------ */ + +static int64_t now_ms(void) { + struct timespec ts; + clock_gettime(CLOCK_REALTIME, &ts); + return (int64_t)ts.tv_sec * 1000 + ts.tv_nsec / 1000000; +} + +static void parked_unlink(wo_vm *vm, wo_fiber *fb) { + wo_fiber **pp = &vm->parked; + while (*pp && *pp != fb) pp = &(*pp)->pnext; + if (*pp) { + *pp = fb->pnext; + fb->pnext = NULL; + vm->nparked--; + } +} + +/* wake: parked -> run queue (fib_enqueue lives in vm.c; the tiny mirror + * here keeps park.c free of vm.c internals) */ +static void wake(wo_vm *vm, wo_fiber *fb) { + parked_unlink(vm, fb); + fb->state = WO_FIB_RUNNABLE; + fb->next = NULL; + if (vm->qtail) vm->qtail->next = fb; + else vm->qhead = fb; + vm->qtail = fb; +} + +/* ---- API --------------------------------------------------------------- */ + +int wo_io_init(wo_vm *vm) { + const char *force = getenv("WO_IO"); + vm->io_fd = -1; + if (!force || strcmp(force, "epoll") != 0) { + if (uring_init(vm) == 0) return 0; + if (force && strcmp(force, "uring") == 0) return -1; /* forced, absent */ + } + int ep = epoll_create1(0); + if (ep < 0) return -1; + vm->io_kind = 1; + vm->io_fd = ep; + return 0; +} + +void wo_io_destroy(wo_vm *vm) { + if (vm->io_fd < 0) return; + if (vm->io_kind == 0) { + munmap(vm->io_sq, vm->io_sq_len); + munmap(vm->io_cq, vm->io_cq_len); + munmap(vm->io_sqes, vm->io_sqes_len); + } + close(vm->io_fd); + vm->io_fd = -1; +} + +int wo_io_arm(wo_vm *vm, wo_fiber *fb) { + fb->state = WO_FIB_PARKED; + fb->pnext = vm->parked; + vm->parked = fb; + vm->nparked++; + if (vm->io_kind == 0) { + struct io_uring_sqe sqe; + memset(&sqe, 0, sizeof sqe); + sqe.user_data = (uint64_t)(uintptr_t)fb; + if (fb->park_fd >= 0) { + sqe.opcode = IORING_OP_POLL_ADD; + sqe.fd = fb->park_fd; + sqe.poll32_events = (uint32_t)(uint16_t)fb->park_events; + } else { + int64_t rel = fb->park_deadline - now_ms(); + if (rel < 0) rel = 0; + fb->park_ts.sec = rel / 1000; + fb->park_ts.nsec = (rel % 1000) * 1000000LL; + sqe.opcode = IORING_OP_TIMEOUT; + sqe.fd = -1; + sqe.addr = (uint64_t)(uintptr_t)&fb->park_ts; + sqe.len = 1; + } + if (uring_submit(vm, &sqe) != 0) { + parked_unlink(vm, fb); + return -1; + } + return 0; + } + /* epoll fallback: fds registered oneshot; deadlines live on the parked + * list and become the wait timeout */ + if (fb->park_fd >= 0) { + struct epoll_event ev; + memset(&ev, 0, sizeof ev); + ev.events = (uint32_t)(uint16_t)fb->park_events | EPOLLONESHOT; + ev.data.ptr = fb; + if (epoll_ctl(vm->io_fd, EPOLL_CTL_ADD, fb->park_fd, &ev) != 0 + && (errno != EEXIST || epoll_ctl(vm->io_fd, EPOLL_CTL_MOD, fb->park_fd, &ev) != 0)) { + parked_unlink(vm, fb); + return -1; + } + } + return 0; +} + +int wo_io_wait(wo_vm *vm) { + for (;;) { + if (wo_sys_stop_pending()) return WO_IO_STOP; + if (vm->io_kind == 0) { + rings r = ring_ptrs(vm); + uint32_t head = *r.cq_head; + uint32_t tail = __atomic_load_n(r.cq_tail, __ATOMIC_ACQUIRE); + if (head == tail) { + long rc = syscall(SYS_io_uring_enter, vm->io_fd, 0u, 1u, + IORING_ENTER_GETEVENTS, NULL, 0); + if (rc < 0 && errno == EINTR) continue; /* stop checked on loop */ + if (rc < 0) return -1; + tail = __atomic_load_n(r.cq_tail, __ATOMIC_ACQUIRE); + } + int woke = 0; + while (head != tail) { + struct io_uring_cqe *cqe = &r.cqes[head & *r.cq_mask]; + wo_fiber *fb = (wo_fiber *)(uintptr_t)cqe->user_data; + if (fb && fb->state == WO_FIB_PARKED) { + wake(vm, fb); + woke = 1; + } + head++; + } + __atomic_store_n(r.cq_head, head, __ATOMIC_RELEASE); + if (woke) return 0; + continue; + } + /* epoll: timeout from the nearest sleep deadline */ + int timeout = -1; + int64_t now = now_ms(); + for (wo_fiber *fb = vm->parked; fb; fb = fb->pnext) + if (fb->park_fd < 0) { + int64_t rel = fb->park_deadline - now; + if (rel < 0) rel = 0; + if (timeout < 0 || rel < timeout) timeout = (int)rel; + } + struct epoll_event evs[16]; + int n = epoll_wait(vm->io_fd, evs, 16, timeout); + if (n < 0 && errno == EINTR) continue; + if (n < 0) return -1; + int woke = 0; + for (int i = 0; i < n; i++) { + wo_fiber *fb = (wo_fiber *)evs[i].data.ptr; + if (fb && fb->state == WO_FIB_PARKED) { + epoll_ctl(vm->io_fd, EPOLL_CTL_DEL, fb->park_fd, NULL); + wake(vm, fb); + woke = 1; + } + } + now = now_ms(); + wo_fiber *fb = vm->parked; + while (fb) { + wo_fiber *nx = fb->pnext; + if (fb->park_fd < 0 && fb->park_deadline <= now) { + wake(vm, fb); + woke = 1; + } + fb = nx; + } + if (woke) return 0; + } +} diff --git a/runtime/src/park.h b/runtime/src/park.h new file mode 100644 index 0000000..934efe7 --- /dev/null +++ b/runtime/src/park.h @@ -0,0 +1,34 @@ +/* park.h — the per-shard I/O plane (the 8+11 arc, stage 1 Task 4). + * + * One event loop per shard, io_uring-FIRST (developer directive; the linux + * reference project's "single event loop" card): parked fibers wait as + * POLL_ADD (fd readiness — resume RE-EXECUTES the now-ready builtin) or + * TIMEOUT (sleep — resume continues PAST the builtin) submissions on a raw + * ring. epoll is the PORTABILITY fallback, selected by a startup probe + * (seccomp'd containers routinely deny io_uring) or forced with + * WO_IO=uring|epoll so CI proves both paths on one kernel. + * + * Raw io_uring_setup/io_uring_enter syscalls, libc-only — struct layouts + * mirrored from include/uapi/linux/io_uring.h. Ops restricted to the + * TIMEOUT floor (Linux 5.4); multishot variants are recorded future work. + */ +#ifndef WO_PARK_H +#define WO_PARK_H + +#include "vm.h" + +/* Probe (or WO_IO-force) the backend. 0 ok; nonzero = no backend (fatal). */ +int wo_io_init(wo_vm *vm); +void wo_io_destroy(wo_vm *vm); + +/* Register the just-parked fiber's wait (fb->park_* already filled by the + * builtin). 0 ok; -1 = arming failed (caller traps the builtin as IO). */ +int wo_io_arm(wo_vm *vm, wo_fiber *fb); + +/* Block until at least one parked fiber wakes; woken fibers move to the + * run queue. 0 = something woke; WO_IO_STOP = the stop flag interrupted + * the wait (caller unwinds everything); -1 = fatal backend error. */ +#define WO_IO_STOP (-2) +int wo_io_wait(wo_vm *vm); + +#endif /* WO_PARK_H */ diff --git a/runtime/src/sysio.c b/runtime/src/sysio.c index d120bf3..67dda0d 100644 --- a/runtime/src/sysio.c +++ b/runtime/src/sysio.c @@ -14,11 +14,12 @@ * documented beside each case below. Absence is the zero word, like * every other `?T`. */ -#define _POSIX_C_SOURCE 200809L +#define _GNU_SOURCE /* accept4, plus everything 200809L gave */ #include #include #include +#include #include #include #include @@ -89,6 +90,8 @@ static void on_stop(int sig) { * wait. A regular-file read is not one of them and keeps its plain retry. */ static int stop_pending(void) { return stop_flag != 0; } +int wo_sys_stop_pending(void) { return stop_flag != 0; } + static void install_stop_handlers(void) { if (stop_installed) return; stop_installed = 1; @@ -266,16 +269,20 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { } /* ---- time -------------------------------------------------------- */ case WO_B_TIME_SLEEP: { + /* arc T4: sleep parks against the I/O plane (deadline); with one + * fiber the plane's wait IS the blocking sleep — same path. The + * result is preset and park_done=1, so resume continues PAST the + * builtin (re-executing would restart the full duration). */ int64_t ms = (int64_t)R[B]; - if (ms > 0) { - struct timespec ts = {ms / 1000, (ms % 1000) * 1000000L}, rem; - while (nanosleep(&ts, &rem) != 0 && errno == EINTR) { - if (stop_pending()) return WO_SYS_STOPPED; - ts = rem; - } - } R[A] = 0; - return 0; + if (ms <= 0) return 0; + struct timespec now; + clock_gettime(CLOCK_REALTIME, &now); + vm->cur->park_fd = -1; + vm->cur->park_deadline = + (int64_t)now.tv_sec * 1000 + now.tv_nsec / 1000000 + ms; + vm->cur->park_done = 1; + return WO_SYS_PARKED; } case WO_B_TIME_LOCAL: { /* Parts: 0 year, 1 month (1..12), 2 day, 3 hour, * 4 minute, 5 second, 6 dow (0 = Sunday) */ @@ -358,6 +365,7 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { addr.sin_port = htons((uint16_t)port); addr.sin_addr.s_addr = !strcmp(path, "0.0.0.0") ? (in_addr_t)INADDR_ANY : inet_addr(path); + fcntl(fd, F_SETFL, fcntl(fd, F_GETFL, 0) | O_NONBLOCK); /* arc T4 */ if (bind(fd, (struct sockaddr *)&addr, sizeof addr) != 0 || listen(fd, 64) != 0) { *msg = strerror(errno); close(fd); @@ -369,10 +377,17 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { case WO_B_NET_ACCEPT: { int fd; for (;;) { - fd = accept((int)R[B], NULL, NULL); + fd = accept4((int)R[B], NULL, NULL, SOCK_NONBLOCK); if (fd >= 0 || errno != EINTR) break; if (stop_pending()) return WO_SYS_STOPPED; } + if (fd < 0 && (errno == EAGAIN || errno == EWOULDBLOCK)) { + /* arc T4: park until the listener is readable, then retry */ + vm->cur->park_fd = (int)R[B]; + vm->cur->park_events = POLLIN; + vm->cur->park_done = 0; + return WO_SYS_PARKED; + } if (fd < 0) { *msg = strerror(errno); return WO_T_IO; @@ -398,6 +413,15 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { return WO_SYS_STOPPED; } } + if (n < 0 && (errno == EAGAIN || errno == EWOULDBLOCK)) { + /* arc T4: nothing readable yet — free the buffer (the retry + * re-allocates) and park until the fd is readable */ + wo_str_free(rt, s); + vm->cur->park_fd = (int)R[B]; + vm->cur->park_events = POLLIN; + vm->cur->park_done = 0; + return WO_SYS_PARKED; + } if (n < 0) { wo_str_free(rt, s); *msg = strerror(errno); @@ -424,7 +448,11 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { *msg = "not a text value"; return WO_T_BOUNDS; } - uint32_t at = 0; + /* arc T4: a partial write's progress survives the park via + * park_wr_at — the retry re-executes this builtin with the same + * arguments and resumes at the saved offset */ + uint32_t at = vm->cur->park_wr_at; + vm->cur->park_wr_at = 0; while (at < body->len) { ssize_t n = write((int)R[B], body->data + at, body->len - at); if (n < 0) { @@ -432,6 +460,13 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { if (stop_pending()) return WO_SYS_STOPPED; continue; } + if (errno == EAGAIN || errno == EWOULDBLOCK) { + vm->cur->park_wr_at = at; + vm->cur->park_fd = (int)R[B]; + vm->cur->park_events = POLLOUT; + vm->cur->park_done = 0; + return WO_SYS_PARKED; + } *msg = strerror(errno); return WO_T_IO; } diff --git a/runtime/src/vm.c b/runtime/src/vm.c index 17e5cb1..1181a6d 100644 --- a/runtime/src/vm.c +++ b/runtime/src/vm.c @@ -9,6 +9,7 @@ #include "builtin.h" #include "cont.h" #include "gc.h" +#include "park.h" uint32_t wo_vm_depth(const wo_vm *vm) { return vm->cur->depth; } @@ -25,6 +26,7 @@ int wo_vm_init(wo_vm *vm, const wo_module *mod, size_t heap_cap) { } } vm->budget = vm->budget0; + if (wo_io_init(vm) != 0) return -1; /* no I/O plane at all: fatal */ return wo_rt_init(&vm->rt, heap_cap, mod->classes, mod->class_cnt); } @@ -44,6 +46,7 @@ void wo_vm_destroy(wo_vm *vm) { a = nx; } vm->actors = NULL; + wo_io_destroy(vm); wo_rt_destroy(&vm->rt); } @@ -107,10 +110,17 @@ static void fib_reap(wo_vm *vm, wo_fiber *fb) { } } -/* Main finished (return or stop): every remaining fiber unwinds clean. */ +/* Main finished (return or stop): every remaining fiber — queued AND + * parked — unwinds clean. */ static void fib_reap_all(wo_vm *vm) { wo_fiber *fb; while ((fb = fib_dequeue(vm)) != NULL) fib_reap(vm, fb); + while ((fb = vm->parked) != NULL) { + vm->parked = fb->pnext; + fb->pnext = NULL; + vm->nparked--; + fib_reap(vm, fb); + } } /* ---- actors (arc stage 1 Task 3) -------------------------------------- */ @@ -272,6 +282,8 @@ static void vm_gc_roots(wo_vm *vm) { vm_gc_roots_fiber(vm, vm->cur); for (const wo_fiber *fb = vm->qhead; fb; fb = fb->next) vm_gc_roots_fiber(vm, fb); + for (const wo_fiber *fb = vm->parked; fb; fb = fb->pnext) + vm_gc_roots_fiber(vm, fb); /* actors: moved-in state, queued messages, and the in-flight message * are runtime-owned — none sits in any frame's masks */ for (const wo_actor *a = vm->actors; a; a = a->next_all) { @@ -439,10 +451,9 @@ static int vm_run(wo_vm *vm, uint64_t *ret, wo_err *err) { "wovm: fiber trap %d at %s:%d: %s\n", \ err->code, err->method, err->line, err->msg); \ wo_fiber *dead = vm->cur; \ - vm->cur = fib_dequeue(vm); \ vm->nfibers--; \ free(dead); \ - vm->budget = vm->budget0; \ + NEXT_RUNNABLE(); \ RELOAD(); \ NEXT(); \ } \ @@ -450,6 +461,35 @@ static int vm_run(wo_vm *vm, uint64_t *ret, wo_err *err) { return -1; \ } while (0) +/* Pick the next runnable fiber; when the queue is empty, wait on the I/O + * plane for a parked one. A stop interrupting the wait unwinds EVERYTHING + * and returns 1 (the WO_SYS_STOPPED contract). The queue-and-parked-both- + * empty case cannot be reached from a live fiber (main is always one of + * cur/queued/parked). */ +#define NEXT_RUNNABLE() \ + do { \ + vm->cur = fib_dequeue(vm); \ + while (!vm->cur) { \ + int iorc_ = wo_io_wait(vm); \ + if (iorc_ == WO_IO_STOP) { \ + fib_reap_all(vm); \ + vm->cur = &vm->f0; \ + return 1; \ + } \ + if (iorc_ != 0) { \ + fib_reap_all(vm); \ + vm->cur = &vm->f0; \ + if (err) { \ + err->code = WO_T_IO; \ + snprintf(err->msg, sizeof err->msg, "I/O plane failed"); \ + } \ + return -1; \ + } \ + vm->cur = fib_dequeue(vm); \ + } \ + vm->budget = vm->budget0; \ + } while (0) + /* Collector safepoint (iteration 7b): placed at allocations, calls, and * loop back-edges — the pcs that already carry drop-table entries, so the * root snapshot's masks are exact. Costs one predictable branch when the @@ -625,6 +665,7 @@ dispatch: while (vm->cur->ncatch && vm->cur->catches[vm->cur->ncatch - 1].depth > vm->cur->depth) \ vm->cur->ncatch-- + /* A fiber's last frame returned. Main ending IS the program ending: every * other fiber unwinds through its drop maps (clean, ASan-proven) and the * program's value is main's. A spawned fiber ending just leaves the @@ -661,17 +702,15 @@ dispatch: memset(dead->regs + 2, 0, (size_t)(sme_->reg_cnt - 2) * 8u); \ dead->cur_msg = m_; \ fib_enqueue(vm, dead); \ - vm->cur = fib_dequeue(vm); \ - vm->budget = vm->budget0; \ + NEXT_RUNNABLE(); \ RELOAD(); \ NEXT(); \ } \ a->active = NULL; \ } \ - vm->cur = fib_dequeue(vm); \ vm->nfibers--; \ free(dead); \ - vm->budget = vm->budget0; \ + NEXT_RUNNABLE(); \ RELOAD(); \ NEXT(); \ } while (0) @@ -815,6 +854,21 @@ dispatch: * The stack is unwound exactly as an uncaught trap unwinds it, so * every live value is still released on the way out; the CLI turns * this into the same exit status a clean `return 0` gives. */ + if (brc == WO_SYS_PARKED) { + /* arc T4: the builtin filled cur->park_*. Resume either + * RE-EXECUTES it (park_done=0: fd readiness — accept/read/ + * write retry against a now-ready fd) or continues PAST it + * (park_done=1: sleep — result preset before parking). */ + vm->cur->frames[vm->cur->depth - 1].pc = vm->cur->park_done ? pc : pc - 1; + wo_fiber *pk = vm->cur; + if (wo_io_arm(vm, pk) != 0) { + pk->state = WO_FIB_RUNNABLE; + TRAPF(WO_T_IO, "%s", "cannot arm the I/O wait"); + } + NEXT_RUNNABLE(); + RELOAD(); + NEXT(); + } if (brc == WO_SYS_STOPPED) { vm->cur->frames[vm->cur->depth - 1].pc = pc - 1; vm->cur->ncatch = 0; diff --git a/runtime/src/vm.h b/runtime/src/vm.h index e40035b..e7c0a7b 100644 --- a/runtime/src/vm.h +++ b/runtime/src/vm.h @@ -59,6 +59,21 @@ typedef struct wo_fiber { wo_err caught; wo_fib_state state; struct wo_fiber *next; /* intrusive FIFO link (run queue) */ + /* parking (arc T4): what this fiber waits on while PARKED. park_done + * says how it resumes — 0 = re-execute the builtin (fd readiness: + * accept/read/write retry, now ready), 1 = continue PAST it (sleep: + * the result was preset before parking). park_wr_at carries a partial + * net.write's progress across the retry. park_ts must outlive the + * ring submission (TIMEOUT reads it asynchronously). */ + struct wo_fiber *pnext; /* parked-list link */ + int park_fd; /* -1 = deadline-only (sleep) */ + short park_events; /* POLLIN / POLLOUT */ + int park_done; + int64_t park_deadline; /* wall ms, sleep only */ + uint32_t park_wr_at; + struct { + long long sec, nsec; + } park_ts; /* arc actors: when this fiber is an actor's delivery fiber, `actor` * points at it and `cur_msg` is the message the current receive call * borrows — the RUNTIME owns it and drops it after the call returns. */ @@ -88,6 +103,13 @@ typedef struct wo_vm { int64_t budget0; /* reductions per slice (WO_REDUCTIONS, default 4000) */ int64_t budget; /* countdown for the live fiber */ wo_actor *actors; /* every spawned actor (torn down at destroy) */ + /* the I/O plane (arc T4, park.c): io_uring primary, epoll fallback */ + wo_fiber *parked; /* fibers waiting on the plane */ + uint32_t nparked; + int io_kind; /* 0 = uring, 1 = epoll */ + int io_fd; /* ring fd or epoll fd */ + void *io_sq, *io_cq, *io_sqes; /* uring mmaps (NULL under epoll) */ + size_t io_sq_len, io_cq_len, io_sqes_len; } wo_vm; /* arc: the spawn/send builtins' runtime halves (vm.c owns the scheduler). */