feat: monitor + time.after (ids 89/90) — the lifecycle slice completes

- monitor(watched, observer, msg): registration lives on the watched
  actor's home thread (kind-7 envelope cross-shard); actor_die walks
  the list; already-dead fires NOW; the notice msg moves; a full
  observer's notice drops with a stderr line (no fiber to trap)
- time.after(ms, addr, msg): per-shard timer list riding the deadline
  machinery (uring tick min + epoll timeout both include timers;
  fired from the same sweep); ms <= 0 delivers now; NO cancel — the
  generation-counter idiom is pinned by run/timer-generation
- runtime_notify: one runtime-sourced delivery path (notices, timers) —
  reserve-or-drop, cross-shard via kind-0 envelopes
- compiler: monitor typed as a bespoke free fn (notice typed against
  the OBSERVER's mailbox — the three-argument deviation, disclosed);
  time.after as a stdlib row whose msg arg is EXEMPT from the module-
  call fresh-arg drop (it moves — the double-own bug the timer fixture
  caught); owner move slots for both
- corpus: run/monitor-death (trap-death + already-dead notices),
  run/timer-delivery (armed + immediate), run/timer-generation
- teardown drops undelivered notices and unfired timers; battery 13/13

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
shoney.arickathil 2026-08-23 08:56:48 +02:00
parent 9661c08696
commit 4092074201
15 changed files with 431 additions and 10 deletions

View file

@ -289,6 +289,7 @@ let b_sha1 = 85
let b_sha256 = 86 let b_sha256 = 86
let b_hmac_sha256 = 87 let b_hmac_sha256 = 87
let b_call = 88 let b_call = 88
let b_monitor = 89
let b_split = 28 let b_split = 28
let b_split_ws = 29 let b_split_ws = 29
let b_join = 30 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"; "substr"; "trim"; "to_lower"; "char_of"; "parse_int"; "split"; "split_ws"; "join"; "slice";
"pop"; "shift"; "sort"; "reverse"; "remove"; "key_at"; "val_at"; "pop"; "shift"; "sort"; "reverse"; "remove"; "key_at"; "val_at";
(* the concurrency arc *) (* the concurrency arc *)
"send"; "call"; "send"; "call"; "monitor";
(* iteration 19: Float bridges and Bytes surface *) (* iteration 19: Float bridges and Bytes surface *)
"float"; "trunc"; "parse_float"; "float_to_text"; "float_cmp"; "bytes_len"; "bytes_at"; "float"; "trunc"; "parse_float"; "float_to_text"; "float_cmp"; "bytes_len"; "bytes_at";
"bytes_slice"; "bytes_eq"; "bytes_concat"; "base64_encode"; "base64_decode"; "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); put f (ins_abc op_builtin dst base sm.Types.sm_builtin);
(* every stdlib member only READS its arguments, so one that was (* every stdlib member only READS its arguments, so one that was
freshly built here (`net.write(c, head .. resp.body)`) has no 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 List.iteri
(fun i (a : Ast.expr) -> (fun i (a : Ast.expr) ->
drop_fresh_owned ~keep:dst p f (base + i) a; if not (moves i) then begin
drop_fresh_text ~keep:dst p f (base + i) a) drop_fresh_owned ~keep:dst p f (base + i) a;
drop_fresh_text ~keep:dst p f (base + i) a
end)
args args
end) end)
| Some u -> ( | 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, dangle the value just read) and the stores, which either copy (Text,
handled by copied_container_call) or take ownership (OWNED/GCREF). *) handled by copied_container_call) or take ownership (OWNED/GCREF). *)
let reader = List.mem name [ "get"; "latest"; "key_at"; "val_at" ] in 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 List.iteri
(fun i (a : Ast.expr) -> (fun i (a : Ast.expr) ->
(* a reader's result points into arg0 (the container) — dropping (* 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 match name with
| "send" -> fixed b_send (* arc: msg (arg1) moved to the runtime — never dropped here *) | "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 *) | "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 | "now" -> fixed b_now
| "print" -> fixed b_print | "print" -> fixed b_print
| "print_int" -> fixed b_print_int | "print_int" -> fixed b_print_int

View file

@ -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. *) 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 "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 | 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 | _ -> false
in in
List.iteri List.iteri
@ -1365,7 +1369,8 @@ and analyze_call (ctx : ctx) (call_e : Ast.expr) (callee : Ast.expr) (args : Ast
transfer ctx p transfer ctx p
~what: ~what:
(match callee.kind with (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 sent — a message moves to the receiver"
| _ -> "cannot be stored in a container") | _ -> "cannot be stored in a container")
then record_move ctx p (MvArg "element")) then record_move ctx p (MvArg "element"))

View file

@ -309,6 +309,8 @@ let stdlib_members : stdlib_member list =
m "net" "write_dl" 3 93 (Some (TScalar "Bool")) None; m "net" "write_dl" 3 93 (Some (TScalar "Bool")) None;
m "net" "listen_unix" 1 94 (Some (TScalar "Int")) None; m "net" "listen_unix" 1 94 (Some (TScalar "Int")) None;
m "net" "peer" 1 95 (Some (TScalar "Text")) 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 *) (* proc *)
m "proc" "run" 2 56 (Some (TNullable (TScalar proc_record_name))) (Some proc_record_name); 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 (* 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" ()) ~message:"`call`'s first argument must be an `actor M` address" ())
| None -> ()) | 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 -> | None ->
let confident_types = List.map (confident_typ cenv) args in let confident_types = List.map (confident_typ cenv) args in
check_builtin_call ~file collector name e.pos args confident_types) check_builtin_call ~file collector name e.pos args confident_types)

