diff --git a/compiler/src/types.ml b/compiler/src/types.ml index 6e28cc7..869c727 100644 --- a/compiler/src/types.ml +++ b/compiler/src/types.ml @@ -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 diff --git a/runtime/src/builtin.c b/runtime/src/builtin.c index 474af5b..1afaa23 100644 --- a/runtime/src/builtin.c +++ b/runtime/src/builtin.c @@ -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); diff --git a/runtime/src/loader.c b/runtime/src/loader.c index 0ad36c3..a28fc76 100644 --- a/runtime/src/loader.c +++ b/runtime/src/loader.c @@ -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, diff --git a/runtime/src/sysio.c b/runtime/src/sysio.c index fc92d4e..5ddd4b3 100644 --- a/runtime/src/sysio.c +++ b/runtime/src/sysio.c @@ -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; diff --git a/runtime/src/vm.c b/runtime/src/vm.c index 3b45882..d48194e 100644 --- a/runtime/src/vm.c +++ b/runtime/src/vm.c @@ -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; diff --git a/runtime/src/vm.h b/runtime/src/vm.h index eb29fa7..53f338a 100644 --- a/runtime/src/vm.h +++ b/runtime/src/vm.h @@ -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 diff --git a/runtime/src/wob.h b/runtime/src/wob.h index 406d26a..422e662 100644 --- a/runtime/src/wob.h +++ b/runtime/src/wob.h @@ -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 diff --git a/runtime/test/test_proc.c b/runtime/test/test_proc.c index 3aa3164..25e8fd2 100644 --- a/runtime/test/test_proc.c +++ b/runtime/test/test_proc.c @@ -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"); }