feat(rt2): signal.on — latched signals become Signal records for actors
- mechanics amendment to the spec (recorded at close-out): no signalfd —
the stop-latch pattern generalized. An async-signal-safe handler
latches the number, bumps a sequence and pokes shard 0's wake eventfd;
wo_io_wait's loop head drains latches into fresh Signal{sig} records
delivered via runtime_notify (exported as wo_actor_notify)
- payloads must be heap objects (vm.c drops them unconditionally) — the
Signal record exists exactly for that; class id rides the call as the
appended record operand (sm_record drives it even with no return)
- offerable: WINCH/CHLD/HUP/USR1/USR2; SIGTERM/SIGINT refused naming the
stop latch; shard-0-only registration; coalescing disclosed
- stdlib_modules gains `signal` (and `term`, next task)
- test_term: a real child kills the test process with USR1; the actor's
multi holds one coalesced delivery; refusal leg verbatim. 14/0
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
(cherry picked from commit 14e03a6a4a343b97c9b47fab1a5e3c4bb69d8201)
This commit is contained in:
parent
0c7d0530e9
commit
c55d6e1d33
6 changed files with 251 additions and 1 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
142
runtime/test/test_term.c
Normal file
142
runtime/test/test_term.c
Normal file
|
|
@ -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 <errno.h>
|
||||
#include <signal.h>
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <sys/wait.h>
|
||||
#include <unistd.h>
|
||||
|
||||
#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");
|
||||
}
|
||||
Loading…
Reference in a new issue