feat(rt2): proc.spawn/wait_dl/signal — the streaming child

- a child is fds: Child {id, stdin, stdout, stderr}, driven by the
  existing net verbs (echo leg proves cat round-trip through write_dl/
  read_dl); caller owns the fds, the runtime owns pid + pidfd
- wait_dl parks on the pidfd: code on exit, nil at the deadline with the
  child untouched; one waiter per id, a second refuses by name; stale
  ids refused via a generation counter in the handle
- proc.signal through pidfd_send_signal; actor_die kills the streaming
  children the dying actor owns; dead fibers cannot linger as waiters
- ids 97-107 registered wholesale (wob.h, loader arities, dispatch
  bound); Child + Signal predeclared records in types.ml; unimplemented
  ids trap at the default case until their task lands
- test_proc 168/0 (echo, wait trio, one-waiter refusal, 200-round churn
  fd-flat), suite ASan clean, woc-test green

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
(cherry picked from commit 9be87f159f1bf9cdd509ceed160e7ea518fde46c)
This commit is contained in:
shoney.arickathil 2026-09-01 23:52:40 +02:00
parent 3190b609af
commit 803ff0b790
8 changed files with 595 additions and 6 deletions

View file

@ -257,9 +257,26 @@ let proc_record_name = "Proc"
let proc_record_fields : (string * field_ty) list =
[ ("code", Scalar "Int"); ("out", Scalar "Text"); ("err", Scalar "Text") ]
(* runtime-v2 1: the streaming child. The fds are ordinary conn-shaped
Ints the net verbs drive; stderr is -1 on a PTY child (master carries
both streams). The id refuses stale handles by name at runtime. *)
let child_record_name = "Child"
let child_record_fields : (string * field_ty) list =
[ ("id", Scalar "Int"); ("stdin", Scalar "Int"); ("stdout", Scalar "Int");
("stderr", Scalar "Int") ]
(* runtime-v2 3: what signal.on delivers — a fresh record per arrival
(message payloads must be heap objects; the runtime drops them). *)
let signal_record_name = "Signal"
let signal_record_fields : (string * field_ty) list = [ ("sig", Scalar "Int") ]
let predeclared_records : (string * (string * field_ty) list) list =
[ (error_record_name, error_record_fields); (stat_record_name, stat_record_fields);
(time_record_name, time_record_fields); (proc_record_name, proc_record_fields) ]
(time_record_name, time_record_fields); (proc_record_name, proc_record_fields);
(child_record_name, child_record_fields);
(signal_record_name, signal_record_fields) ]
(* One member of a reserved stdlib module (`fs.stat`, `net.write`, ...).
[sm_builtin] is its .wob builtin id (runtime/src/wob.h); [sm_record] names
@ -317,6 +334,12 @@ let stdlib_members : stdlib_member list =
(<= 0 picks the default: 30 000 ms / 1 MiB / 64 KiB). A bound
violation kills the child and traps WO_T_IO naming the bound. *)
m "proc" "run_dl" 5 96 (Some (TNullable (TScalar proc_record_name))) (Some proc_record_name);
(* runtime-v2 1: the streaming child — fds the net verbs drive; the
caller closes them with net.close. wait_dl: nil = still running at
the deadline (child untouched); one waiter per id. *)
m "proc" "spawn" 2 97 (Some (TNullable (TScalar child_record_name))) (Some child_record_name);
m "proc" "wait_dl" 2 98 (Some (TNullable (TScalar "Int"))) None;
m "proc" "signal" 2 99 None None;
(* json — both members are lowered specially (emit.ml): encode needs its
argument's static kind, and decode has no type until an `as` names one,
so neither goes through the generic builtin path. They are listed here

View file

@ -172,7 +172,7 @@ int wo_builtin(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) {
if (C == WO_B_JSON_ENCODE || C == WO_B_JSON_DECODE)
return wo_builtin_json(vm, R, ins, msg);
if ((C >= WO_B_SYS_FIRST && C <= WO_B_PROC_RUN) || C == WO_B_TIME_TICKS
|| (C >= WO_B_NET_READ_DL && C <= WO_B_PROC_RUN_DL))
|| (C >= WO_B_NET_READ_DL && C <= WO_B_NET_CONNECT_UNIX))
return wo_builtin_sys(vm, R, ins, msg);
if (C >= WO_B_SHA1 && C <= WO_B_HMAC_SHA256)
return wo_builtin_crypto(vm, R, ins, msg);

View file

@ -73,6 +73,11 @@ static const uint8_t b_arity[WO_B_MAX + 1] = {
[WO_B_NET_ACCEPT] = 1, [WO_B_NET_READ] = 2, [WO_B_NET_WRITE] = 2,
[WO_B_NET_CLOSE] = 1, [WO_B_PROC_RUN] = 3,
[WO_B_PROC_RUN_DL] = 6, /* iteration 42: cmd, argv, dl, ocap, ecap, cls */
/* runtime-v2 (ids 97-107) */
[WO_B_PROC_SPAWN] = 3, [WO_B_PROC_WAIT_DL] = 2, [WO_B_PROC_SIGNAL] = 2,
[WO_B_PROC_SPAWN_PTY] = 5, [WO_B_PROC_RESIZE] = 3,
[WO_B_SIGNAL_ON] = 3, [WO_B_TERM_RAW] = 1, [WO_B_TERM_RESTORE] = 1,
[WO_B_NET_SEND_FD] = 2, [WO_B_NET_RECV_FD] = 1, [WO_B_NET_CONNECT_UNIX] = 1,
/* json (json.c): encode takes the value's static kind, decode the class
id to build */
[WO_B_JSON_ENCODE] = 2, [WO_B_JSON_DECODE] = 2, [WO_B_MAP_GET_OPT] = 2,

View file

@ -174,16 +174,20 @@ static wo_str *read_range(wo_rt *rt, int fd, off_t off, size_t want, const char
* per thread */
static _Thread_local char proc_msg[96];
/* release everything a slot holds; the child must already be reaped */
/* release everything a slot holds; the child must already be reaped.
* A streaming slot's stdio fds are the CALLER's (never closed here —
* fd numbers get recycled); the master dup is the slot's own. */
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);
if (ch->master_dup > 0) close(ch->master_dup);
free(ch->obuf);
free(ch->ebuf);
if (ch->owner) ch->owner->proc_st = NULL;
memset(ch, 0, sizeof *ch);
ch->master_dup = -1;
vm->nchildren--;
}
@ -197,6 +201,17 @@ static void proc_slot_kill(wo_vm *vm, wo_child *ch) {
void wo_proc_abandon(wo_vm *vm, wo_fiber *fb) {
if (fb->proc_st) proc_slot_kill(vm, fb->proc_st);
/* a dead fiber must not linger as a streaming child's waiter */
for (uint32_t i = 0; i < WO_PROC_MAX; i++)
if (vm->children[i].used && vm->children[i].waiter == fb)
vm->children[i].waiter = NULL;
}
/* runtime-v2 1: a dying actor's streaming children die with it */
void wo_proc_abandon_actor(wo_vm *vm, struct wo_actor *a) {
for (uint32_t i = 0; i < WO_PROC_MAX; i++)
if (vm->children[i].used && vm->children[i].owner_actor == a)
proc_slot_kill(vm, &vm->children[i]);
}
void wo_proc_reap_all(wo_vm *vm) {
@ -204,6 +219,57 @@ void wo_proc_reap_all(wo_vm *vm) {
if (vm->children[i].used) proc_slot_kill(vm, &vm->children[i]);
}
/* the language-visible child id: (gen << 6) | slot index. Stale or
* foreign ids refuse by name instead of touching a recycled slot. */
static wo_child *proc_slot_by_id(wo_vm *vm, uint64_t id, const char **msg) {
uint32_t idx = (uint32_t)(id & 63u);
wo_child *ch = idx < WO_PROC_MAX ? &vm->children[idx] : NULL;
if (!ch || !ch->used || !ch->streaming || ch->gen != (uint32_t)(id >> 6)) {
*msg = "process id is not a live child";
return NULL;
}
return ch;
}
/* argv marshalling shared by the streaming spawn forms. argv[0] is the
* command; the multi supplies the rest; buffers are the caller's. */
static int proc_argv(uint64_t vcmd, uint64_t vargv, char *path, size_t pathcap,
char (*argbuf)[512], char **argv, const char **msg) {
if (cstr_of(vcmd, path, pathcap, msg)) return -1;
wo_multi *m = (wo_multi *)(uintptr_t)vargv;
if (!m || m->h.class_id != WO_CLS_MULTI || m->elem_kind != WO_K_TEXT) {
*msg = "`proc.spawn` needs a `multi Text` of arguments";
return -1;
}
if (m->len > 62) {
*msg = "too many process arguments";
return -1;
}
argv[0] = path;
for (uint32_t i = 0; i < m->len; i++) {
const wo_str *a = (const wo_str *)(uintptr_t)m->items[i];
if (!a || a->h.class_id != WO_CLS_STR || a->len + 1 > 512) {
*msg = "process argument is not a short text";
return -1;
}
memcpy(argbuf[i], a->data, a->len);
argbuf[i][a->len] = '\0';
argv[i + 1] = argbuf[i];
}
argv[m->len + 1] = NULL;
return 0;
}
/* claim a slot or refuse by name (shared by run/run_dl/spawn forms) */
static wo_child *proc_slot_claim(wo_vm *vm, const char **msg) {
for (uint32_t i = 0; i < WO_PROC_MAX; i++)
if (!vm->children[i].used) return &vm->children[i];
snprintf(proc_msg, sizeof proc_msg,
"process ceiling: %u live children on this shard", WO_PROC_MAX);
*msg = proc_msg;
return NULL;
}
/* 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,
@ -1077,6 +1143,148 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) {
fb->park_done = 0;
return WO_SYS_PARKED;
}
/* ---- runtime-v2 1: the streaming child --------------------------- */
case WO_B_PROC_SPAWN: { /* Child: 0 id, 1 stdin, 2 stdout, 3 stderr.
* The caller owns the three fds (net verbs
* drive them, net.close releases them); the
* runtime owns pid + pidfd. */
char argbuf[62][512];
char *argv[64];
if (proc_argv(R[B], R[B + 1], path, sizeof path, argbuf, argv, msg))
return WO_T_BOUNDS;
wo_child *ch = proc_slot_claim(vm, msg);
if (!ch) return WO_T_IO;
int ip[2], op[2], ep[2];
if (pipe(ip) != 0) {
*msg = strerror(errno);
return WO_T_IO;
}
if (pipe(op) != 0) {
close(ip[0]); close(ip[1]);
*msg = strerror(errno);
return WO_T_IO;
}
if (pipe(ep) != 0) {
close(ip[0]); close(ip[1]); close(op[0]); close(op[1]);
*msg = strerror(errno);
return WO_T_IO;
}
pid_t pid = fork();
if (pid < 0) {
close(ip[0]); close(ip[1]); close(op[0]); close(op[1]);
close(ep[0]); close(ep[1]);
*msg = strerror(errno);
return WO_T_IO;
}
if (pid == 0) {
dup2(ip[0], STDIN_FILENO);
dup2(op[1], STDOUT_FILENO);
dup2(ep[1], STDERR_FILENO);
close(ip[0]); close(ip[1]); close(op[0]); close(op[1]);
close(ep[0]); close(ep[1]);
execvp(path, argv);
_exit(127);
}
close(ip[0]);
close(op[1]);
close(ep[1]);
fcntl(ip[1], F_SETFL, fcntl(ip[1], F_GETFL, 0) | O_NONBLOCK);
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);
if (pidfd < 0) {
int e = errno;
kill(pid, SIGKILL);
int st;
while (waitpid(pid, &st, 0) < 0 && errno == EINTR) {}
close(ip[1]); close(op[0]); close(ep[0]);
*msg = strerror(e);
return WO_T_IO;
}
memset(ch, 0, sizeof *ch);
ch->used = 1;
ch->streaming = 1;
ch->pid = (int)pid;
ch->pidfd = pidfd;
ch->epfd = -1;
ch->ofd = ch->efd = -1;
ch->master_dup = -1;
ch->gen = ++vm->proc_gen;
ch->owner_actor = vm->cur->actor; /* NULL = the program */
vm->nchildren++;
wo_hdr *o = record_of(vm, R[B + 2], 4, msg);
if (!o) {
close(ip[1]); close(op[0]); close(ep[0]);
proc_slot_kill(vm, ch);
return R[B + 2] >= vm->mod->class_cnt ? WO_T_BOUNDS : WO_T_OOM;
}
uint64_t *fp = wo_fields(o);
fp[0] = ((uint64_t)ch->gen << 6) | (uint64_t)(ch - vm->children);
fp[1] = (uint64_t)ip[1];
fp[2] = (uint64_t)op[0];
fp[3] = (uint64_t)ep[0];
R[A] = (uint64_t)(uintptr_t)o;
return 0;
}
case WO_B_PROC_WAIT_DL: { /* (id, ms) -> ?Int code; nil = deadline,
* child untouched. One waiter per id. */
wo_fiber *fb = vm->cur;
wo_child *ch = proc_slot_by_id(vm, R[B], msg);
if (!ch) {
fb->dl_active = 0;
return WO_T_IO;
}
if (ch->waiter && ch->waiter != fb) {
fb->dl_active = 0;
*msg = "child already has a waiter";
return WO_T_IO;
}
struct timespec dts;
clock_gettime(CLOCK_REALTIME, &dts);
int64_t dnow = (int64_t)dts.tv_sec * 1000 + dts.tv_nsec / 1000000;
if (!fb->dl_active) {
int64_t ms = (int64_t)R[B + 1];
fb->dl_active = 1;
fb->dl_at = ms > 0 ? dnow + ms : 0;
}
int status = 0;
pid_t r = waitpid(ch->pid, &status, WNOHANG);
if (r == (pid_t)ch->pid) {
fb->dl_active = 0;
proc_slot_close(vm, ch);
R[A] = (uint64_t)(int64_t)(WIFEXITED(status) ? WEXITSTATUS(status)
: -1);
return 0;
}
if (stop_pending()) {
fb->dl_active = 0;
proc_slot_kill(vm, ch);
return WO_SYS_STOPPED;
}
if (fb->dl_at > 0 && dnow >= fb->dl_at) {
/* the deadline answers nil; the CHILD is untouched */
fb->dl_active = 0;
ch->waiter = NULL;
R[A] = WO_NIL_SCALAR;
return 0;
}
ch->waiter = fb;
fb->park_fd = ch->pidfd;
fb->park_events = POLLIN;
fb->park_deadline = fb->dl_at;
fb->park_done = 0;
return WO_SYS_PARKED;
}
case WO_B_PROC_SIGNAL: { /* (id, sig) -> 0 through the pidfd */
wo_child *ch = proc_slot_by_id(vm, R[B], msg);
if (!ch) return WO_T_IO;
if (syscall(SYS_pidfd_send_signal, ch->pidfd, (int)R[B + 1], NULL, 0) != 0) {
*msg = strerror(errno);
return WO_T_IO;
}
R[A] = 0;
return 0;
}
default:
*msg = "unknown stdlib builtin";
return WO_T_EXPLICIT;

View file

@ -987,6 +987,7 @@ static void call_reply_to(wo_vm *vm, wo_fiber *caller, uint32_t caller_shard,
* (addresses are copyable scalars that may still be sent to). */
static void actor_die(wo_vm *vm, wo_actor *a, wo_fiber *delivery) {
a->dead = 1;
wo_proc_abandon_actor(vm, a); /* runtime-v2 1: its children die with it */
if (delivery->cur_msg) {
wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)delivery->cur_msg);
delivery->cur_msg = 0;

View file

@ -178,13 +178,27 @@ typedef struct wo_child {
size_t olen, elen, oalloc, ealloc;
uint64_t out_cap, err_cap;
struct wo_fiber *owner;
/* runtime-v2 1/2: the STREAMING child. The caller owns the stdio fds
* (Child.stdin/stdout/stderr, closed with net.close — the slot never
* touches them: fd numbers get recycled); the slot owns pid + pidfd
* and, for a PTY child, a private dup of the master for resize.
* `gen` makes the language-visible id ((gen << 6) | index) refuse
* stale handles by name. One waiter at a time parks on the pidfd. */
int streaming;
uint32_t gen;
int master_dup; /* -1 = pipe child */
struct wo_fiber *waiter; /* the one wait_dl parker, NULL when none */
struct wo_actor *owner_actor; /* NULL = the program owns it */
} 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). */
* every live child on the shard (wo_vm_destroy / engine stop).
* runtime-v2: abandon_actor kills the streaming children a dying actor
* owns (actor_die); wo_proc_abandon also clears a dead fiber's waiter. */
struct wo_vm;
void wo_proc_abandon(struct wo_vm *vm, wo_fiber *fb);
void wo_proc_abandon_actor(struct wo_vm *vm, struct wo_actor *a);
void wo_proc_reap_all(struct wo_vm *vm);
/* iteration 24: the one mailbox cap (default 1024, WO_MAILBOX overrides
@ -232,6 +246,7 @@ typedef struct wo_vm {
/* iteration 42: this shard's live children (proc.run in flight) */
wo_child children[WO_PROC_MAX];
uint32_t nchildren;
uint32_t proc_gen; /* runtime-v2 1: claim counter behind child ids */
/* iteration 35, uring backend: the shard's ONE deadline tick — a
* 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

View file

@ -510,9 +510,40 @@ enum {
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 */
/* ---- runtime-v2: the runtime beyond sockets (track spec
* 2026-09-01). A child/received fd is an ORDINARY fd the existing
* net verbs drive; these are acquisition verbs only. ---- */
WO_B_PROC_SPAWN = 97, /* (cmd, args, cls) -> Child {id, stdin,
* stdout, stderr}: streaming child, pipes;
* caller owns the three fds (net.close),
* the runtime owns pid+pidfd */
WO_B_PROC_WAIT_DL = 98, /* (id, ms) -> ?Int exit code; nil = still
* running at the deadline (child untouched).
* One waiter per id — a second refuses */
WO_B_PROC_SIGNAL = 99, /* (id, sig) -> 0: pidfd_send_signal */
WO_B_PROC_SPAWN_PTY = 100,/* (cmd, args, cols, rows, cls) -> Child:
* stdin==stdout=PTY master (raw), stderr -1 */
WO_B_PROC_RESIZE = 101, /* (id, cols, rows) -> 0: TIOCSWINSZ; refuses
* by name on a pipe child */
WO_B_SIGNAL_ON = 102, /* (sig, addr, cls) -> 0: standing
* subscription; each arrival delivers a
* fresh Signal {sig} record (payloads must
* be heap objects — vm.c drops them).
* SIGTERM/SIGINT refused: the stop latch
* stays the engine's */
WO_B_TERM_RAW = 103, /* (fd) -> 0: save termios, cfmakeraw.
* Restore is a RUNTIME obligation on
* unwind/stop — no wrecked tty */
WO_B_TERM_RESTORE = 104, /* (fd) -> 0: restore the saved termios */
WO_B_NET_SEND_FD = 105, /* (conn, fd) -> Bool: SCM_RIGHTS, one fd;
* unix sockets only, refuses by name */
WO_B_NET_RECV_FD = 106, /* (conn) -> ?Int: the received fd, nil if
* the peer sent plain bytes */
WO_B_NET_CONNECT_UNIX = 107, /* (path) -> Int: AF_UNIX client fd,
* nonblocking */
};
#define WO_B_MAX 96u
#define WO_B_MAX 107u
/* ids at or above this one live in sysio.c, not builtin.c */
#define WO_B_SYS_FIRST WO_B_FS_EXISTS

View file

@ -579,6 +579,308 @@ static void test_stop_kills_child(void) {
free(img);
}
/* ---- runtime-v2 1: the streaming child ---------------------------------
* Child record class (id, stdin, stdout, stderr — all scalar) is built
* into each module; the fds are driven by the NET verbs, which is the
* whole design. */
/* echo module: spawn `cat`, write "hi\n" to Child.stdin, read it back
* from Child.stdout, close all three fds, return the text. */
static uint8_t *stream_echo_module(size_t *len) {
wb_t *b = wb_new();
uint32_t kchild = wb_const_text(b, "Child");
uint32_t km = wb_const_text(b, "main");
uint32_t kcmd = wb_const_text(b, "cat");
uint32_t khi = wb_const_text(b, "hi\n");
uint8_t kinds[4] = {WO_K_SCALAR, WO_K_SCALAR, WO_K_SCALAR, WO_K_SCALAR};
uint32_t cls = wb_class(b, kchild, 0, kinds, 4);
uint32_t kcls = wb_const_int(b, (int64_t)cls);
uint32_t kms = wb_const_int(b, 3000);
uint32_t kmax = wb_const_int(b, 16);
uint32_t code[32];
uint32_t n = 0;
code[n++] = wo_ins_abx(WOP_LOADK, 1, (uint16_t)kcmd);
code[n++] = wo_ins_abc(WOP_BUILTIN, 2, WO_K_TEXT, WO_B_MULTI_NEW);
code[n++] = wo_ins_abx(WOP_LOADK, 3, (uint16_t)kcls);
code[n++] = wo_ins_abc(WOP_BUILTIN, 0, 1, WO_B_PROC_SPAWN);
code[n++] = wo_ins_abc(WOP_DROP, 2, 0, 0); /* argv multi */
code[n++] = wo_ins_abc(WOP_GETF, 1, 0, 1); /* stdin */
code[n++] = wo_ins_abc(WOP_GETF, 2, 0, 2); /* stdout */
code[n++] = wo_ins_abc(WOP_GETF, 3, 0, 3); /* stderr */
/* write "hi\n" with a deadline */
code[n++] = wo_ins_abc(WOP_MOVE, 4, 1, 0);
code[n++] = wo_ins_abx(WOP_LOADK, 5, (uint16_t)khi);
code[n++] = wo_ins_abx(WOP_LOADK, 6, (uint16_t)kms);
code[n++] = wo_ins_abc(WOP_BUILTIN, 7, 4, WO_B_NET_WRITE_DL);
/* read it back */
code[n++] = wo_ins_abc(WOP_MOVE, 4, 2, 0);
code[n++] = wo_ins_abx(WOP_LOADK, 5, (uint16_t)kmax);
code[n++] = wo_ins_abx(WOP_LOADK, 6, (uint16_t)kms);
code[n++] = wo_ins_abc(WOP_BUILTIN, 7, 4, WO_B_NET_READ_DL);
/* close the caller-owned fds */
code[n++] = wo_ins_abc(WOP_BUILTIN, 8, 1, WO_B_NET_CLOSE);
code[n++] = wo_ins_abc(WOP_BUILTIN, 8, 2, WO_B_NET_CLOSE);
code[n++] = wo_ins_abc(WOP_BUILTIN, 8, 3, WO_B_NET_CLOSE);
code[n++] = wo_ins_abc(WOP_DROP, 0, 0, 0); /* the Child record */
code[n++] = wo_ins_abc(WOP_RET, 7, 0, 0);
wb_method(b, km, WOB_NONE, 0, 9, code, n, NULL, 0, NULL, 0);
return wb_finish(b, len);
}
static void test_stream_echo(void) {
size_t len;
uint8_t *img = stream_echo_module(&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;
wo_err err;
memset(&err, 0, sizeof err);
T_EQ(wo_vm_call(&VM, 0, NULL, 0, &ret, &err), 0);
const wo_str *s = (const wo_str *)(uintptr_t)ret;
T_CHECK(s != NULL);
if (s) {
T_EQ(s->len, 3u);
T_CHECK(memcmp(s->data, "hi\n", 3) == 0);
wo_str_free(&VM.rt, (wo_str *)s);
}
wo_vm_destroy(&VM); /* cat (EOF'd) reaped here if still live */
int st;
T_EQ(waitpid(-1, &st, WNOHANG), -1);
T_EQ(errno, ECHILD);
wo_module_free(&mod);
free(img);
}
/* wait module: spawn `cmd arg`, close the fds, then EITHER one wait
* (ms1) returning its result, OR wait(ms1) -> signal(sig) -> wait(ms2)
* returning the second result. */
static uint8_t *stream_wait_module(const char *cmd, const char *arg,
int64_t ms1, int64_t sig, int64_t ms2,
size_t *len) {
wb_t *b = wb_new();
uint32_t kchild = wb_const_text(b, "Child");
uint32_t km = wb_const_text(b, "main");
uint32_t kcmd = wb_const_text(b, cmd);
uint32_t karg = arg ? wb_const_text(b, arg) : 0;
uint8_t kinds[4] = {WO_K_SCALAR, WO_K_SCALAR, WO_K_SCALAR, WO_K_SCALAR};
uint32_t cls = wb_class(b, kchild, 0, kinds, 4);
uint32_t kcls = wb_const_int(b, (int64_t)cls);
uint32_t kms1 = wb_const_int(b, ms1);
uint32_t ksig = wb_const_int(b, sig);
uint32_t kms2 = wb_const_int(b, ms2);
uint32_t code[40];
uint32_t n = 0;
code[n++] = wo_ins_abx(WOP_LOADK, 1, (uint16_t)kcmd);
code[n++] = wo_ins_abc(WOP_BUILTIN, 2, WO_K_TEXT, WO_B_MULTI_NEW);
if (arg) {
code[n++] = wo_ins_abx(WOP_LOADK, 3, (uint16_t)karg);
code[n++] = wo_ins_abc(WOP_BUILTIN, 4, 2, WO_B_MULTI_PUSH);
}
code[n++] = wo_ins_abx(WOP_LOADK, 3, (uint16_t)kcls);
code[n++] = wo_ins_abc(WOP_BUILTIN, 0, 1, WO_B_PROC_SPAWN);
code[n++] = wo_ins_abc(WOP_DROP, 2, 0, 0);
code[n++] = wo_ins_abc(WOP_GETF, 1, 0, 1);
code[n++] = wo_ins_abc(WOP_BUILTIN, 8, 1, WO_B_NET_CLOSE);
code[n++] = wo_ins_abc(WOP_GETF, 1, 0, 2);
code[n++] = wo_ins_abc(WOP_BUILTIN, 8, 1, WO_B_NET_CLOSE);
code[n++] = wo_ins_abc(WOP_GETF, 1, 0, 3);
code[n++] = wo_ins_abc(WOP_BUILTIN, 8, 1, WO_B_NET_CLOSE);
code[n++] = wo_ins_abc(WOP_GETF, 4, 0, 0); /* id stays in r4 */
code[n++] = wo_ins_abx(WOP_LOADK, 5, (uint16_t)kms1);
code[n++] = wo_ins_abc(WOP_BUILTIN, 7, 4, WO_B_PROC_WAIT_DL);
if (sig > 0) {
code[n++] = wo_ins_abx(WOP_LOADK, 5, (uint16_t)ksig);
code[n++] = wo_ins_abc(WOP_BUILTIN, 8, 4, WO_B_PROC_SIGNAL);
code[n++] = wo_ins_abx(WOP_LOADK, 5, (uint16_t)kms2);
code[n++] = wo_ins_abc(WOP_BUILTIN, 7, 4, WO_B_PROC_WAIT_DL);
}
code[n++] = wo_ins_abc(WOP_DROP, 0, 0, 0);
code[n++] = wo_ins_abc(WOP_RET, 7, 0, 0);
wb_method(b, km, WOB_NONE, 0, 9, code, n, NULL, 0, NULL, 0);
return wb_finish(b, len);
}
static int64_t run_stream_wait(const char *cmd, const char *arg, int64_t ms1,
int64_t sig, int64_t ms2, int expect_rc) {
size_t len;
uint8_t *img = stream_wait_module(cmd, arg, ms1, sig, ms2, &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;
wo_err err;
memset(&err, 0, sizeof err);
T_EQ(wo_vm_call(&VM, 0, NULL, 0, &ret, &err), expect_rc);
wo_vm_destroy(&VM);
int st;
T_EQ(waitpid(-1, &st, WNOHANG), -1);
T_EQ(errno, ECHILD);
wo_module_free(&mod);
free(img);
return (int64_t)ret;
}
static void test_stream_wait(void) {
/* a fast child answers its code */
T_EQ(run_stream_wait("true", NULL, 3000, 0, 0, 0), 0);
/* a slow child answers nil at the deadline (and is untouched, then
* swept by destroy — ECHILD proves the sweep) */
T_EQ((uint64_t)run_stream_wait("sleep", "10", 100, 0, 0, 0),
WO_NIL_SCALAR);
/* signal SIGKILL, then the wait observes the signal death (-1) */
T_EQ(run_stream_wait("sleep", "10", 100, 9, 3000, 0), -1);
}
/* double-wait: a fiber parks as the waiter; main tries second, refused */
static uint8_t *double_wait_module(size_t *len) {
wb_t *b = wb_new();
uint32_t kchild = wb_const_text(b, "Child");
uint32_t kspawn = wb_const_text(b, "spawner");
uint32_t kw = wb_const_text(b, "waiter");
uint32_t ksec = wb_const_text(b, "second");
uint32_t kcmd = wb_const_text(b, "sleep");
uint32_t karg = wb_const_text(b, "1");
uint8_t kinds[4] = {WO_K_SCALAR, WO_K_SCALAR, WO_K_SCALAR, WO_K_SCALAR};
uint32_t cls = wb_class(b, kchild, 0, kinds, 4);
uint32_t kcls = wb_const_int(b, (int64_t)cls);
uint32_t k3000 = wb_const_int(b, 3000);
uint32_t k200 = wb_const_int(b, 200);
uint32_t k100 = wb_const_int(b, 100);
{ /* spawner (method 0): spawn sleep 1, close fds, RET id */
uint32_t code[24];
uint32_t n = 0;
code[n++] = wo_ins_abx(WOP_LOADK, 1, (uint16_t)kcmd);
code[n++] = wo_ins_abc(WOP_BUILTIN, 2, WO_K_TEXT, WO_B_MULTI_NEW);
code[n++] = wo_ins_abx(WOP_LOADK, 3, (uint16_t)karg);
code[n++] = wo_ins_abc(WOP_BUILTIN, 4, 2, WO_B_MULTI_PUSH);
code[n++] = wo_ins_abx(WOP_LOADK, 3, (uint16_t)kcls);
code[n++] = wo_ins_abc(WOP_BUILTIN, 0, 1, WO_B_PROC_SPAWN);
code[n++] = wo_ins_abc(WOP_DROP, 2, 0, 0);
code[n++] = wo_ins_abc(WOP_GETF, 1, 0, 1);
code[n++] = wo_ins_abc(WOP_BUILTIN, 8, 1, WO_B_NET_CLOSE);
code[n++] = wo_ins_abc(WOP_GETF, 1, 0, 2);
code[n++] = wo_ins_abc(WOP_BUILTIN, 8, 1, WO_B_NET_CLOSE);
code[n++] = wo_ins_abc(WOP_GETF, 1, 0, 3);
code[n++] = wo_ins_abc(WOP_BUILTIN, 8, 1, WO_B_NET_CLOSE);
code[n++] = wo_ins_abc(WOP_GETF, 4, 0, 0);
code[n++] = wo_ins_abc(WOP_DROP, 0, 0, 0);
code[n++] = wo_ins_abc(WOP_RET, 4, 0, 0);
wb_method(b, kspawn, WOB_NONE, 0, 9, code, n, NULL, 0, NULL, 0);
}
{ /* waiter (method 1, argc 1: r0 = id): wait 3000, RET0 */
uint32_t code[4];
code[0] = wo_ins_abx(WOP_LOADK, 1, (uint16_t)k3000);
code[1] = wo_ins_abc(WOP_BUILTIN, 2, 0, WO_B_PROC_WAIT_DL);
code[2] = wo_ins_abc(WOP_RET0, 0, 0, 0);
wb_method(b, kw, WOB_NONE, 1, 3, code, 3, NULL, 0, NULL, 0);
}
{ /* second (method 2, argc 1): sleep 200 so the fiber parks first,
* then wait 100 — must trap "already has a waiter" */
uint32_t code[6];
code[0] = wo_ins_abx(WOP_LOADK, 1, (uint16_t)k200);
code[1] = wo_ins_abc(WOP_BUILTIN, 2, 1, WO_B_TIME_SLEEP);
code[2] = wo_ins_abx(WOP_LOADK, 1, (uint16_t)k100);
code[3] = wo_ins_abc(WOP_BUILTIN, 2, 0, WO_B_PROC_WAIT_DL);
code[4] = wo_ins_abc(WOP_RET0, 0, 0, 0);
wb_method(b, ksec, WOB_NONE, 1, 3, code, 5, NULL, 0, NULL, 0);
}
return wb_finish(b, len);
}
static void test_stream_one_waiter(void) {
size_t len;
uint8_t *img = double_wait_module(&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 id = 0;
wo_err err;
memset(&err, 0, sizeof err);
T_EQ(wo_vm_call(&VM, 0, NULL, 0, &id, &err), 0); /* spawner */
uint64_t warg[1] = {id};
T_CHECK(wo_vm_spawn_fiber(&VM, 1, warg, 1) != NULL); /* waiter */
uint64_t ret = 0;
T_EQ(wo_vm_call(&VM, 2, warg, 1, &ret, &err), -1); /* second */
T_EQ(err.code, WO_T_IO);
T_CHECK(strstr(err.msg, "waiter") != NULL);
wo_vm_destroy(&VM);
int st;
T_EQ(waitpid(-1, &st, WNOHANG), -1);
T_EQ(errno, ECHILD);
wo_module_free(&mod);
free(img);
}
/* churn: 200 spawn/wait/close rounds leave the fd table flat */
static uint8_t *stream_churn_module(size_t *len) {
wb_t *b = wb_new();
uint32_t kchild = wb_const_text(b, "Child");
uint32_t km = wb_const_text(b, "main");
uint32_t kcmd = wb_const_text(b, "true");
uint8_t kinds[4] = {WO_K_SCALAR, WO_K_SCALAR, WO_K_SCALAR, WO_K_SCALAR};
uint32_t cls = wb_class(b, kchild, 0, kinds, 4);
uint32_t kcls = wb_const_int(b, (int64_t)cls);
uint32_t k5000 = wb_const_int(b, 5000);
uint32_t kz = wb_const_int(b, 0);
uint32_t klim = wb_const_int(b, 200);
uint32_t k1 = wb_const_int(b, 1);
uint32_t code[40];
uint32_t n = 0;
code[n++] = wo_ins_abx(WOP_LOADK, 10, (uint16_t)kz);
code[n++] = wo_ins_abx(WOP_LOADK, 11, (uint16_t)klim);
uint32_t loop_pc = n;
code[n++] = wo_ins_abc(WOP_LT, 12, 10, 11);
uint32_t jz_pc = n;
code[n++] = wo_ins_asbx(WOP_JZ, 12, 0); /* patched */
code[n++] = wo_ins_abx(WOP_LOADK, 1, (uint16_t)kcmd);
code[n++] = wo_ins_abc(WOP_BUILTIN, 2, WO_K_TEXT, WO_B_MULTI_NEW);
code[n++] = wo_ins_abx(WOP_LOADK, 3, (uint16_t)kcls);
code[n++] = wo_ins_abc(WOP_BUILTIN, 0, 1, WO_B_PROC_SPAWN);
code[n++] = wo_ins_abc(WOP_DROP, 2, 0, 0);
code[n++] = wo_ins_abc(WOP_GETF, 1, 0, 1);
code[n++] = wo_ins_abc(WOP_BUILTIN, 8, 1, WO_B_NET_CLOSE);
code[n++] = wo_ins_abc(WOP_GETF, 1, 0, 2);
code[n++] = wo_ins_abc(WOP_BUILTIN, 8, 1, WO_B_NET_CLOSE);
code[n++] = wo_ins_abc(WOP_GETF, 1, 0, 3);
code[n++] = wo_ins_abc(WOP_BUILTIN, 8, 1, WO_B_NET_CLOSE);
code[n++] = wo_ins_abc(WOP_GETF, 4, 0, 0);
code[n++] = wo_ins_abx(WOP_LOADK, 5, (uint16_t)k5000);
code[n++] = wo_ins_abc(WOP_BUILTIN, 7, 4, WO_B_PROC_WAIT_DL);
code[n++] = wo_ins_abc(WOP_DROP, 0, 0, 0);
code[n++] = wo_ins_abx(WOP_LOADK, 12, (uint16_t)k1);
code[n++] = wo_ins_abc(WOP_ADD, 10, 10, 12);
uint32_t jmp_pc = n;
code[n++] = wo_ins_asbx(WOP_JMP, 0, (int)loop_pc - ((int)jmp_pc + 1));
uint32_t exit_pc = n;
code[n++] = wo_ins_abc(WOP_RET0, 0, 0, 0);
code[jz_pc] = wo_ins_asbx(WOP_JZ, 12, (int)exit_pc - ((int)jz_pc + 1));
wb_method(b, km, WOB_NONE, 0, 13, code, n, NULL, 0, NULL, 0);
return wb_finish(b, len);
}
static void test_stream_churn(void) {
size_t len;
uint8_t *img = stream_churn_module(&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);
int fds0 = fd_count();
uint64_t ret = 0;
wo_err err;
memset(&err, 0, sizeof err);
T_EQ(wo_vm_call(&VM, 0, NULL, 0, &ret, &err), 0);
T_EQ(fd_count(), fds0);
T_EQ(VM.nchildren, 0u);
wo_vm_destroy(&VM);
wo_module_free(&mod);
free(img);
}
static void on_alarm(int sig) { (void)sig; } /* interrupt, don't die */
static int64_t mono_ms(void) {
@ -625,6 +927,10 @@ int main(void) {
test_ceiling_fails_closed();
test_unwind_reaps_child();
test_thousand_spawns_fd_flat();
test_stop_kills_child();
test_stream_echo();
test_stream_wait();
test_stream_one_waiter();
test_stream_churn();
test_stop_kills_child(); /* last: it latches the stop flag */
return t_report("test_proc");
}