diff --git a/compiler/src/emit.ml b/compiler/src/emit.ml index 271f71f..32e9a45 100644 --- a/compiler/src/emit.ml +++ b/compiler/src/emit.ml @@ -289,6 +289,7 @@ let b_sha1 = 85 let b_sha256 = 86 let b_hmac_sha256 = 87 let b_call = 88 +let b_monitor = 89 let b_split = 28 let b_split_ws = 29 let b_join = 30 @@ -1100,7 +1101,7 @@ let is_builtin_name (n : string) = "substr"; "trim"; "to_lower"; "char_of"; "parse_int"; "split"; "split_ws"; "join"; "slice"; "pop"; "shift"; "sort"; "reverse"; "remove"; "key_at"; "val_at"; (* the concurrency arc *) - "send"; "call"; + "send"; "call"; "monitor"; (* iteration 19: Float bridges and Bytes surface *) "float"; "trunc"; "parse_float"; "float_to_text"; "float_cmp"; "bytes_len"; "bytes_at"; "bytes_slice"; "bytes_eq"; "bytes_concat"; "base64_encode"; "base64_decode"; @@ -3419,11 +3420,18 @@ and emit_call (p : pctx) (f : fstate) (v : views) ~(dst : int) ?expected (e : As put f (ins_abc op_builtin dst base sm.Types.sm_builtin); (* every stdlib member only READS its arguments, so one that was freshly built here (`net.write(c, head .. resp.body)`) has no - other owner and dies with the call *) + other owner and dies with the call. The ONE exception: + `time.after`'s message (arg 2) MOVES to the runtime — the + timer owns it until delivery (iteration 24 T5). *) + let moves i = + alias = "time" && mname = "after" && i = 2 + in List.iteri (fun i (a : Ast.expr) -> - drop_fresh_owned ~keep:dst p f (base + i) a; - drop_fresh_text ~keep:dst p f (base + i) a) + if not (moves i) then begin + drop_fresh_owned ~keep:dst p f (base + i) a; + drop_fresh_text ~keep:dst p f (base + i) a + end) args end) | Some u -> ( @@ -3700,7 +3708,7 @@ and emit_builtin (p : pctx) (f : fstate) (v : views) ~(dst : int) ?expected (e : dangle the value just read) and the stores, which either copy (Text, handled by copied_container_call) or take ownership (OWNED/GCREF). *) let reader = List.mem name [ "get"; "latest"; "key_at"; "val_at" ] in - (if not (List.mem name [ "push"; "set"; "send"; "call" ]) then + (if not (List.mem name [ "push"; "set"; "send"; "call"; "monitor" ]) then List.iteri (fun i (a : Ast.expr) -> (* a reader's result points into arg0 (the container) — dropping @@ -3732,6 +3740,7 @@ and emit_builtin (p : pctx) (f : fstate) (v : views) ~(dst : int) ?expected (e : match name with | "send" -> fixed b_send (* arc: msg (arg1) moved to the runtime — never dropped here *) | "call" -> fixed b_call (* iteration 24: same move; the SCALAR reply lands in dst *) + | "monitor" -> fixed b_monitor (* T4: notice msg (arg2) moves to the runtime *) | "now" -> fixed b_now | "print" -> fixed b_print | "print_int" -> fixed b_print_int diff --git a/compiler/src/owner.ml b/compiler/src/owner.ml index bfe8f29..3db6180 100644 --- a/compiler/src/owner.ml +++ b/compiler/src/owner.ml @@ -1348,6 +1348,10 @@ and analyze_call (ctx : ctx) (call_e : Ast.expr) (callee : Ast.expr) (args : Ast iteration 24: call(addr, msg) moves its message identically. *) | Ident "send" -> i = 1 && Types.StringMap.find_opt "send" ctx.syms.Types.free_fns = None | Ident "call" -> i = 1 && Types.StringMap.find_opt "call" ctx.syms.Types.free_fns = None + (* T4/T5: the notice / timer message moves to the runtime too *) + | Ident "monitor" -> + i = 2 && Types.StringMap.find_opt "monitor" ctx.syms.Types.free_fns = None + | Field ({ kind = Ident "time"; _ }, "after") -> i = 2 | _ -> false in List.iteri @@ -1365,7 +1369,8 @@ and analyze_call (ctx : ctx) (call_e : Ast.expr) (callee : Ast.expr) (args : Ast transfer ctx p ~what: (match callee.kind with - | Ident "send" | Ident "call" -> + | Ident "send" | Ident "call" | Ident "monitor" + | Field ({ kind = Ident "time"; _ }, "after") -> "cannot be sent — a message moves to the receiver" | _ -> "cannot be stored in a container") then record_move ctx p (MvArg "element")) diff --git a/compiler/src/types.ml b/compiler/src/types.ml index 4afc6cd..a56489a 100644 --- a/compiler/src/types.ml +++ b/compiler/src/types.ml @@ -309,6 +309,8 @@ let stdlib_members : stdlib_member list = m "net" "write_dl" 3 93 (Some (TScalar "Bool")) None; m "net" "listen_unix" 1 94 (Some (TScalar "Int")) None; m "net" "peer" 1 95 (Some (TScalar "Text")) None; + (* iteration 24 T5: one-shot timer — the msg MOVES to the runtime *) + m "time" "after" 3 90 None None; (* proc *) m "proc" "run" 2 56 (Some (TNullable (TScalar proc_record_name))) (Some proc_record_name); (* json — both members are lowered specially (emit.ml): encode needs its @@ -1869,6 +1871,48 @@ let typecheck_program ~file ~(module_of : string -> string) ~message:"`call`'s first argument must be an `actor M` address" ()) | None -> ()) | _ -> ()) + | None when name = "monitor" -> + (* iteration 24 T4: monitor(watched, observer, msg) — the + notice msg is typed against the OBSERVER's mailbox + (three-argument form: the caller may be main, which has + no mailbox). msg moves like send's. *) + (if List.length args <> 3 then + Diag.Collector.add collector + (Diag.error ~code:bad_arity_code ~file ~line:e.pos.line ~col:e.pos.col + ~message: + (Printf.sprintf + "`monitor` takes 3 arguments (watched, observer, notice), given %d" + (List.length args)) + ()) + else + match args with + | [ w; o; m ] -> ( + (match confident_typ cenv w with + | Some (TActor _) | None -> () + | Some _ -> + Diag.Collector.add collector + (Diag.error ~code:type_mismatch_code ~file ~line:w.pos.line + ~col:w.pos.col + ~message:"`monitor`'s first argument must be an `actor M` address" ())); + match confident_typ cenv o with + | Some (TActor want) -> ( + match confident_typ cenv m with + | Some (TScalar got) when got <> want -> + Diag.Collector.add collector + (Diag.error ~code:type_mismatch_code ~file ~line:m.pos.line + ~col:m.pos.col + ~message: + (Printf.sprintf + "the observer receives `%s` — the notice is a `%s`" want got) + ()) + | _ -> ()) + | Some _ -> + Diag.Collector.add collector + (Diag.error ~code:type_mismatch_code ~file ~line:o.pos.line + ~col:o.pos.col + ~message:"`monitor`'s second argument must be an `actor M` address" ()) + | None -> ()) + | _ -> ()) | None -> let confident_types = List.map (confident_typ cenv) args in check_builtin_call ~file collector name e.pos args confident_types) diff --git a/docs/plan/oop-vm/08-builtin-surface.md b/docs/plan/oop-vm/08-builtin-surface.md index 4dd458b..ec4f0c3 100644 --- a/docs/plan/oop-vm/08-builtin-surface.md +++ b/docs/plan/oop-vm/08-builtin-surface.md @@ -277,6 +277,8 @@ unset `env.get` are nil. | `sha1(bytes)` | `-> Bytes` | 20-byte digest (id 85, iteration 34) — exists because RFC 6455's Sec-WebSocket-Accept demands SHA-1 | | `sha256(bytes)` | `-> Bytes` | 32-byte digest (id 86, iteration 34) | | `hmac_sha256(key, msg)` | `-> Bytes` | RFC 2104 over SHA-256, both args Bytes (id 87, iteration 34); key > 64 bytes hashed first | +| `monitor(watched, observer, msg)` | — | iteration 24 (id 89): the observer's own M-typed msg (MOVED) is delivered when watched dies (trap-death); already-dead delivers now; a full observer's notice is dropped with a stderr line — no fiber to trap | +| `time.after(ms, addr, msg)` | — | iteration 24 (id 90): one-shot timer — msg (MOVED) arrives as an ordinary send after ms on the arming shard; ms <= 0 delivers now; NO cancel — the generation-counter idiom (run/timer-generation) is the answer | | `call(addr, msg)` | `-> R` | send that WAITS (id 88, iteration 24): the message moves like `send`'s, the caller's fiber parks until the receive's return value arrives. R = the receive's declared return type — every `receive(msg: M)` program-wide must agree on it and it must be a copyable scalar in v1 (WO-E226 otherwise). A dead callee traps WO_T_ACTOR, immediately or mid-call — a `call` never hangs | | `env.get(name)` | `-> ?Text` | unset is nil | | `env.stopping()` | `-> Bool` | SIGTERM/SIGINT latch, handlers installed on first use | diff --git a/runtime/src/builtin.c b/runtime/src/builtin.c index 095367e..ff0861f 100644 --- a/runtime/src/builtin.c +++ b/runtime/src/builtin.c @@ -200,6 +200,18 @@ int wo_builtin(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { } case WO_B_CALL: /* iteration 24: park/reply protocol lives in vm.c */ return wo_vm_actor_call(vm, R, ins, msg); + case WO_B_MONITOR: { + int rc = wo_vm_actor_monitor(vm, R[B], R[B + 1], R[B + 2], msg); + if (rc) return rc; + R[A] = 0; + return 0; + } + case WO_B_TIME_AFTER: { + int rc = wo_vm_timer_after(vm, (int64_t)R[B], R[B + 1], R[B + 2], msg); + if (rc) return rc; + R[A] = 0; + return 0; + } case WO_B_NOW: { /* wall-clock milliseconds */ struct timespec ts; clock_gettime(CLOCK_REALTIME, &ts); diff --git a/runtime/src/park.c b/runtime/src/park.c index 6fbe11a..599a4b1 100644 --- a/runtime/src/park.c +++ b/runtime/src/park.c @@ -296,6 +296,8 @@ static void tick_arm_uring(wo_vm *vm, int64_t now) { for (wo_fiber *fb = vm->parked; fb; fb = fb->pnext) if (fb->state == WO_FIB_PARKED && fb->park_fd >= 0 && fb->park_deadline > 0) if (next == 0 || fb->park_deadline < next) next = fb->park_deadline; + int64_t tn = wo_vm_timers_next(vm); /* iteration 24 T5: armed timers */ + if (tn > 0 && (next == 0 || tn < next)) next = tn; if (next == 0) return; if (vm->tick_armed && vm->tick_at <= next) return; int64_t rel = next - now; @@ -371,7 +373,9 @@ int wo_io_wait(wo_vm *vm) { head++; } __atomic_store_n(r.cq_head, head, __ATOMIC_RELEASE); - if (deadline_sweep_uring(vm, now_ms()) && woke != 2) woke = 1; + int64_t swnow = now_ms(); + if (wo_vm_timers_fire(vm, swnow) && woke != 2) woke = 1; + if (deadline_sweep_uring(vm, swnow) && woke != 2) woke = 1; if (woke == 2) return 1; /* adopt-needed */ if (woke) return 0; continue; @@ -389,6 +393,14 @@ int wo_io_wait(wo_vm *vm) { } int timeout = -1; int64_t now = now_ms(); + { + int64_t tn = wo_vm_timers_next(vm); /* iteration 24 T5 */ + if (tn > 0) { + int64_t rel = tn - now; + if (rel < 0) rel = 0; + timeout = (int)rel; + } + } for (wo_fiber *fb = vm->parked; fb; fb = fb->pnext) if (fb->park_fd == -1 || (fb->park_fd >= 0 && fb->park_deadline > 0)) { @@ -417,6 +429,7 @@ int wo_io_wait(wo_vm *vm) { } } now = now_ms(); + if (wo_vm_timers_fire(vm, now)) woke = 1; wo_fiber *fb = vm->parked; while (fb) { wo_fiber *nx = fb->pnext; diff --git a/runtime/src/vm.c b/runtime/src/vm.c index 1ca9271..344ede9 100644 --- a/runtime/src/vm.c +++ b/runtime/src/vm.c @@ -78,6 +78,9 @@ static int actor_push(wo_actor *a, wo_msg m); static void call_reply_to(wo_vm *vm, wo_fiber *caller, uint32_t caller_shard, uint64_t reply, int status); static void actor_drop_payload(wo_vm *vm, uint64_t payload); +static void monitors_fire(wo_vm *vm, wo_actor *a); +static void runtime_notify(wo_vm *vm, wo_actor *target, uint64_t msg_val, + const char *what); /* the owning thread drains its inbox: adopt actors, deliver sends, * execute home-routed frees. Returns how many envelopes were handled. */ @@ -133,6 +136,25 @@ static int wo_vm_adopt(wo_vm *vm) { } break; } + case 7: { /* iteration 24 T4: a cross-shard monitor registration — + WE are the watched actor's home. Dead already = the + notice fires now; else it joins the list. */ + wo_actor *ob = (wo_actor *)(uintptr_t)e->from_fiber; + if (e->actor->dead) { + runtime_notify(vm, ob, e->payload, "death notice"); + break; + } + wo_monitor *mn = calloc(1, sizeof *mn); + if (!mn) { + actor_drop_payload(vm, e->payload); + break; + } + mn->observer = ob; + mn->msg = e->payload; + mn->next = e->actor->monitors; + e->actor->monitors = mn; + break; + } case 6: /* iteration 24: a call reply landing on the caller's shard — fill the slot and wake the parked fiber; the re-executed builtin consumes it (status != 0 makes it trap). */ @@ -566,10 +588,25 @@ void wo_vm_destroy(wo_vm *vm) { uint64_t m = a->msgs[(a->mhead + i) % a->mcap].payload; if (m) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)m); } + wo_monitor *mo = a->monitors; + while (mo) { /* undelivered notices are the runtime's to drop */ + wo_monitor *mnx = mo->next; + if (mo->msg) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)mo->msg); + free(mo); + mo = mnx; + } free(a->msgs); free(a); a = nx; } + wo_timer *tt = vm->timers; + vm->timers = NULL; + while (tt) { /* unfired timers likewise */ + wo_timer *tnx = tt->next; + if (tt->msg) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)tt->msg); + free(tt); + tt = tnx; + } vm->actors = NULL; wo_io_destroy(vm); wo_rt_destroy(&vm->rt); @@ -760,6 +797,7 @@ static void actor_die(wo_vm *vm, wo_actor *a, wo_fiber *delivery) { a->instance = 0; } a->active = NULL; + monitors_fire(vm, a); } /* Mailbox nonempty, no delivery fiber: start one on the next message. @@ -877,6 +915,57 @@ int wo_vm_actor_send(wo_vm *vm, uint64_t addr, uint64_t msg_val, const char **ms return 0; } +/* iteration 24 T4/T5: a RUNTIME-sourced delivery (death notice, timer). + * No fiber to trap: a full or dead target drops the message with a + * stderr line (spec'd disclosure), never silently. Runs on any thread — + * cross-shard targets ride the ordinary kind-0 envelope. */ +static void runtime_notify(wo_vm *vm, wo_actor *target, uint64_t msg_val, + const char *what) { + if (!target || !msg_val) return; + if (target->dead) { + actor_drop_payload(vm, msg_val); + return; /* send-to-dead: silent by contract */ + } + if (wo_mbox_reserve(target) != 0) { + fprintf(stderr, "wovm: %s dropped — the observer's mailbox is full\n", what); + actor_drop_payload(vm, msg_val); + return; + } + if (target->home != vm->shard_id) { + wo_envelope *e = calloc(1, sizeof *e); + if (!e) { + wo_mbox_release(target); + actor_drop_payload(vm, msg_val); + return; + } + e->kind = 0; + e->actor = target; + e->payload = msg_val; + inbox_push_to(target->home, e); + return; + } + wo_msg m0 = { msg_val, NULL, 0 }; + if (actor_push(target, m0) != 0) { + wo_mbox_release(target); + actor_drop_payload(vm, msg_val); + return; + } + if (!target->active) (void)actor_activate(vm, target); +} + +/* iteration 24 T4: the death walk — every registered observer gets its + * chosen notice, then the list is gone (an actor dies once). */ +static void monitors_fire(wo_vm *vm, wo_actor *a) { + wo_monitor *m = a->monitors; + a->monitors = NULL; + while (m) { + wo_monitor *nx = m->next; + runtime_notify(vm, m->observer, m->msg, "death notice"); + free(m); + m = nx; + } +} + /* iteration 24: call — send that waits. First entry enqueues with the * caller attached and parks (WO_PARK_INBOX, the DB-RPC park); the resume * RE-EXECUTES this builtin and consumes the scalar reply. No hangs, ever: @@ -946,6 +1035,105 @@ int wo_vm_actor_call(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { return WO_SYS_PARKED; } +int wo_vm_actor_monitor(wo_vm *vm, uint64_t watched, uint64_t observer, + uint64_t msg_val, const char **msg) { + wo_actor *w = (wo_actor *)(uintptr_t)watched; + wo_actor *o = (wo_actor *)(uintptr_t)observer; + if (!w || !o) { + *msg = "monitor: nil actor address"; + return WO_T_BOUNDS; + } + if (!msg_val) { + *msg = "monitor: nil notice message"; + return WO_T_BOUNDS; + } + /* the registration belongs to the WATCHED actor's home thread */ + if (w->home != vm->shard_id) { + wo_envelope *e = calloc(1, sizeof *e); + if (!e) { + actor_drop_payload(vm, msg_val); + *msg = "out of memory"; + return WO_T_OOM; + } + e->kind = 7; + e->actor = w; + e->payload = msg_val; + e->from_fiber = (wo_fiber *)o; /* reused slot: the observer */ + inbox_push_to(w->home, e); + return 0; + } + if (w->dead) { /* monitoring the dead: the notice fires NOW */ + runtime_notify(vm, o, msg_val, "death notice"); + return 0; + } + wo_monitor *m = calloc(1, sizeof *m); + if (!m) { + actor_drop_payload(vm, msg_val); + *msg = "out of memory"; + return WO_T_OOM; + } + m->observer = o; + m->msg = msg_val; + m->next = w->monitors; + w->monitors = m; + return 0; +} + +int wo_vm_timer_after(wo_vm *vm, int64_t ms, uint64_t addr, uint64_t msg_val, + const char **msg) { + wo_actor *a = (wo_actor *)(uintptr_t)addr; + if (!a) { + *msg = "time.after: nil actor address"; + return WO_T_BOUNDS; + } + if (!msg_val) { + *msg = "time.after: nil message"; + return WO_T_BOUNDS; + } + if (ms <= 0) { /* no wait to arm: deliver now */ + runtime_notify(vm, a, msg_val, "timer message"); + return 0; + } + wo_timer *t = calloc(1, sizeof *t); + if (!t) { + actor_drop_payload(vm, msg_val); + *msg = "out of memory"; + return WO_T_OOM; + } + struct timespec now; + clock_gettime(CLOCK_REALTIME, &now); + t->at = (int64_t)now.tv_sec * 1000 + now.tv_nsec / 1000000 + ms; + t->target = a; + t->msg = msg_val; + t->next = vm->timers; + vm->timers = t; + return 0; +} + +int wo_vm_timers_fire(wo_vm *vm, int64_t now) { + int fired = 0; + wo_timer **pp = &vm->timers; + while (*pp) { + wo_timer *t = *pp; + if (t->at <= now) { + *pp = t->next; + runtime_notify(vm, t->target, t->msg, "timer message"); + free(t); + fired++; + } else { + pp = &t->next; + } + } + return fired; +} + +int64_t wo_vm_timers_next(wo_vm *vm) { + int64_t next = 0; + for (wo_timer *t = vm->timers; t; t = t->next) + if (next == 0 || t->at < next) next = t->at; + return next; +} + /* The drop-table entry governing instruction [pc]: the last one recorded * at or before it. NULL = nothing live there. */ static const wo_dropent *vm_dropent(const wo_methodrec *me, uint32_t pc) { diff --git a/runtime/src/vm.h b/runtime/src/vm.h index 78675cb..94ea82b 100644 --- a/runtime/src/vm.h +++ b/runtime/src/vm.h @@ -120,6 +120,27 @@ typedef struct wo_msg { * guarantee). Death (iteration 24): a receive trapping uncaught marks * the actor dead — sends to it drop silently, calls trap, queued * callers are error-unparked; the state and mailbox are released. */ +/* iteration 24 T4: one death-notice registration. The runtime owns the + * moved-in notice message until delivery (or drops it if the observer is + * unreachable). The list lives on the WATCHED actor, owned by its home + * thread. */ +typedef struct wo_monitor { + struct wo_actor *observer; + uint64_t msg; + struct wo_monitor *next; +} wo_monitor; + +/* iteration 24 T5: one armed one-shot timer — fires as an ordinary + * runtime send of the moved message when `at` passes. The list lives on + * the ARMING fiber's shard and is scanned by the same deadline machinery + * that serves fd-park deadlines. */ +typedef struct wo_timer { + int64_t at; /* wall ms */ + struct wo_actor *target; + uint64_t msg; + struct wo_timer *next; +} wo_timer; + typedef struct wo_actor { uint64_t instance; /* the moved-in state object (runtime-owned) */ uint32_t method; /* receive's method index (self + msg = 2 args) */ @@ -134,6 +155,7 @@ typedef struct wo_actor { * overshoot by at most the number of in-flight sends — disclosed. */ uint32_t pending; wo_fiber *active; /* the delivery fiber, NULL when idle */ + wo_monitor *monitors; /* iteration 24 T4: who wants the death notice */ struct wo_actor *next_all; /* the vm's all-actors list */ } wo_actor; @@ -176,6 +198,9 @@ typedef struct wo_vm { * freed memory is the UAF this prevents. Steady-state pool size = the * peak live fiber count; the pool dies with the vm. */ wo_fiber *fib_pool; + /* iteration 24 T5: this shard's armed timers (unsorted list — the + * deadline scan is already linear; a wheel is measured-later work) */ + wo_timer *timers; /* 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 @@ -206,6 +231,20 @@ int wo_vm_actor_send(wo_vm *vm, uint64_t addr, uint64_t msg_val, const char **ms * caller attached and parks (WO_SYS_PARKED); the re-execution consumes the * scalar reply into R[A] (vm.c owns the protocol, builtin.c dispatches). */ int wo_vm_actor_call(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg); +/* iteration 24 T4: register a death notice — monitor(watched, observer, + * msg). The msg MOVES to the runtime; an already-dead watched actor + * delivers it immediately. */ +int wo_vm_actor_monitor(wo_vm *vm, uint64_t watched, uint64_t observer, + uint64_t msg_val, const char **msg); +/* iteration 24 T5: arm a one-shot timer on THIS shard — time.after(ms, + * addr, msg). ms <= 0 delivers now. */ +int wo_vm_timer_after(wo_vm *vm, int64_t ms, uint64_t addr, uint64_t msg_val, + const char **msg); +/* iteration 24 T5: fire every timer at or past `now` (park.c's deadline + * machinery calls this beside the fd-park sweep). Returns fired count. */ +int wo_vm_timers_fire(wo_vm *vm, int64_t now); +/* The nearest armed timer's deadline, 0 = none (park.c's tick/timeout). */ +int64_t wo_vm_timers_next(wo_vm *vm); /* ---- the shard engine (arc stage 2) ------------------------------------ * One pinned thread per shard, each a full wo_vm (own arena, GC, I/O @@ -231,7 +270,11 @@ typedef struct wo_envelope { * from_shard/from_fiber = the parked caller), * 6 = CALL_REPLY (payload = the SCALAR reply, from_fiber = * the caller to unpark; status 0 = ok, WO_T_ACTOR = - * the callee was/went dead — the caller traps) */ + * the callee was/went dead — the caller traps), + * 7 = MONITOR (iteration 24 T4: actor = the WATCHED one, + * from_fiber REUSED as the observer wo_actor*, payload = + * the moved notice — registered on the watched actor's + * home thread; already-dead delivers the notice now) */ struct wo_actor *actor; uint64_t payload; uint32_t from_shard; diff --git a/runtime/src/wob.h b/runtime/src/wob.h index 506d59a..d15e0b9 100644 --- a/runtime/src/wob.h +++ b/runtime/src/wob.h @@ -468,8 +468,17 @@ enum { * return value arrives. R is a SCALAR (v1, * compiler-enforced WO-E226). Dead callee = * WO_T_ACTOR, immediately or mid-call. */ - /* ids 89 (monitor) and 90 (time.after) are RESERVED for the rest of - * the lifecycle slice — do not reuse. */ + WO_B_MONITOR = 89, /* (watched, observer, msg) -> (): the + * observer's own M-typed msg is delivered + * when watched dies (trap-death); already + * dead delivers NOW; msg MOVES. A full + * observer's notice is dropped with a + * stderr line (no fiber to trap). */ + WO_B_TIME_AFTER = 90, /* (ms, addr, msg) -> (): one-shot timer — + * msg (MOVED) arrives as an ordinary send + * after ms; no cancel (the generation- + * counter idiom is the documented answer); + * ms <= 0 delivers now. */ /* ---- iteration 35: net seams (sysio.c). Deadlines are per-CALL (no * hidden fd state); a timeout is an EXPECTED outcome, so it answers * nil/false, never a trap. ms <= 0 = no deadline (the old behavior, diff --git a/tests/corpus/run/monitor-death/fixture.out b/tests/corpus/run/monitor-death/fixture.out new file mode 100644 index 0000000..0c676c2 --- /dev/null +++ b/tests/corpus/run/monitor-death/fixture.out @@ -0,0 +1,3 @@ +died: boom +died: late +done diff --git a/tests/corpus/run/monitor-death/fixture.wo b/tests/corpus/run/monitor-death/fixture.wo new file mode 100644 index 0000000..025196b --- /dev/null +++ b/tests/corpus/run/monitor-death/fixture.wo @@ -0,0 +1,36 @@ +use time + +-- iteration 24 T4: actor death is OBSERVABLE. The observer names its own +-- notice message; the watched actor trapping uncaught (the runtime's +-- stderr line) delivers it. Monitoring an ALREADY dead actor fires +-- immediately. WO_SHARDS=1 (the runner) keeps the order deterministic. +class Note { + who: Text +} + +class Watch { + pad: Int + fn receive(msg: Note) { + print("died: ${msg.who}"); + } +} + +class Boom { + pad: Int + fn receive(msg: Note) { + let z = len(msg.who) - len(msg.who); + let q = 1 / z; + } +} + +fn main() -> Int { + let obs: actor Note = spawn Watch { pad: 0 }; + let b: actor Note = spawn Boom { pad: 0 }; + monitor(b, obs, Note { who: "boom" }); + send(b, Note { who: "x" }); + time.sleep(100); + monitor(b, obs, Note { who: "late" }); + time.sleep(100); + print("done"); + return 0; +} diff --git a/tests/corpus/run/timer-delivery/fixture.out b/tests/corpus/run/timer-delivery/fixture.out new file mode 100644 index 0000000..d88608c --- /dev/null +++ b/tests/corpus/run/timer-delivery/fixture.out @@ -0,0 +1,3 @@ +tick: now +tick: armed +done diff --git a/tests/corpus/run/timer-delivery/fixture.wo b/tests/corpus/run/timer-delivery/fixture.wo new file mode 100644 index 0000000..ead06e5 --- /dev/null +++ b/tests/corpus/run/timer-delivery/fixture.wo @@ -0,0 +1,23 @@ +use time + +-- iteration 24 T5: a timer is a MESSAGE. time.after arms a one-shot on +-- this shard; the target receives it like any send. ms <= 0 delivers now. +class Tick { + tag: Text +} + +class Sink { + pad: Int + fn receive(msg: Tick) { + print("tick: ${msg.tag}"); + } +} + +fn main() -> Int { + let a: actor Tick = spawn Sink { pad: 0 }; + time.after(30, a, Tick { tag: "armed" }); + time.after(0, a, Tick { tag: "now" }); + time.sleep(150); + print("done"); + return 0; +} diff --git a/tests/corpus/run/timer-generation/fixture.out b/tests/corpus/run/timer-generation/fixture.out new file mode 100644 index 0000000..5b60035 --- /dev/null +++ b/tests/corpus/run/timer-generation/fixture.out @@ -0,0 +1,3 @@ +stale gen 1 ignored +fired gen 2 +done diff --git a/tests/corpus/run/timer-generation/fixture.wo b/tests/corpus/run/timer-generation/fixture.wo new file mode 100644 index 0000000..0686e37 --- /dev/null +++ b/tests/corpus/run/timer-generation/fixture.wo @@ -0,0 +1,28 @@ +use time + +-- iteration 24 T5: the CANCEL idiom — no cancel builtin, a generation +-- counter instead. The actor bumps its generation; a stale timer's +-- message names the old one and is recognized and ignored on arrival. +class Timer { + gen: Int +} + +class Gate { + gen: Int + fn receive(msg: Timer) { + if msg.gen == self.gen { + print("fired gen ${msg.gen}"); + } else { + print("stale gen ${msg.gen} ignored"); + } + } +} + +fn main() -> Int { + let g: actor Timer = spawn Gate { gen: 2 }; + time.after(30, g, Timer { gen: 1 }); + time.after(60, g, Timer { gen: 2 }); + time.sleep(200); + print("done"); + return 0; +}