feat(runtime): the per-shard I/O plane — io_uring-first fiber parking (arc T4)
- 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) <noreply@anthropic.com>
This commit is contained in:
parent
ac7c80e1b9
commit
897c8442f0
6 changed files with 465 additions and 18 deletions
|
|
@ -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);
|
||||
|
||||
|
|
|
|||
296
runtime/src/park.c
Normal file
296
runtime/src/park.c
Normal file
|
|
@ -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 <errno.h>
|
||||
#include <poll.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <sys/epoll.h>
|
||||
#include <sys/mman.h>
|
||||
#include <sys/syscall.h>
|
||||
#include <time.h>
|
||||
#include <unistd.h>
|
||||
|
||||
#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;
|
||||
}
|
||||
}
|
||||
34
runtime/src/park.h
Normal file
34
runtime/src/park.h
Normal file
|
|
@ -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 */
|
||||
|
|
@ -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 <dirent.h>
|
||||
#include <errno.h>
|
||||
#include <fcntl.h>
|
||||
#include <poll.h>
|
||||
#include <arpa/inet.h>
|
||||
#include <netinet/in.h>
|
||||
#include <signal.h>
|
||||
|
|
@ -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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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). */
|
||||
|
|
|
|||
Loading…
Reference in a new issue