feat(lang42): proc.run parks — pidfd + epoll bundle + child registry
- deadlock proven first: chatty child (200 KB stdout, stderr held open) hung the old sequential drain 5.0 s into the alarm, code -1, stdout truncated at 8192; the leg demands completion under 4 s - rework: nonblocking pipe read ends + pidfd_open behind one epoll fd the fiber parks on (the _dl retry mould); both pipes drain on readiness, so the deadlock is gone structurally — leg passes in 15 ms - wo_child slot table in wo_vm (32/shard) carries cross-park state; caps refuse by name (kill + WO_T_IO), deadline armed via dl_active/dl_at, defaults 30 s / 1 MiB / 64 KiB - WO_B_PROC_RUN_DL = 96 shares the case (per-call deadline_ms/out_cap/ err_cap; compiler row lands in a later task) - fib_reap kills a reaped fiber's child; wo_vm_destroy sweeps the table - raw syscalls for pidfd_open/pidfd_send_signal: glibc 2.35 build floor has no wrappers - all 19 suites green under ASan+UBSan Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
parent
b8d9f9f594
commit
258222c3a1
5 changed files with 457 additions and 83 deletions
|
|
@ -20,6 +20,8 @@
|
||||||
#include <errno.h>
|
#include <errno.h>
|
||||||
#include <fcntl.h>
|
#include <fcntl.h>
|
||||||
#include <poll.h>
|
#include <poll.h>
|
||||||
|
#include <sys/epoll.h>
|
||||||
|
#include <sys/syscall.h>
|
||||||
#include <arpa/inet.h>
|
#include <arpa/inet.h>
|
||||||
#include <netinet/in.h>
|
#include <netinet/in.h>
|
||||||
#include <signal.h>
|
#include <signal.h>
|
||||||
|
|
@ -144,6 +146,107 @@ static wo_str *read_range(wo_rt *rt, int fd, off_t off, size_t want, const char
|
||||||
return exact;
|
return exact;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/* ---- iteration 42: bounded subprocess --------------------------------
|
||||||
|
* proc.run parks instead of blocking: the two pipe read ends and a pidfd
|
||||||
|
* for the child sit behind ONE epoll fd the fiber parks on (the plane
|
||||||
|
* watches one fd per fiber; the bundle turns three waits into it). The
|
||||||
|
* cross-park state is a wo_child slot in the shard's vm — the _dl retry
|
||||||
|
* protocol re-executes the builtin and the slot is how the re-entry
|
||||||
|
* remembers buffers, fds and caps. Every bound violation KILLS the child
|
||||||
|
* and traps WO_T_IO naming the bound; a zombie or an orphan is a bug by
|
||||||
|
* definition (fib_reap and wo_vm_destroy sweep the slots).
|
||||||
|
*
|
||||||
|
* glibc 2.35 (the release build floor) has no pidfd wrappers — raw
|
||||||
|
* syscalls, numbers guarded for older headers. */
|
||||||
|
#ifndef SYS_pidfd_open
|
||||||
|
#define SYS_pidfd_open 434
|
||||||
|
#endif
|
||||||
|
#ifndef SYS_pidfd_send_signal
|
||||||
|
#define SYS_pidfd_send_signal 424
|
||||||
|
#endif
|
||||||
|
|
||||||
|
#define WO_PROC_DL_DEFAULT 30000
|
||||||
|
#define WO_PROC_OUT_DEFAULT (1u << 20)
|
||||||
|
#define WO_PROC_ERR_DEFAULT (1u << 16)
|
||||||
|
|
||||||
|
/* bound-violation messages carry values; the buffer must outlive the
|
||||||
|
* return (wo_err copies later, on the trap path) — per-thread, one shard
|
||||||
|
* per thread */
|
||||||
|
static _Thread_local char proc_msg[96];
|
||||||
|
|
||||||
|
/* release everything a slot holds; the child must already be reaped */
|
||||||
|
static void proc_slot_close(wo_vm *vm, wo_child *ch) {
|
||||||
|
if (ch->pidfd >= 0) close(ch->pidfd);
|
||||||
|
if (ch->epfd >= 0) close(ch->epfd);
|
||||||
|
if (ch->ofd >= 0) close(ch->ofd);
|
||||||
|
if (ch->efd >= 0) close(ch->efd);
|
||||||
|
free(ch->obuf);
|
||||||
|
free(ch->ebuf);
|
||||||
|
if (ch->owner) ch->owner->proc_st = NULL;
|
||||||
|
memset(ch, 0, sizeof *ch);
|
||||||
|
vm->nchildren--;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* SIGKILL through the pidfd (no pid-reuse race), reap, release */
|
||||||
|
static void proc_slot_kill(wo_vm *vm, wo_child *ch) {
|
||||||
|
syscall(SYS_pidfd_send_signal, ch->pidfd, SIGKILL, NULL, 0);
|
||||||
|
int st;
|
||||||
|
while (waitpid(ch->pid, &st, 0) < 0 && errno == EINTR) {}
|
||||||
|
proc_slot_close(vm, ch);
|
||||||
|
}
|
||||||
|
|
||||||
|
void wo_proc_abandon(wo_vm *vm, wo_fiber *fb) {
|
||||||
|
if (fb->proc_st) proc_slot_kill(vm, fb->proc_st);
|
||||||
|
}
|
||||||
|
|
||||||
|
void wo_proc_reap_all(wo_vm *vm) {
|
||||||
|
for (uint32_t i = 0; i < WO_PROC_MAX; i++)
|
||||||
|
if (vm->children[i].used) proc_slot_kill(vm, &vm->children[i]);
|
||||||
|
}
|
||||||
|
|
||||||
|
/* append a chunk, growing by doubling up to the cap.
|
||||||
|
* 0 ok; -1 cap exceeded; -2 oom */
|
||||||
|
static int proc_buf_append(char **buf, size_t *len, size_t *alloc,
|
||||||
|
uint64_t cap, const char *chunk, size_t n) {
|
||||||
|
if (*len + n > (size_t)cap) return -1;
|
||||||
|
if (*len + n > *alloc) {
|
||||||
|
size_t want = *alloc ? *alloc : 4096;
|
||||||
|
while (want < *len + n) want *= 2;
|
||||||
|
if (want > (size_t)cap) want = (size_t)cap;
|
||||||
|
char *nb = realloc(*buf, want);
|
||||||
|
if (!nb) return -2;
|
||||||
|
*buf = nb;
|
||||||
|
*alloc = want;
|
||||||
|
}
|
||||||
|
memcpy(*buf + *len, chunk, n);
|
||||||
|
*len += n;
|
||||||
|
return 0;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* drain one pipe until EAGAIN or EOF. 0 ok (fd may now be -1),
|
||||||
|
* -1 cap exceeded, -2 oom, -3 read error (errno kept) */
|
||||||
|
static int proc_drain_fd(int *fd, char **buf, size_t *len, size_t *alloc,
|
||||||
|
uint64_t cap) {
|
||||||
|
char chunk[4096];
|
||||||
|
while (*fd >= 0) {
|
||||||
|
ssize_t n = read(*fd, chunk, sizeof chunk);
|
||||||
|
if (n > 0) {
|
||||||
|
int rc = proc_buf_append(buf, len, alloc, cap, chunk, (size_t)n);
|
||||||
|
if (rc != 0) return rc;
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
if (n == 0) { /* EOF: the child closed its end (or died) */
|
||||||
|
close(*fd);
|
||||||
|
*fd = -1;
|
||||||
|
return 0;
|
||||||
|
}
|
||||||
|
if (errno == EINTR) continue;
|
||||||
|
if (errno == EAGAIN || errno == EWOULDBLOCK) return 0;
|
||||||
|
return -3;
|
||||||
|
}
|
||||||
|
return 0;
|
||||||
|
}
|
||||||
|
|
||||||
int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) {
|
int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) {
|
||||||
wo_rt *rt = &vm->rt;
|
wo_rt *rt = &vm->rt;
|
||||||
uint8_t A = wo_ins_a(ins), B = wo_ins_b(ins), C = wo_ins_c(ins);
|
uint8_t A = wo_ins_a(ins), B = wo_ins_b(ins), C = wo_ins_c(ins);
|
||||||
|
|
@ -719,97 +822,260 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) {
|
||||||
R[A] = (uint64_t)(uintptr_t)out;
|
R[A] = (uint64_t)(uintptr_t)out;
|
||||||
return 0;
|
return 0;
|
||||||
}
|
}
|
||||||
/* ---- proc -------------------------------------------------------- */
|
/* ---- proc (iteration 42: bounded + parked) ------------------------
|
||||||
case WO_B_PROC_RUN: { /* Proc: 0 code, 1 out, 2 err. argv[0] is the
|
* Proc: 0 code, 1 out, 2 err. argv[0] is the command itself; the
|
||||||
* command itself; the `multi Text` argument
|
* `multi Text` argument supplies the rest. First entry validates,
|
||||||
* supplies the rest. stdout and stderr are
|
* forks and claims a wo_child slot; every entry drains whatever is
|
||||||
* captured through one pipe each, capped. */
|
* ready and either finishes (child reaped), refuses (a bound hit,
|
||||||
if (cstr_of(R[B], path, sizeof path, msg)) return WO_T_BOUNDS;
|
* child killed), or parks on the slot's epoll bundle with the
|
||||||
wo_multi *argv_m = (wo_multi *)(uintptr_t)R[B + 1];
|
* deadline armed. RUN uses the named defaults; RUN_DL states them
|
||||||
if (!argv_m || argv_m->h.class_id != WO_CLS_MULTI || argv_m->elem_kind != WO_K_TEXT) {
|
* per call (<= 0 picks the default). */
|
||||||
*msg = "`proc.run` needs a `multi Text` of arguments";
|
case WO_B_PROC_RUN:
|
||||||
return WO_T_BOUNDS;
|
case WO_B_PROC_RUN_DL: {
|
||||||
}
|
wo_fiber *fb = vm->cur;
|
||||||
if (argv_m->len > 62) {
|
int isdl = (C == WO_B_PROC_RUN_DL);
|
||||||
*msg = "too many process arguments";
|
uint64_t cls_id = isdl ? R[B + 5] : R[B + 2];
|
||||||
return WO_T_BOUNDS;
|
struct timespec dts;
|
||||||
}
|
clock_gettime(CLOCK_REALTIME, &dts);
|
||||||
char *argv[64];
|
int64_t dnow = (int64_t)dts.tv_sec * 1000 + dts.tv_nsec / 1000000;
|
||||||
char argbuf[62][512];
|
|
||||||
argv[0] = path;
|
if (!fb->proc_st) { /* ---- first entry: validate, fork, claim */
|
||||||
for (uint32_t i = 0; i < argv_m->len; i++) {
|
if (cstr_of(R[B], path, sizeof path, msg)) return WO_T_BOUNDS;
|
||||||
const wo_str *a = (const wo_str *)(uintptr_t)argv_m->items[i];
|
wo_multi *argv_m = (wo_multi *)(uintptr_t)R[B + 1];
|
||||||
if (!a || a->h.class_id != WO_CLS_STR || a->len + 1 > sizeof argbuf[0]) {
|
if (!argv_m || argv_m->h.class_id != WO_CLS_MULTI ||
|
||||||
*msg = "process argument is not a short text";
|
argv_m->elem_kind != WO_K_TEXT) {
|
||||||
|
*msg = "`proc.run` needs a `multi Text` of arguments";
|
||||||
return WO_T_BOUNDS;
|
return WO_T_BOUNDS;
|
||||||
}
|
}
|
||||||
memcpy(argbuf[i], a->data, a->len);
|
if (argv_m->len > 62) {
|
||||||
argbuf[i][a->len] = '\0';
|
*msg = "too many process arguments";
|
||||||
argv[i + 1] = argbuf[i];
|
return WO_T_BOUNDS;
|
||||||
}
|
}
|
||||||
argv[argv_m->len + 1] = NULL;
|
char *argv[64];
|
||||||
int op[2], ep[2];
|
char argbuf[62][512];
|
||||||
if (pipe(op) != 0) {
|
argv[0] = path;
|
||||||
*msg = strerror(errno);
|
for (uint32_t i = 0; i < argv_m->len; i++) {
|
||||||
return WO_T_IO;
|
const wo_str *a = (const wo_str *)(uintptr_t)argv_m->items[i];
|
||||||
}
|
if (!a || a->h.class_id != WO_CLS_STR ||
|
||||||
if (pipe(ep) != 0) {
|
a->len + 1 > sizeof argbuf[0]) {
|
||||||
close(op[0]);
|
*msg = "process argument is not a short text";
|
||||||
|
return WO_T_BOUNDS;
|
||||||
|
}
|
||||||
|
memcpy(argbuf[i], a->data, a->len);
|
||||||
|
argbuf[i][a->len] = '\0';
|
||||||
|
argv[i + 1] = argbuf[i];
|
||||||
|
}
|
||||||
|
argv[argv_m->len + 1] = NULL;
|
||||||
|
|
||||||
|
int64_t dl_ms = WO_PROC_DL_DEFAULT;
|
||||||
|
uint64_t out_cap = WO_PROC_OUT_DEFAULT, err_cap = WO_PROC_ERR_DEFAULT;
|
||||||
|
if (isdl) {
|
||||||
|
if ((int64_t)R[B + 2] > 0) dl_ms = (int64_t)R[B + 2];
|
||||||
|
if ((int64_t)R[B + 3] > 0) out_cap = R[B + 3];
|
||||||
|
if ((int64_t)R[B + 4] > 0) err_cap = R[B + 4];
|
||||||
|
}
|
||||||
|
|
||||||
|
wo_child *ch = NULL;
|
||||||
|
for (uint32_t i = 0; i < WO_PROC_MAX; i++)
|
||||||
|
if (!vm->children[i].used) {
|
||||||
|
ch = &vm->children[i];
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
if (!ch) { /* the ceiling fails CLOSED, by name */
|
||||||
|
snprintf(proc_msg, sizeof proc_msg,
|
||||||
|
"process ceiling: %u live children on this shard",
|
||||||
|
WO_PROC_MAX);
|
||||||
|
*msg = proc_msg;
|
||||||
|
return WO_T_IO;
|
||||||
|
}
|
||||||
|
|
||||||
|
int op[2], ep[2];
|
||||||
|
if (pipe(op) != 0) {
|
||||||
|
*msg = strerror(errno);
|
||||||
|
return WO_T_IO;
|
||||||
|
}
|
||||||
|
if (pipe(ep) != 0) {
|
||||||
|
close(op[0]);
|
||||||
|
close(op[1]);
|
||||||
|
*msg = strerror(errno);
|
||||||
|
return WO_T_IO;
|
||||||
|
}
|
||||||
|
pid_t pid = fork();
|
||||||
|
if (pid < 0) {
|
||||||
|
close(op[0]);
|
||||||
|
close(op[1]);
|
||||||
|
close(ep[0]);
|
||||||
|
close(ep[1]);
|
||||||
|
*msg = strerror(errno);
|
||||||
|
return WO_T_IO;
|
||||||
|
}
|
||||||
|
if (pid == 0) {
|
||||||
|
dup2(op[1], STDOUT_FILENO);
|
||||||
|
dup2(ep[1], STDERR_FILENO);
|
||||||
|
close(op[0]);
|
||||||
|
close(op[1]);
|
||||||
|
close(ep[0]);
|
||||||
|
close(ep[1]);
|
||||||
|
execvp(path, argv);
|
||||||
|
_exit(127); /* exec failed: the same code a shell reports */
|
||||||
|
}
|
||||||
close(op[1]);
|
close(op[1]);
|
||||||
*msg = strerror(errno);
|
|
||||||
return WO_T_IO;
|
|
||||||
}
|
|
||||||
pid_t pid = fork();
|
|
||||||
if (pid < 0) {
|
|
||||||
close(op[0]);
|
|
||||||
close(op[1]);
|
|
||||||
close(ep[0]);
|
|
||||||
close(ep[1]);
|
close(ep[1]);
|
||||||
*msg = strerror(errno);
|
/* NONBLOCK on the parent's read ends only — the child keeps
|
||||||
|
* ordinary blocking pipes */
|
||||||
|
fcntl(op[0], F_SETFL, fcntl(op[0], F_GETFL, 0) | O_NONBLOCK);
|
||||||
|
fcntl(ep[0], F_SETFL, fcntl(ep[0], F_GETFL, 0) | O_NONBLOCK);
|
||||||
|
int pidfd = (int)syscall(SYS_pidfd_open, pid, 0);
|
||||||
|
int epfd = pidfd >= 0 ? epoll_create1(0) : -1;
|
||||||
|
if (epfd >= 0) {
|
||||||
|
struct epoll_event ev;
|
||||||
|
memset(&ev, 0, sizeof ev);
|
||||||
|
ev.events = EPOLLIN;
|
||||||
|
ev.data.fd = op[0];
|
||||||
|
epoll_ctl(epfd, EPOLL_CTL_ADD, op[0], &ev);
|
||||||
|
ev.data.fd = ep[0];
|
||||||
|
epoll_ctl(epfd, EPOLL_CTL_ADD, ep[0], &ev);
|
||||||
|
ev.data.fd = pidfd;
|
||||||
|
epoll_ctl(epfd, EPOLL_CTL_ADD, pidfd, &ev);
|
||||||
|
}
|
||||||
|
if (epfd < 0) { /* pidfd_open or epoll failed: no orphan */
|
||||||
|
if (pidfd >= 0) close(pidfd);
|
||||||
|
kill(pid, SIGKILL);
|
||||||
|
int st;
|
||||||
|
while (waitpid(pid, &st, 0) < 0 && errno == EINTR) {}
|
||||||
|
close(op[0]);
|
||||||
|
close(ep[0]);
|
||||||
|
*msg = strerror(errno);
|
||||||
|
return WO_T_IO;
|
||||||
|
}
|
||||||
|
ch->used = 1;
|
||||||
|
ch->pid = (int)pid;
|
||||||
|
ch->pidfd = pidfd;
|
||||||
|
ch->epfd = epfd;
|
||||||
|
ch->ofd = op[0];
|
||||||
|
ch->efd = ep[0];
|
||||||
|
ch->obuf = ch->ebuf = NULL;
|
||||||
|
ch->olen = ch->elen = ch->oalloc = ch->ealloc = 0;
|
||||||
|
ch->out_cap = out_cap;
|
||||||
|
ch->err_cap = err_cap;
|
||||||
|
ch->owner = fb;
|
||||||
|
fb->proc_st = ch;
|
||||||
|
vm->nchildren++;
|
||||||
|
/* the _dl protocol: arm once, the retry remembers */
|
||||||
|
fb->dl_active = 1;
|
||||||
|
fb->dl_at = dnow + dl_ms;
|
||||||
|
}
|
||||||
|
|
||||||
|
wo_child *ch = fb->proc_st;
|
||||||
|
/* ---- drain whatever is ready, caps enforced */
|
||||||
|
int drc = proc_drain_fd(&ch->ofd, &ch->obuf, &ch->olen, &ch->oalloc,
|
||||||
|
ch->out_cap);
|
||||||
|
uint64_t hit_cap = ch->out_cap;
|
||||||
|
const char *hit_name = "stdout";
|
||||||
|
if (drc == 0) {
|
||||||
|
drc = proc_drain_fd(&ch->efd, &ch->ebuf, &ch->elen, &ch->ealloc,
|
||||||
|
ch->err_cap);
|
||||||
|
hit_cap = ch->err_cap;
|
||||||
|
hit_name = "stderr";
|
||||||
|
}
|
||||||
|
if (drc != 0) {
|
||||||
|
int rderr = errno;
|
||||||
|
fb->dl_active = 0;
|
||||||
|
if (drc == -1) {
|
||||||
|
snprintf(proc_msg, sizeof proc_msg,
|
||||||
|
"process %s cap %llu bytes exceeded", hit_name,
|
||||||
|
(unsigned long long)hit_cap);
|
||||||
|
*msg = proc_msg;
|
||||||
|
proc_slot_kill(vm, ch);
|
||||||
|
return WO_T_IO;
|
||||||
|
}
|
||||||
|
proc_slot_kill(vm, ch);
|
||||||
|
if (drc == -2) {
|
||||||
|
*msg = "out of memory";
|
||||||
|
return WO_T_OOM;
|
||||||
|
}
|
||||||
|
*msg = strerror(rderr);
|
||||||
return WO_T_IO;
|
return WO_T_IO;
|
||||||
}
|
}
|
||||||
if (pid == 0) {
|
|
||||||
dup2(op[1], STDOUT_FILENO);
|
|
||||||
dup2(ep[1], STDERR_FILENO);
|
|
||||||
close(op[0]);
|
|
||||||
close(op[1]);
|
|
||||||
close(ep[0]);
|
|
||||||
close(ep[1]);
|
|
||||||
execvp(path, argv);
|
|
||||||
_exit(127); /* exec failed: the same code a shell reports */
|
|
||||||
}
|
|
||||||
close(op[1]);
|
|
||||||
close(ep[1]);
|
|
||||||
char obuf[8192], ebuf[4096];
|
|
||||||
size_t olen = 0, elen = 0;
|
|
||||||
ssize_t n;
|
|
||||||
while (olen < sizeof obuf && (n = read(op[0], obuf + olen, sizeof obuf - olen)) > 0)
|
|
||||||
olen += (size_t)n;
|
|
||||||
while (elen < sizeof ebuf && (n = read(ep[0], ebuf + elen, sizeof ebuf - elen)) > 0)
|
|
||||||
elen += (size_t)n;
|
|
||||||
close(op[0]);
|
|
||||||
close(ep[0]);
|
|
||||||
int status = 0;
|
int status = 0;
|
||||||
while (waitpid(pid, &status, 0) < 0 && errno == EINTR) {
|
pid_t r = waitpid(ch->pid, &status, WNOHANG);
|
||||||
if (stop_pending()) return WO_SYS_STOPPED;
|
if (r == (pid_t)ch->pid) { /* ---- exited: final drain, answer */
|
||||||
|
/* the write ends died with the child; what remains in the
|
||||||
|
* pipes reads out then EOFs — still cap-bounded */
|
||||||
|
int frc = proc_drain_fd(&ch->ofd, &ch->obuf, &ch->olen,
|
||||||
|
&ch->oalloc, ch->out_cap);
|
||||||
|
uint64_t fcap = ch->out_cap;
|
||||||
|
const char *fname = "stdout";
|
||||||
|
if (frc == 0) {
|
||||||
|
frc = proc_drain_fd(&ch->efd, &ch->ebuf, &ch->elen,
|
||||||
|
&ch->ealloc, ch->err_cap);
|
||||||
|
fcap = ch->err_cap;
|
||||||
|
fname = "stderr";
|
||||||
|
}
|
||||||
|
fb->dl_active = 0;
|
||||||
|
if (frc != 0) {
|
||||||
|
int rderr = errno;
|
||||||
|
proc_slot_close(vm, ch);
|
||||||
|
if (frc == -1) {
|
||||||
|
snprintf(proc_msg, sizeof proc_msg,
|
||||||
|
"process %s cap %llu bytes exceeded", fname,
|
||||||
|
(unsigned long long)fcap);
|
||||||
|
*msg = proc_msg;
|
||||||
|
return WO_T_IO;
|
||||||
|
}
|
||||||
|
if (frc == -2) {
|
||||||
|
*msg = "out of memory";
|
||||||
|
return WO_T_OOM;
|
||||||
|
}
|
||||||
|
*msg = strerror(rderr);
|
||||||
|
return WO_T_IO;
|
||||||
|
}
|
||||||
|
wo_hdr *o = record_of(vm, cls_id, 3, msg);
|
||||||
|
if (!o) {
|
||||||
|
proc_slot_close(vm, ch);
|
||||||
|
return cls_id >= vm->mod->class_cnt ? WO_T_BOUNDS : WO_T_OOM;
|
||||||
|
}
|
||||||
|
wo_str *out = wo_str_new(rt, ch->obuf ? ch->obuf : "", (uint32_t)ch->olen);
|
||||||
|
wo_str *errs = wo_str_new(rt, ch->ebuf ? ch->ebuf : "", (uint32_t)ch->elen);
|
||||||
|
proc_slot_close(vm, ch);
|
||||||
|
if (!out || !errs) {
|
||||||
|
if (out) wo_str_free(rt, out);
|
||||||
|
if (errs) wo_str_free(rt, errs);
|
||||||
|
wo_drop_obj(rt, o);
|
||||||
|
*msg = "out of memory";
|
||||||
|
return WO_T_OOM;
|
||||||
|
}
|
||||||
|
uint64_t *fp = wo_fields(o);
|
||||||
|
fp[0] = (uint64_t)(int64_t)(WIFEXITED(status) ? WEXITSTATUS(status) : -1);
|
||||||
|
fp[1] = (uint64_t)(uintptr_t)out;
|
||||||
|
fp[2] = (uint64_t)(uintptr_t)errs;
|
||||||
|
R[A] = (uint64_t)(uintptr_t)o;
|
||||||
|
return 0;
|
||||||
}
|
}
|
||||||
wo_hdr *o = record_of(vm, R[B + 2], 3, msg);
|
|
||||||
if (!o) return R[B + 2] >= vm->mod->class_cnt ? WO_T_BOUNDS : WO_T_OOM;
|
if (stop_pending()) { /* told to stop: no orphan survives it */
|
||||||
wo_str *out = wo_str_new(rt, obuf, (uint32_t)olen);
|
fb->dl_active = 0;
|
||||||
wo_str *errs = wo_str_new(rt, ebuf, (uint32_t)elen);
|
proc_slot_kill(vm, ch);
|
||||||
if (!out || !errs) {
|
return WO_SYS_STOPPED;
|
||||||
if (out) wo_str_free(rt, out);
|
|
||||||
if (errs) wo_str_free(rt, errs);
|
|
||||||
wo_drop_obj(rt, o);
|
|
||||||
*msg = "out of memory";
|
|
||||||
return WO_T_OOM;
|
|
||||||
}
|
}
|
||||||
uint64_t *fp = wo_fields(o);
|
if (fb->dl_at > 0 && dnow >= fb->dl_at) { /* ---- deadline: refuse */
|
||||||
fp[0] = (uint64_t)(int64_t)(WIFEXITED(status) ? WEXITSTATUS(status) : -1);
|
fb->dl_active = 0;
|
||||||
fp[1] = (uint64_t)(uintptr_t)out;
|
snprintf(proc_msg, sizeof proc_msg,
|
||||||
fp[2] = (uint64_t)(uintptr_t)errs;
|
"process deadline exceeded after %lld ms",
|
||||||
R[A] = (uint64_t)(uintptr_t)o;
|
(long long)(isdl && (int64_t)R[B + 2] > 0
|
||||||
return 0;
|
? (int64_t)R[B + 2]
|
||||||
|
: WO_PROC_DL_DEFAULT));
|
||||||
|
*msg = proc_msg;
|
||||||
|
proc_slot_kill(vm, ch);
|
||||||
|
return WO_T_IO;
|
||||||
|
}
|
||||||
|
/* ---- child alive, nothing more ready: park on the bundle */
|
||||||
|
fb->park_fd = ch->epfd;
|
||||||
|
fb->park_events = POLLIN;
|
||||||
|
fb->park_deadline = fb->dl_at;
|
||||||
|
fb->park_done = 0;
|
||||||
|
return WO_SYS_PARKED;
|
||||||
}
|
}
|
||||||
default:
|
default:
|
||||||
*msg = "unknown stdlib builtin";
|
*msg = "unknown stdlib builtin";
|
||||||
|
|
|
||||||
|
|
@ -754,6 +754,7 @@ int wo_vm_init(wo_vm *vm, const wo_module *mod, size_t heap_cap) {
|
||||||
}
|
}
|
||||||
|
|
||||||
void wo_vm_destroy(wo_vm *vm) {
|
void wo_vm_destroy(wo_vm *vm) {
|
||||||
|
wo_proc_reap_all(vm); /* iteration 42: no child outlives its shard */
|
||||||
/* iteration 35: the fiber pool dies with the vm */
|
/* iteration 35: the fiber pool dies with the vm */
|
||||||
while (vm->fib_pool) {
|
while (vm->fib_pool) {
|
||||||
wo_fiber *fb = vm->fib_pool;
|
wo_fiber *fb = vm->fib_pool;
|
||||||
|
|
@ -862,6 +863,7 @@ wo_fiber *wo_vm_spawn_fiber(wo_vm *vm, uint32_t method_idx, const uint64_t *args
|
||||||
* queued fibers die as cleanly as trapped ones), then free it if it is a
|
* queued fibers die as cleanly as trapped ones), then free it if it is a
|
||||||
* spawned one. `vm->cur` is borrowed to do it, restored after. */
|
* spawned one. `vm->cur` is borrowed to do it, restored after. */
|
||||||
static void fib_reap(wo_vm *vm, wo_fiber *fb) {
|
static void fib_reap(wo_vm *vm, wo_fiber *fb) {
|
||||||
|
wo_proc_abandon(vm, fb); /* iteration 42: its child dies with it */
|
||||||
wo_fiber *save = vm->cur;
|
wo_fiber *save = vm->cur;
|
||||||
vm->cur = fb;
|
vm->cur = fb;
|
||||||
vm_unwind(vm, 0);
|
vm_unwind(vm, 0);
|
||||||
|
|
|
||||||
|
|
@ -99,6 +99,11 @@ typedef struct wo_fiber {
|
||||||
* the original deadline. */
|
* the original deadline. */
|
||||||
int dl_active;
|
int dl_active;
|
||||||
int64_t dl_at; /* wall ms */
|
int64_t dl_at; /* wall ms */
|
||||||
|
/* iteration 42: the in-flight child while parked inside proc.run — a
|
||||||
|
* slot in the vm's children table (sysio.c owns the protocol). NULL
|
||||||
|
* when no run is in flight. A reaped fiber's child is killed with it
|
||||||
|
* (wo_proc_abandon from fib_reap). */
|
||||||
|
struct wo_child *proc_st;
|
||||||
} wo_fiber;
|
} wo_fiber;
|
||||||
|
|
||||||
/* arc stage 3: park_fd sentinel — PARKED with NO plane wait; the wake is
|
/* arc stage 3: park_fd sentinel — PARKED with NO plane wait; the wake is
|
||||||
|
|
@ -159,6 +164,29 @@ typedef struct wo_actor {
|
||||||
struct wo_actor *next_all; /* the vm's all-actors list */
|
struct wo_actor *next_all; /* the vm's all-actors list */
|
||||||
} wo_actor;
|
} wo_actor;
|
||||||
|
|
||||||
|
/* iteration 42: one live child process (proc.run in flight). The slot is
|
||||||
|
* the cross-park state: the _dl retry protocol re-executes the builtin,
|
||||||
|
* and this is where a re-entry finds its buffers, fds and caps. Slots
|
||||||
|
* live in the owning shard's vm (no locks — one thread), capped at
|
||||||
|
* WO_PROC_MAX; the claim failing closed IS the concurrency ceiling. */
|
||||||
|
typedef struct wo_child {
|
||||||
|
int used;
|
||||||
|
int pid;
|
||||||
|
int pidfd, epfd; /* pidfd_open handle; the epoll bundle the fiber parks on */
|
||||||
|
int ofd, efd; /* pipe read ends, O_NONBLOCK; -1 once EOF-closed */
|
||||||
|
char *obuf, *ebuf;
|
||||||
|
size_t olen, elen, oalloc, ealloc;
|
||||||
|
uint64_t out_cap, err_cap;
|
||||||
|
struct wo_fiber *owner;
|
||||||
|
} wo_child;
|
||||||
|
#define WO_PROC_MAX 32u
|
||||||
|
|
||||||
|
/* sysio.c: kill+reap the fiber's in-flight child, if any (fib_reap), and
|
||||||
|
* every live child on the shard (wo_vm_destroy / engine stop). */
|
||||||
|
struct wo_vm;
|
||||||
|
void wo_proc_abandon(struct wo_vm *vm, wo_fiber *fb);
|
||||||
|
void wo_proc_reap_all(struct wo_vm *vm);
|
||||||
|
|
||||||
/* iteration 24: the one mailbox cap (default 1024, WO_MAILBOX overrides
|
/* iteration 24: the one mailbox cap (default 1024, WO_MAILBOX overrides
|
||||||
* at boot — soak tests shrink it to force the fail-fast policy). */
|
* at boot — soak tests shrink it to force the fail-fast policy). */
|
||||||
extern uint32_t wo_mailbox_cap;
|
extern uint32_t wo_mailbox_cap;
|
||||||
|
|
@ -201,6 +229,9 @@ typedef struct wo_vm {
|
||||||
/* iteration 24 T5: this shard's armed timers (unsorted list — the
|
/* iteration 24 T5: this shard's armed timers (unsorted list — the
|
||||||
* deadline scan is already linear; a wheel is measured-later work) */
|
* deadline scan is already linear; a wheel is measured-later work) */
|
||||||
wo_timer *timers;
|
wo_timer *timers;
|
||||||
|
/* iteration 42: this shard's live children (proc.run in flight) */
|
||||||
|
wo_child children[WO_PROC_MAX];
|
||||||
|
uint32_t nchildren;
|
||||||
/* iteration 35, uring backend: the shard's ONE deadline tick — a
|
/* iteration 35, uring backend: the shard's ONE deadline tick — a
|
||||||
* TIMEOUT op with a sentinel user_data armed for the nearest fd-park
|
* TIMEOUT op with a sentinel user_data armed for the nearest fd-park
|
||||||
* deadline (fd parks keep exactly one POLL op each; expiry wakes them
|
* deadline (fd parks keep exactly one POLL op each; expiry wakes them
|
||||||
|
|
|
||||||
|
|
@ -503,9 +503,16 @@ enum {
|
||||||
* never needs manual cleanup) */
|
* never needs manual cleanup) */
|
||||||
WO_B_NET_PEER = 95, /* (fd) -> Text: "ip:port" for TCP peers,
|
WO_B_NET_PEER = 95, /* (fd) -> Text: "ip:port" for TCP peers,
|
||||||
* "unix" for unix-socket peers, "" on error */
|
* "unix" for unix-socket peers, "" on error */
|
||||||
|
/* ---- iteration 42: bounded subprocess. proc.run keeps id 56 and
|
||||||
|
* gains the same machinery with the named defaults (30 000 ms, 1 MiB
|
||||||
|
* out, 64 KiB err). Bound violations KILL the child and trap WO_T_IO
|
||||||
|
* naming the bound — never nil, never truncation. ---- */
|
||||||
|
WO_B_PROC_RUN_DL = 96, /* (cmd, multi Text args, deadline_ms,
|
||||||
|
* out_cap, err_cap, cls) -> Proc {code, out,
|
||||||
|
* err}; ms/caps <= 0 pick the default */
|
||||||
};
|
};
|
||||||
|
|
||||||
#define WO_B_MAX 95u
|
#define WO_B_MAX 96u
|
||||||
/* ids at or above this one live in sysio.c, not builtin.c */
|
/* ids at or above this one live in sysio.c, not builtin.c */
|
||||||
#define WO_B_SYS_FIRST WO_B_FS_EXISTS
|
#define WO_B_SYS_FIRST WO_B_FS_EXISTS
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -4,9 +4,13 @@
|
||||||
* stdout, the execvp-failed 127) so the parked rework in Task 3 has a
|
* stdout, the execvp-failed 127) so the parked rework in Task 3 has a
|
||||||
* baseline that must not move. Later tasks add the bounds legs: deadline,
|
* baseline that must not move. Later tasks add the bounds legs: deadline,
|
||||||
* output caps, ceiling, fd hygiene, stop/unwind reaping. */
|
* output caps, ceiling, fd hygiene, stop/unwind reaping. */
|
||||||
|
#define _POSIX_C_SOURCE 200809L /* sigaction/clock_gettime under -std=c11 */
|
||||||
|
#include <signal.h>
|
||||||
#include <stdio.h>
|
#include <stdio.h>
|
||||||
#include <stdlib.h>
|
#include <stdlib.h>
|
||||||
#include <string.h>
|
#include <string.h>
|
||||||
|
#include <time.h>
|
||||||
|
#include <unistd.h>
|
||||||
|
|
||||||
#include "gc.h"
|
#include "gc.h"
|
||||||
#include "loader.h"
|
#include "loader.h"
|
||||||
|
|
@ -84,6 +88,34 @@ static int run_proc(const char *cmd, const char **args, uint32_t nargs,
|
||||||
return rc;
|
return rc;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/* Like run_proc but keeps only the exit code and output LENGTHS — for legs
|
||||||
|
* whose output is too big to copy out. */
|
||||||
|
static int run_proc_lens(const char *cmd, const char **args, uint32_t nargs,
|
||||||
|
int64_t *pcode, size_t *outlen, size_t *errlen,
|
||||||
|
wo_err *err) {
|
||||||
|
size_t len;
|
||||||
|
uint8_t *img = run_module(cmd, args, nargs, &len);
|
||||||
|
wo_module mod;
|
||||||
|
char lerr[256];
|
||||||
|
T_EQ(wo_load_buf(&mod, img, len, lerr, sizeof lerr), 0);
|
||||||
|
T_EQ(wo_vm_init(&VM, &mod, 1 << 20), 0);
|
||||||
|
uint64_t ret = 0;
|
||||||
|
int rc = wo_vm_call(&VM, 0, NULL, 0, &ret, err);
|
||||||
|
if (rc == 0) {
|
||||||
|
wo_hdr *o = (wo_hdr *)(uintptr_t)ret;
|
||||||
|
T_CHECK(o != NULL);
|
||||||
|
uint64_t *fp = wo_fields(o);
|
||||||
|
*pcode = (int64_t)fp[0];
|
||||||
|
*outlen = ((const wo_str *)(uintptr_t)fp[1])->len;
|
||||||
|
*errlen = ((const wo_str *)(uintptr_t)fp[2])->len;
|
||||||
|
wo_drop_obj(&VM.rt, o);
|
||||||
|
}
|
||||||
|
wo_vm_destroy(&VM);
|
||||||
|
wo_module_free(&mod);
|
||||||
|
free(img);
|
||||||
|
return rc;
|
||||||
|
}
|
||||||
|
|
||||||
static char OUTBUF[1 << 16], ERRBUF[1 << 16];
|
static char OUTBUF[1 << 16], ERRBUF[1 << 16];
|
||||||
|
|
||||||
/* a child that exits 0 with known stdout */
|
/* a child that exits 0 with known stdout */
|
||||||
|
|
@ -120,9 +152,45 @@ static void test_missing(void) {
|
||||||
T_EQ(code, 127);
|
T_EQ(code, 127);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
static void on_alarm(int sig) { (void)sig; } /* interrupt, don't die */
|
||||||
|
|
||||||
|
static int64_t mono_ms(void) {
|
||||||
|
struct timespec ts;
|
||||||
|
clock_gettime(CLOCK_MONOTONIC, &ts);
|
||||||
|
return (int64_t)ts.tv_sec * 1000 + ts.tv_nsec / 1000000;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* The Task 3 deadlock repro: a child that writes 200 KB to stdout and only
|
||||||
|
* THEN a line to stderr. The old sequential drain stopped reading stdout at
|
||||||
|
* its 8 KiB cap, the child blocked on the full pipe and never reached
|
||||||
|
* stderr, and the parent sat in the stderr read forever. The leg demands
|
||||||
|
* completion well under the 5 s alarm — the alarm exists so the OLD code
|
||||||
|
* FAILS instead of hanging the suite. */
|
||||||
|
static void test_chatty_child_completes(void) {
|
||||||
|
struct sigaction sa;
|
||||||
|
memset(&sa, 0, sizeof sa);
|
||||||
|
sa.sa_handler = on_alarm; /* no SA_RESTART: reads must EINTR */
|
||||||
|
sigaction(SIGALRM, &sa, NULL);
|
||||||
|
alarm(5);
|
||||||
|
int64_t t0 = mono_ms();
|
||||||
|
const char *args[] = {"-c", "head -c 200000 /dev/zero; echo done >&2"};
|
||||||
|
int64_t code = -99;
|
||||||
|
size_t outlen = 0, errlen = 0;
|
||||||
|
wo_err err;
|
||||||
|
int rc = run_proc_lens("sh", args, 2, &code, &outlen, &errlen, &err);
|
||||||
|
int64_t elapsed = mono_ms() - t0;
|
||||||
|
alarm(0);
|
||||||
|
T_EQ(rc, 0);
|
||||||
|
T_EQ(code, 0);
|
||||||
|
T_CHECK(elapsed < 4000); /* the old drain needs the alarm to escape */
|
||||||
|
T_EQ(outlen, 200000u);
|
||||||
|
T_EQ(errlen, 5u); /* "done\n" */
|
||||||
|
}
|
||||||
|
|
||||||
int main(void) {
|
int main(void) {
|
||||||
test_echo();
|
test_echo();
|
||||||
test_false();
|
test_false();
|
||||||
test_missing();
|
test_missing();
|
||||||
|
test_chatty_child_completes();
|
||||||
return t_report("test_proc");
|
return t_report("test_proc");
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue