diff --git a/compiler/src/types.ml b/compiler/src/types.ml index a67f131..bb6a52a 100644 --- a/compiler/src/types.ml +++ b/compiler/src/types.ml @@ -203,7 +203,7 @@ let numeric_world (t : string) : [ `Int | `Float | `Other ] = single-segment names (`check_use_edges` below treats any one-segment `use` path whose name is in this list as stdlib, unconditionally, never as a project directory search). *) -let stdlib_modules = [ "fs"; "proc"; "net"; "time"; "json"; "env" ] +let stdlib_modules = [ "fs"; "proc"; "net"; "time"; "json"; "env"; "signal"; "term" ] let is_stdlib_module (name : string) : bool = List.mem name stdlib_modules @@ -344,6 +344,10 @@ let stdlib_members : stdlib_member list = resize refuses by name on a pipe child *) m "proc" "spawn_pty" 4 100 (Some (TNullable (TScalar child_record_name))) (Some child_record_name); m "proc" "resize" 3 101 None None; + (* runtime-v2 3: standing subscription; each arrival delivers a fresh + Signal {sig} record to the actor. SIGTERM/SIGINT refused (the stop + latch). Coalescing disclosed. *) + m "signal" "on" 2 102 None (Some signal_record_name); (* 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/park.c b/runtime/src/park.c index af4345a..8574a22 100644 --- a/runtime/src/park.c +++ b/runtime/src/park.c @@ -325,6 +325,10 @@ static void efd_drain(wo_vm *vm) { int wo_io_wait(wo_vm *vm) { for (;;) { + /* runtime-v2 3: latched signals become Signal-record sends. The + * handler's EINTR (or its wake-eventfd write) lands the plane + * here, so a parked shard delivers promptly. */ + wo_vm_signals_drain(vm); if (wo_sys_stop_pending()) { /* iteration 24 (the drain): a STOP does not kill parked fibers * from the outside — it WAKES them all, and each blocking diff --git a/runtime/src/sysio.c b/runtime/src/sysio.c index 6bff5e7..7cf4184 100644 --- a/runtime/src/sysio.c +++ b/runtime/src/sysio.c @@ -262,6 +262,44 @@ static int proc_argv(uint64_t vcmd, uint64_t vargv, char *path, size_t pathcap, return 0; } +/* ---- runtime-v2 3: signals as events --------------------------------- + * The stop-latch pattern generalized: an async-signal-safe handler + * latches the number and pokes shard 0's wake eventfd; the drain (every + * wo_io_wait pass) turns latches into fresh Signal{sig} records + * delivered as ordinary sends. Kernel-style coalescing is disclosed: + * N arrivals between drains deliver once. */ +static volatile sig_atomic_t sig_pending[32]; +static volatile sig_atomic_t sig_seq; +static int sig_wake_efd = -1; + +static void on_subscribed_signal(int sig) { + if (sig > 0 && sig < 32) sig_pending[sig] = 1; + sig_seq = sig_seq + 1; + if (sig_wake_efd > 0) { + uint64_t one = 1; + ssize_t r = write(sig_wake_efd, &one, sizeof one); + (void)r; + } +} + +void wo_vm_signals_drain(wo_vm *vm) { + if (vm->shard_id != 0 || vm->nsigsubs == 0) return; + if (vm->sig_seen == (uint32_t)sig_seq) return; + vm->sig_seen = (uint32_t)sig_seq; + for (int s = 1; s < 32; s++) { + if (!sig_pending[s]) continue; + sig_pending[s] = 0; + for (uint32_t i = 0; i < vm->nsigsubs; i++) { + if (vm->sigsubs[i].sig != s) continue; + wo_hdr *o = wo_obj_new(&vm->rt, vm->sigsubs[i].cls); + if (!o) continue; /* oom: this delivery is lost, latch cleared */ + wo_fields(o)[0] = (uint64_t)s; + wo_actor_notify(vm, vm->sigsubs[i].target, (uint64_t)(uintptr_t)o, + "signal message"); + } + } +} + /* 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++) @@ -1388,6 +1426,43 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { R[A] = 0; return 0; } + /* ---- runtime-v2 3: signal.on -------------------------------------- */ + case WO_B_SIGNAL_ON: { /* (sig, addr, cls) -> 0: standing subscription */ + int sig = (int)(int64_t)R[B]; + if (sig == SIGTERM || sig == SIGINT) { + *msg = "SIGTERM/SIGINT belong to the stop latch, not signal.on"; + return WO_T_IO; + } + if (sig != SIGWINCH && sig != SIGCHLD && sig != SIGHUP && + sig != SIGUSR1 && sig != SIGUSR2) { + *msg = "signal.on offers SIGWINCH/SIGCHLD/SIGHUP/SIGUSR1/SIGUSR2"; + return WO_T_IO; + } + if (vm->shard_id != 0) { + *msg = "signal.on registers on shard 0"; + return WO_T_IO; + } + if (!R[B + 1]) { + *msg = "signal.on: nil actor address"; + return WO_T_BOUNDS; + } + if (vm->nsigsubs >= 8) { + *msg = "signal.on: subscription table full (8)"; + return WO_T_IO; + } + vm->sigsubs[vm->nsigsubs].sig = sig; + vm->sigsubs[vm->nsigsubs].cls = (uint32_t)R[B + 2]; + vm->sigsubs[vm->nsigsubs].target = (struct wo_actor *)(uintptr_t)R[B + 1]; + vm->nsigsubs++; + sig_wake_efd = vm->wake_efd; + struct sigaction sa; + memset(&sa, 0, sizeof sa); + sa.sa_handler = on_subscribed_signal; /* no SA_RESTART: waits must + * EINTR so the drain runs */ + sigaction(sig, &sa, NULL); + 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 d48194e..013a55d 100644 --- a/runtime/src/vm.c +++ b/runtime/src/vm.c @@ -1317,6 +1317,13 @@ int wo_vm_timer_after(wo_vm *vm, int64_t ms, uint64_t addr, uint64_t msg_val, return 0; } +/* runtime-v2 3: the signal drain's delivery path (sysio.c cannot see the + * static runtime_notify) */ +void wo_actor_notify(wo_vm *vm, wo_actor *target, uint64_t payload, + const char *what) { + runtime_notify(vm, target, payload, what); +} + int wo_vm_timers_fire(wo_vm *vm, int64_t now) { int fired = 0; wo_timer **pp = &vm->timers; diff --git a/runtime/src/vm.h b/runtime/src/vm.h index 53f338a..05e885c 100644 --- a/runtime/src/vm.h +++ b/runtime/src/vm.h @@ -201,6 +201,13 @@ 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); +/* runtime-v2 3 (sysio.c): turn latched signals into Signal-record sends. + * Cheap when nothing arrived; called from wo_io_wait and the inbox + * drain. wo_actor_notify is vm.c's runtime_notify, exported for it. */ +void wo_vm_signals_drain(struct wo_vm *vm); +void wo_actor_notify(struct wo_vm *vm, struct wo_actor *target, + uint64_t payload, const char *what); + /* iteration 24: the one mailbox cap (default 1024, WO_MAILBOX overrides * at boot — soak tests shrink it to force the fail-fast policy). */ extern uint32_t wo_mailbox_cap; @@ -247,6 +254,17 @@ typedef struct wo_vm { wo_child children[WO_PROC_MAX]; uint32_t nchildren; uint32_t proc_gen; /* runtime-v2 1: claim counter behind child ids */ + /* runtime-v2 3: signal subscriptions (shard 0 only). The handler + * latches sig_pending and bumps a sequence; the drain (called each + * wo_io_wait pass and on inbox adoption) turns latches into fresh + * Signal records delivered as ordinary sends. Coalescing disclosed. */ + struct { + int sig; + uint32_t cls; /* the Signal record's class id */ + struct wo_actor *target; + } sigsubs[8]; + uint32_t nsigsubs; + uint32_t sig_seen; /* 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/test/test_term.c b/runtime/test/test_term.c new file mode 100644 index 0000000..c3a9040 --- /dev/null +++ b/runtime/test/test_term.c @@ -0,0 +1,142 @@ +/* test_term — runtime-v2 3/4/5: signals as events, termios adoption, + * fd passing. Named for the terminal-facing half of the track. */ +#define _POSIX_C_SOURCE 200809L +#include +#include +#include +#include +#include +#include +#include + +#include "cont.h" +#include "gc.h" +#include "loader.h" +#include "t.h" +#include "vm.h" +#include "wob_build.h" + +static wo_vm VM; + +/* ---- runtime-v2 3: signal.on ------------------------------------------- + * classes: 0 Signal {sig}, 1 Collector {m multi}, 2 Proc {code,out,err}. + * methods: 0 receive(self, msg) pushes msg.sig into self.m; + * 1 main(addr): signal.on(SIGUSR1, addr) then a child kills + * the test process with USR1; 2 refuse(addr): SIGTERM. */ +static uint8_t *signal_module(size_t *len) { + wb_t *b = wb_new(); + uint32_t ksig = wb_const_text(b, "Signal"); + uint32_t kcol = wb_const_text(b, "Collector"); + uint32_t kproc = wb_const_text(b, "Proc"); + uint32_t krecv = wb_const_text(b, "receive"); + uint32_t kmain = wb_const_text(b, "main"); + uint32_t kref = wb_const_text(b, "refuse"); + uint32_t kcmd = wb_const_text(b, "sh"); + uint32_t kdc = wb_const_text(b, "-c"); + uint32_t kkill = wb_const_text(b, "kill -USR1 $PPID"); + uint8_t sk[1] = {WO_K_SCALAR}; + uint32_t cls_sig = wb_class(b, ksig, 0, sk, 1); + uint8_t ck[1] = {WO_K_MULTI}; + uint32_t cls_col = wb_class(b, kcol, 0, ck, 1); + (void)cls_col; + uint8_t pk[3] = {WO_K_SCALAR, WO_K_TEXT, WO_K_TEXT}; + uint32_t cls_proc = wb_class(b, kproc, 0, pk, 3); + uint32_t kcls_sig = wb_const_int(b, (int64_t)cls_sig); + uint32_t kcls_proc = wb_const_int(b, (int64_t)cls_proc); + uint32_t kusr1 = wb_const_int(b, SIGUSR1); + uint32_t kterm = wb_const_int(b, SIGTERM); + uint32_t k100 = wb_const_int(b, 100); + { /* receive(self, msg): self.m gets msg.sig */ + uint32_t code[8]; + uint32_t n = 0; + code[n++] = wo_ins_abc(WOP_GETF, 2, 1, 0); /* msg.sig */ + code[n++] = wo_ins_abc(WOP_GETF, 3, 0, 0); /* self.m */ + code[n++] = wo_ins_abc(WOP_MOVE, 4, 2, 0); + code[n++] = wo_ins_abc(WOP_BUILTIN, 5, 3, WO_B_MULTI_PUSH); + code[n++] = wo_ins_abc(WOP_RET0, 0, 0, 0); + wb_method(b, krecv, WOB_NONE, 2, 6, code, n, NULL, 0, NULL, 0); + } + { /* main(addr) */ + uint32_t code[24]; + uint32_t n = 0; + code[n++] = wo_ins_abx(WOP_LOADK, 1, (uint16_t)kusr1); + code[n++] = wo_ins_abc(WOP_MOVE, 2, 0, 0); + code[n++] = wo_ins_abx(WOP_LOADK, 3, (uint16_t)kcls_sig); + code[n++] = wo_ins_abc(WOP_BUILTIN, 4, 1, WO_B_SIGNAL_ON); + /* a child delivers SIGUSR1 to this process */ + 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)kdc); + code[n++] = wo_ins_abc(WOP_BUILTIN, 6, 2, WO_B_MULTI_PUSH); + code[n++] = wo_ins_abx(WOP_LOADK, 3, (uint16_t)kkill); + code[n++] = wo_ins_abc(WOP_BUILTIN, 6, 2, WO_B_MULTI_PUSH); + code[n++] = wo_ins_abx(WOP_LOADK, 3, (uint16_t)kcls_proc); + code[n++] = wo_ins_abc(WOP_BUILTIN, 0, 1, WO_B_PROC_RUN); + code[n++] = wo_ins_abc(WOP_DROP, 2, 0, 0); + code[n++] = wo_ins_abc(WOP_DROP, 0, 0, 0); + /* give the queued delivery a slice */ + code[n++] = wo_ins_abx(WOP_LOADK, 1, (uint16_t)k100); + code[n++] = wo_ins_abc(WOP_BUILTIN, 2, 1, WO_B_TIME_SLEEP); + code[n++] = wo_ins_abc(WOP_RET0, 0, 0, 0); + wb_drop drops[] = {{.pc = 11, .owned = 1u << 2, .gc = 0}}; + wb_method(b, kmain, WOB_NONE, 1, 7, code, n, NULL, 0, drops, 1); + } + { /* refuse(addr): SIGTERM must refuse naming the stop latch */ + uint32_t code[8]; + uint32_t n = 0; + code[n++] = wo_ins_abx(WOP_LOADK, 1, (uint16_t)kterm); + code[n++] = wo_ins_abc(WOP_MOVE, 2, 0, 0); + code[n++] = wo_ins_abx(WOP_LOADK, 3, (uint16_t)kcls_sig); + code[n++] = wo_ins_abc(WOP_BUILTIN, 4, 1, WO_B_SIGNAL_ON); + code[n++] = wo_ins_abc(WOP_RET0, 0, 0, 0); + wb_method(b, kref, WOB_NONE, 1, 5, code, n, NULL, 0, NULL, 0); + } + return wb_finish(b, len); +} + +static void test_signal_on_delivers_record(void) { + size_t len; + uint8_t *img = signal_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); + /* the collector actor: state = Collector{m} */ + wo_multi *m = wo_multi_new(&VM.rt, WO_K_SCALAR); + T_CHECK(m != NULL); + wo_hdr *inst = wo_obj_new(&VM.rt, 1 /* Collector */); + T_CHECK(inst != NULL); + wo_fields(inst)[0] = (uint64_t)(uintptr_t)m; + uint64_t addr = 0; + const char *emsg = NULL; + T_EQ(wo_vm_actor_spawn(&VM, (uint64_t)(uintptr_t)inst, 0, &addr, &emsg), 0); + uint64_t args[1] = {addr}; + uint64_t ret = 0; + wo_err err; + memset(&err, 0, sizeof err); + T_EQ(wo_vm_call(&VM, 1, args, 1, &ret, &err), 0); + T_EQ(m->len, 1u); /* one delivery, coalesced */ + if (m->len == 1) { + uint64_t v = 0; + T_EQ(wo_multi_get(m, 0, &v), 0); + T_EQ(v, (uint64_t)SIGUSR1); + } + /* SIGTERM registration refuses naming the stop latch */ + memset(&err, 0, sizeof err); + T_EQ(wo_vm_call(&VM, 2, args, 1, &ret, &err), -1); + T_EQ(err.code, WO_T_IO); + T_CHECK(strstr(err.msg, "stop latch") != NULL); + wo_vm_destroy(&VM); /* drops the actor state and the multi */ + int st; + T_EQ(waitpid(-1, &st, WNOHANG), -1); + T_EQ(errno, ECHILD); + wo_module_free(&mod); + free(img); + /* leave no handler behind for later suites in this process */ + signal(SIGUSR1, SIG_DFL); +} + +int main(void) { + test_signal_on_delivers_record(); + return t_report("test_term"); +}