View file

@ -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 | | `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) | | `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 | | `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 | | `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.get(name)` | `-> ?Text` | unset is nil |
| `env.stopping()` | `-> Bool` | SIGTERM/SIGINT latch, handlers installed on first use | | `env.stopping()` | `-> Bool` | SIGTERM/SIGINT latch, handlers installed on first use |

View file

@ -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 */ case WO_B_CALL: /* iteration 24: park/reply protocol lives in vm.c */
return wo_vm_actor_call(vm, R, ins, msg); 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 */ case WO_B_NOW: { /* wall-clock milliseconds */
struct timespec ts; struct timespec ts;
clock_gettime(CLOCK_REALTIME, &ts); clock_gettime(CLOCK_REALTIME, &ts);

View file

@ -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) 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 (fb->state == WO_FIB_PARKED && fb->park_fd >= 0 && fb->park_deadline > 0)
if (next == 0 || fb->park_deadline < next) next = fb->park_deadline; 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 (next == 0) return;
if (vm->tick_armed && vm->tick_at <= next) return; if (vm->tick_armed && vm->tick_at <= next) return;
int64_t rel = next - now; int64_t rel = next - now;
@ -371,7 +373,9 @@ int wo_io_wait(wo_vm *vm) {
head++; head++;
} }
__atomic_store_n(r.cq_head, head, __ATOMIC_RELEASE); __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 == 2) return 1; /* adopt-needed */
if (woke) return 0; if (woke) return 0;
continue; continue;
@ -389,6 +393,14 @@ int wo_io_wait(wo_vm *vm) {
} }
int timeout = -1; int timeout = -1;
int64_t now = now_ms(); 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) for (wo_fiber *fb = vm->parked; fb; fb = fb->pnext)
if (fb->park_fd == -1 if (fb->park_fd == -1
|| (fb->park_fd >= 0 && fb->park_deadline > 0)) { || (fb->park_fd >= 0 && fb->park_deadline > 0)) {
@ -417,6 +429,7 @@ int wo_io_wait(wo_vm *vm) {
} }
} }
now = now_ms(); now = now_ms();
if (wo_vm_timers_fire(vm, now)) woke = 1;
wo_fiber *fb = vm->parked; wo_fiber *fb = vm->parked;
while (fb) { while (fb) {
wo_fiber *nx = fb->pnext; wo_fiber *nx = fb->pnext;

View file

@ -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, static void call_reply_to(wo_vm *vm, wo_fiber *caller, uint32_t caller_shard,
uint64_t reply, int status); uint64_t reply, int status);
static void actor_drop_payload(wo_vm *vm, uint64_t payload); 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, /* the owning thread drains its inbox: adopt actors, deliver sends,
* execute home-routed frees. Returns how many envelopes were handled. */ * execute home-routed frees. Returns how many envelopes were handled. */
@ -133,6 +136,25 @@ static int wo_vm_adopt(wo_vm *vm) {
} }
break; 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 — case 6: /* iteration 24: a call reply landing on the caller's shard —
fill the slot and wake the parked fiber; the re-executed fill the slot and wake the parked fiber; the re-executed
builtin consumes it (status != 0 makes it trap). */ 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; uint64_t m = a->msgs[(a->mhead + i) % a->mcap].payload;
if (m) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)m); 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->msgs);
free(a); free(a);
a = nx; 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; vm->actors = NULL;
wo_io_destroy(vm); wo_io_destroy(vm);
wo_rt_destroy(&vm->rt); 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->instance = 0;
} }
a->active = NULL; a->active = NULL;
monitors_fire(vm, a);
} }
/* Mailbox nonempty, no delivery fiber: start one on the next message. /* 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; 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 /* iteration 24: call — send that waits. First entry enqueues with the
* caller attached and parks (WO_PARK_INBOX, the DB-RPC park); the resume * 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: * 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; 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 /* The drop-table entry governing instruction [pc]: the last one recorded
* at or before it. NULL = nothing live there. */ * at or before it. NULL = nothing live there. */
static const wo_dropent *vm_dropent(const wo_methodrec *me, uint32_t pc) { static const wo_dropent *vm_dropent(const wo_methodrec *me, uint32_t pc) {

View file

@ -120,6 +120,27 @@ typedef struct wo_msg {
* guarantee). Death (iteration 24): a receive trapping uncaught marks * guarantee). Death (iteration 24): a receive trapping uncaught marks
* the actor dead — sends to it drop silently, calls trap, queued * the actor dead — sends to it drop silently, calls trap, queued
* callers are error-unparked; the state and mailbox are released. */ * 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 { typedef struct wo_actor {
uint64_t instance; /* the moved-in state object (runtime-owned) */ uint64_t instance; /* the moved-in state object (runtime-owned) */
uint32_t method; /* receive's method index (self + msg = 2 args) */ 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. */ * overshoot by at most the number of in-flight sends — disclosed. */
uint32_t pending; uint32_t pending;
wo_fiber *active; /* the delivery fiber, NULL when idle */ 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 */ struct wo_actor *next_all; /* the vm's all-actors list */
} wo_actor; } wo_actor;
@ -176,6 +198,9 @@ typedef struct wo_vm {
* freed memory is the UAF this prevents. Steady-state pool size = the * freed memory is the UAF this prevents. Steady-state pool size = the
* peak live fiber count; the pool dies with the vm. */ * peak live fiber count; the pool dies with the vm. */
wo_fiber *fib_pool; 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 /* iteration 35, uring backend: the shard's ONE deadline tick — a
* TIMEOUT op with a sentinel user_data armed for the nearest fd-park * 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 * 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 * 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). */ * 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); 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) ------------------------------------ /* ---- the shard engine (arc stage 2) ------------------------------------
* One pinned thread per shard, each a full wo_vm (own arena, GC, I/O * 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), * from_shard/from_fiber = the parked caller),
* 6 = CALL_REPLY (payload = the SCALAR reply, from_fiber = * 6 = CALL_REPLY (payload = the SCALAR reply, from_fiber =
* the caller to unpark; status 0 = ok, WO_T_ACTOR = * 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; struct wo_actor *actor;
uint64_t payload; uint64_t payload;
uint32_t from_shard; uint32_t from_shard;

View file

@ -468,8 +468,17 @@ enum {
* return value arrives. R is a SCALAR (v1, * return value arrives. R is a SCALAR (v1,
* compiler-enforced WO-E226). Dead callee = * compiler-enforced WO-E226). Dead callee =
* WO_T_ACTOR, immediately or mid-call. */ * WO_T_ACTOR, immediately or mid-call. */
/* ids 89 (monitor) and 90 (time.after) are RESERVED for the rest of WO_B_MONITOR = 89, /* (watched, observer, msg) -> (): the
* the lifecycle slice — do not reuse. */ * 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 /* ---- iteration 35: net seams (sysio.c). Deadlines are per-CALL (no
* hidden fd state); a timeout is an EXPECTED outcome, so it answers * hidden fd state); a timeout is an EXPECTED outcome, so it answers
* nil/false, never a trap. ms <= 0 = no deadline (the old behavior, * nil/false, never a trap. ms <= 0 = no deadline (the old behavior,

View file

@ -0,0 +1,3 @@
died: boom
died: late
done

View file

@ -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;
}

View file

@ -0,0 +1,3 @@
tick: now
tick: armed
done

View file

@ -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;
}

View file

@ -0,0 +1,3 @@
stale gen 1 ignored
fired gen 2
done

View file

@ -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;
}