From 8863aee459de9d80da834600f82b5651e9593caf Mon Sep 17 00:00:00 2001 From: "shoney.arickathil" Date: Sun, 23 Aug 2026 06:22:24 +0200 Subject: [PATCH] =?UTF-8?q?feat:=20call/reply=20=E2=80=94=20send=20that=20?= =?UTF-8?q?waits=20(id=2088,=20envelope=20kinds=205/6,=20WO-E226)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - runtime: mailbox slots grow caller metadata (wo_msg), call parks on WO_PARK_INBOX (the DB-RPC protocol) and the resume consumes a SCALAR reply; FIBER_DONE ships the receive's return value home (same-shard unpark or kind-6 envelope); kind-5 carries cross-shard calls - actor death is real now: a receive trapping uncaught marks the actor dead, error-unparks the in-flight caller AND every queued caller, drops queued payloads + state, releases cap slots; send-to-dead drops silently, call-to-dead traps — a call never hangs. Fixes the pre-existing leak/dangle in TRAPF's fiber-death path (cur_msg leaked, a->active dangled, the mailbox rotted) - compiler: reply typing through actor-M erasure — every receive(M) program-wide must agree on one return type and it must be a copyable scalar (v1); WO-E226 names disagreeing classes / void receives / non-scalar replies; call's message moves exactly like send's (owner) - corpus: run/call-echo (park + ordered replies), run/call-dead-trap (mid-call + to-dead, both catchable), compile-fail/call-void-receive, compile-fail/call-reply-disagree; cross-shard call proof rides the chat gate next - battery 12/12 fresh-built (ASan+TSan lanes in fibers/db-actor green) Co-Authored-By: Claude Fable 5 --- compiler/src/emit.ml | 8 +- compiler/src/owner.ml | 7 +- compiler/src/types.ml | 122 ++++++++++ docs/plan/oop-vm/08-builtin-surface.md | 1 + runtime/src/builtin.c | 2 + runtime/src/vm.c | 227 ++++++++++++++++-- runtime/src/vm.h | 41 +++- runtime/src/wob.h | 8 +- .../call-reply-disagree/fixture.code | 1 + .../call-reply-disagree/fixture.wo | 27 +++ .../call-void-receive/fixture.code | 1 + .../compile-fail/call-void-receive/fixture.wo | 19 ++ tests/corpus/run/call-dead-trap/fixture.out | 2 + tests/corpus/run/call-dead-trap/fixture.wo | 29 +++ tests/corpus/run/call-echo/fixture.out | 2 + tests/corpus/run/call-echo/fixture.wo | 24 ++ 16 files changed, 494 insertions(+), 27 deletions(-) create mode 100644 tests/corpus/compile-fail/call-reply-disagree/fixture.code create mode 100644 tests/corpus/compile-fail/call-reply-disagree/fixture.wo create mode 100644 tests/corpus/compile-fail/call-void-receive/fixture.code create mode 100644 tests/corpus/compile-fail/call-void-receive/fixture.wo create mode 100644 tests/corpus/run/call-dead-trap/fixture.out create mode 100644 tests/corpus/run/call-dead-trap/fixture.wo create mode 100644 tests/corpus/run/call-echo/fixture.out create mode 100644 tests/corpus/run/call-echo/fixture.wo diff --git a/compiler/src/emit.ml b/compiler/src/emit.ml index 67fc864..271f71f 100644 --- a/compiler/src/emit.ml +++ b/compiler/src/emit.ml @@ -288,6 +288,7 @@ let b_text_of_bytes = 83 let b_sha1 = 85 let b_sha256 = 86 let b_hmac_sha256 = 87 +let b_call = 88 let b_split = 28 let b_split_ws = 29 let b_join = 30 @@ -1099,7 +1100,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"; + "send"; "call"; (* 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"; @@ -3662,6 +3663,8 @@ and emit_builtin (p : pctx) (f : fstate) (v : views) ~(dst : int) ?expected (e : || id = b_float_cmp || id = b_bytes_at || id = b_bytes_eq || id = b_bytes_concat (* iteration 34, two arguments *) || id = b_hmac_sha256 + (* iteration 24, two arguments *) + || id = b_call then 2 else 3 (* b_bytes_slice lands here with substr's shape: (value, start, len) *) in @@ -3697,7 +3700,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" ]) then + (if not (List.mem name [ "push"; "set"; "send"; "call" ]) then List.iteri (fun i (a : Ast.expr) -> (* a reader's result points into arg0 (the container) — dropping @@ -3728,6 +3731,7 @@ and emit_builtin (p : pctx) (f : fstate) (v : views) ~(dst : int) ?expected (e : in 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 *) | "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 135ee82..bfe8f29 100644 --- a/compiler/src/owner.ml +++ b/compiler/src/owner.ml @@ -1344,8 +1344,10 @@ and analyze_call (ctx : ctx) (call_e : Ast.expr) (callee : Ast.expr) (args : Ast | Ident "set" -> (i = 1 || i = 2) && Types.StringMap.find_opt "set" ctx.syms.Types.free_fns = None (* arc: send(addr, msg) MOVES the message to the runtime — the sender's - binding dies (compile-time move, iteration 8's criterion) *) + binding dies (compile-time move, iteration 8's criterion). + 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 | _ -> false in List.iteri @@ -1363,7 +1365,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" -> "cannot be sent — a message moves to the receiver" + | Ident "send" | Ident "call" -> + "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 e1b4425..dec3da7 100644 --- a/compiler/src/types.ml +++ b/compiler/src/types.ml @@ -433,6 +433,7 @@ let nullable_used_without_check_code = Diag.types_prefix ^ "11" let nullable_assign_mismatch_code = Diag.types_prefix ^ "12" let spawn_no_receive_code = Diag.types_prefix ^ "21" (* WO-E221: spawn target lacks fn receive(msg: M); E219/E220 are taken on the language-surface-strictness branch *) let traced_send_code = Diag.types_prefix ^ "22" (* WO-E222: traced(-containing) type in an actor message or actor state — aliased graphs cannot cross heap boundaries *) +let call_reply_code = Diag.types_prefix ^ "26" (* WO-E226 (iteration 24): `call`'s reply through actor-M erasure — every receive(msg: M) program-wide must declare the SAME return type, and it must be a copyable scalar (v1) *) let pub_read_write_code = Diag.types_prefix ^ "19" (* WO-E219: pub(read) field written outside its class *) let using_collision_code = Diag.types_prefix ^ "20" (* WO-E220: using extension collides with a real method *) @@ -1049,6 +1050,47 @@ let builtin_confident_ret (name : string) (arg0 : typ option) : typ option = | "base64_decode" -> Some (TNullable (TScalar "Bytes")) | _ -> None +(* iteration 24: every receive(msg: M) in the program, as + (class_name, reply typ option). `actor M` erases the class, so `call`'s + static reply type exists only if ALL of them agree — the WO-E226 rule. + The caller analyzes this list; building it is one fold over the class + table (bounded by the program, done per call SITE — call sites are + rare enough that a cache is speculative). *) +let call_receivers (classes : class_info StringMap.t) (mname : string) : + (string * typ option) list = + StringMap.fold + (fun cname (cls : class_info) acc -> + match List.find_opt (fun (m : method_info) -> m.name = "receive") cls.methods with + | Some { params = [ (_, pty, _) ]; ret; _ } -> ( + match typ_of_field_ty pty with + | TScalar n when n = mname -> (cname, Option.map typ_of_field_ty ret) :: acc + | _ -> acc) + | _ -> acc) + classes [] + +(* v1: a call reply must be a copyable WORD — the runtime ships it in an + envelope payload with no ownership transfer machinery. Text/Bytes/ + containers/objects are the extension a real consumer earns later. *) +let call_reply_scalar (t : typ) : bool = + match t with + | TActor _ | TRef _ -> true + | TScalar n -> + n = "Int" || n = "Bool" || n = "Timestamp" || n = "Id" || n = "Float" + || is_stdlib_scalar_type n + | _ -> false + +(* The agreed reply type, when everything agrees and is scalar — the + silent half confident_typ uses; the diagnostics half reports. *) +let call_reply_typ (classes : class_info StringMap.t) (mname : string) : typ option = + match call_receivers classes mname with + | [] -> None + | (_, first) :: rest -> + if List.for_all (fun (_, r) -> r = first) rest then + match first with + | Some r when call_reply_scalar r -> Some r + | _ -> None + else None + (* `use_edge`/`uses_of_program`/`path_str` -- relocated here (hotfix) from their original home in the "Modules" section, much further below, purely so `confident_typ`'s free-fn resolution (inside @@ -1302,6 +1344,16 @@ let typecheck_program ~file ~(module_of : string -> string) name). *) match find_variant syms name with | Some (u, _) -> Some (TScalar u.u_name) + | None when name = "call" -> ( + (* iteration 24: call's reply type through the address's + actor M — only when every receive(M) agrees on one + scalar R (WO-E226's silent half). *) + match args with + | addr :: _ -> ( + match confident_typ cenv addr with + | Some (TActor m) -> call_reply_typ syms.classes m + | _ -> None) + | [] -> None) | None -> ( match List.find_opt (fun (n, _, _) -> n = name) builtin_signatures with | None -> None @@ -1740,6 +1792,76 @@ let typecheck_program ~file ~(module_of : string -> string) ~message:"`send`'s first argument must be an `actor M` address" ()) | None -> ()) | _ -> ()) + | None when name = "call" -> + (* iteration 24: call(addr, msg) — send's shape plus the + reply contract (WO-E226): every receive(M) in the + program must declare the same return type, and it must + be a copyable scalar (v1). *) + (if List.length args <> 2 then + Diag.Collector.add collector + (Diag.error ~code:bad_arity_code ~file ~line:e.pos.line ~col:e.pos.col + ~message: + (Printf.sprintf "`call` takes 2 arguments (address, message), given %d" + (List.length args)) + ()) + else + match args with + | [ a; m ] -> ( + match confident_typ cenv a 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 + "this actor receives `%s` — the message is a `%s`" want got) + ()) + | _ -> ()); + match call_receivers syms.classes want with + | [] -> () + | (c0, r0) :: rest -> ( + match + List.find_opt (fun (_, r) -> r <> r0) rest + with + | Some (c1, _) -> + Diag.Collector.add collector + (Diag.error ~code:call_reply_code ~file ~line:e.pos.line + ~col:e.pos.col + ~message: + (Printf.sprintf + "`call` on `actor %s` needs one reply type, but `%s` and `%s` declare different `receive` returns" + want c0 c1) + ()) + | None -> ( + match r0 with + | None -> + Diag.Collector.add collector + (Diag.error ~code:call_reply_code ~file ~line:e.pos.line + ~col:e.pos.col + ~message: + (Printf.sprintf + "`call` needs a reply: `%s`'s `receive(msg: %s)` declares no return type — use `send`" + c0 want) + ()) + | Some r when not (call_reply_scalar r) -> + Diag.Collector.add collector + (Diag.error ~code:call_reply_code ~file ~line:e.pos.line + ~col:e.pos.col + ~message: + (Printf.sprintf + "`call`'s reply type `%s` is not a copyable scalar — v1 replies are scalars (Int, Bool, Float, an actor address, ...)" + (typ_label r)) + ()) + | Some _ -> ()))) + | Some _ -> + Diag.Collector.add collector + (Diag.error ~code:type_mismatch_code ~file ~line:a.pos.line + ~col:a.pos.col + ~message:"`call`'s first 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 c1e82be..574c1e0 100644 --- a/docs/plan/oop-vm/08-builtin-surface.md +++ b/docs/plan/oop-vm/08-builtin-surface.md @@ -277,6 +277,7 @@ 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 | +| `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 | | `net.listen(host, port)` | `-> Int` | IPv4, SO_REUSEADDR, backlog 64; returns an fd | diff --git a/runtime/src/builtin.c b/runtime/src/builtin.c index 8438f2a..74c199d 100644 --- a/runtime/src/builtin.c +++ b/runtime/src/builtin.c @@ -197,6 +197,8 @@ int wo_builtin(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { R[A] = 0; return 0; } + 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_NOW: { /* wall-clock milliseconds */ struct timespec ts; clock_gettime(CLOCK_REALTIME, &ts); diff --git a/runtime/src/vm.c b/runtime/src/vm.c index 1b236fe..ca83c74 100644 --- a/runtime/src/vm.c +++ b/runtime/src/vm.c @@ -72,9 +72,12 @@ static void inbox_push_to(uint32_t shard, wo_envelope *e) { } static void fib_enqueue(wo_vm *vm, wo_fiber *fb); -static int actor_push(wo_actor *a, uint64_t m); static int actor_activate(wo_vm *vm, wo_actor *a); static void fib_reap_all(wo_vm *vm); +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); /* the owning thread drains its inbox: adopt actors, deliver sends, * execute home-routed frees. Returns how many envelopes were handled. */ @@ -92,15 +95,51 @@ static int wo_vm_adopt(wo_vm *vm) { e->actor->next_all = vm->actors; vm->actors = e->actor; break; - case 0: /* a cross-shard send: mailbox + activation on the HOME thread. - The sender already reserved the cap slot; a failed push - (OOM) must hand it back or the slot leaks forever. */ - if (actor_push(e->actor, e->payload) == 0) { + case 0: { /* a cross-shard send: mailbox + activation on the HOME + thread. The sender already reserved the cap slot; a failed + push (OOM) must hand it back or the slot leaks forever. + A dead target drops the moved message silently (the + send-to-dead rule) and frees the slot. */ + if (e->actor->dead) { + actor_drop_payload(vm, e->payload); + wo_mbox_release(e->actor); + break; + } + wo_msg m0 = { e->payload, NULL, 0 }; + if (actor_push(e->actor, m0) == 0) { if (!e->actor->active) (void)actor_activate(vm, e->actor); } else { + actor_drop_payload(vm, e->payload); wo_mbox_release(e->actor); } break; + } + case 5: { /* iteration 24: a cross-shard call — same enqueue as a + send, but the slot remembers the parked caller. A dead + target answers the error reply instead. */ + if (e->actor->dead) { + actor_drop_payload(vm, e->payload); + wo_mbox_release(e->actor); + call_reply_to(vm, e->from_fiber, e->from_shard, 0, WO_T_ACTOR); + break; + } + wo_msg mc = { e->payload, e->from_fiber, e->from_shard }; + if (actor_push(e->actor, mc) == 0) { + if (!e->actor->active) (void)actor_activate(vm, e->actor); + } else { + actor_drop_payload(vm, e->payload); + wo_mbox_release(e->actor); + call_reply_to(vm, e->from_fiber, e->from_shard, 0, WO_T_ACTOR); + } + 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). */ + e->from_fiber->call_reply = e->payload; + e->from_fiber->call_state = e->status ? 3 : 2; + wo_io_unpark(vm, e->from_fiber); + break; case 2: /* a home-routed free: this arena owns the object */ wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)e->payload); break; @@ -518,7 +557,7 @@ void wo_vm_destroy(wo_vm *vm) { wo_actor *nx = a->next_all; if (a->instance) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)a->instance); for (uint32_t i = 0; i < a->mlen; i++) { - uint64_t m = a->msgs[(a->mhead + i) % a->mcap]; + uint64_t m = a->msgs[(a->mhead + i) % a->mcap].payload; if (m) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)m); } free(a->msgs); @@ -624,18 +663,18 @@ void wo_mbox_release(wo_actor *a) { __atomic_fetch_sub(&a->pending, 1, __ATOMIC_ACQ_REL); } -static uint64_t actor_pop(wo_actor *a) { - uint64_t m = a->msgs[a->mhead]; +static wo_msg actor_pop(wo_actor *a) { + wo_msg m = a->msgs[a->mhead]; a->mhead = (a->mhead + 1) % a->mcap; a->mlen--; wo_mbox_release(a); return m; } -static int actor_push(wo_actor *a, uint64_t m) { +static int actor_push(wo_actor *a, wo_msg m) { if (a->mlen == a->mcap) { uint32_t ncap = a->mcap ? a->mcap * 2 : 8; - uint64_t *nm = malloc((size_t)ncap * 8u); + wo_msg *nm = malloc((size_t)ncap * sizeof(wo_msg)); if (!nm) return -1; for (uint32_t i = 0; i < a->mlen; i++) nm[i] = a->msgs[(a->mhead + i) % a->mcap]; free(a->msgs); @@ -648,16 +687,75 @@ static int actor_push(wo_actor *a, uint64_t m) { return 0; } +/* iteration 24: a message the runtime must discard (dead target, failed + * enqueue). Messages are class instances — wo_drop_obj routes a wrong- + * shard drop home through the free envelope. */ +static void actor_drop_payload(wo_vm *vm, uint64_t payload) { + if (payload) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)payload); +} + +/* iteration 24: answer one parked caller. Same-shard callers unpark + * directly; remote ones get a kind-6 envelope. status 0 delivers the + * scalar reply; WO_T_ACTOR makes the caller's re-executed builtin trap. */ +static void call_reply_to(wo_vm *vm, wo_fiber *caller, uint32_t caller_shard, + uint64_t reply, int status) { + if (!caller) return; + if (caller_shard == vm->shard_id) { + caller->call_reply = reply; + caller->call_state = status ? 3 : 2; + wo_io_unpark(vm, caller); + return; + } + wo_envelope *e = calloc(1, sizeof *e); + if (!e) return; /* OOM: the caller stays parked until stop — leak, not UB */ + e->kind = 6; + e->payload = reply; + e->from_fiber = caller; + e->status = status; + inbox_push_to(caller_shard, e); +} + +/* iteration 24: an actor dies (its receive trapped uncaught). Marked on + * the HOME thread only. The in-flight caller and every QUEUED caller get + * the dead error; queued payloads are the runtime's to drop; the moved-in + * state is released. The wo_actor shell itself stays allocated forever + * (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; + if (delivery->cur_msg) { + wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)delivery->cur_msg); + delivery->cur_msg = 0; + } + call_reply_to(vm, delivery->msg_caller, delivery->msg_caller_shard, 0, WO_T_ACTOR); + delivery->msg_caller = NULL; + while (a->mlen) { + wo_msg m = actor_pop(a); + if (m.payload) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)m.payload); + call_reply_to(vm, m.caller, m.caller_shard, 0, WO_T_ACTOR); + } + if (a->instance) { + wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)a->instance); + a->instance = 0; + } + a->active = NULL; +} + /* Mailbox nonempty, no delivery fiber: start one on the next message. * receive borrows both self and the message; the runtime keeps ownership * of the message (fiber->cur_msg) and drops it when the call returns. */ static int actor_activate(wo_vm *vm, wo_actor *a) { - uint64_t m = actor_pop(a); - uint64_t args[2] = { a->instance, m }; + wo_msg m = actor_pop(a); + uint64_t args[2] = { a->instance, m.payload }; wo_fiber *fb = wo_vm_spawn_fiber(vm, a->method, args, 2); - if (!fb) return -1; + if (!fb) { + actor_drop_payload(vm, m.payload); + call_reply_to(vm, m.caller, m.caller_shard, 0, WO_T_ACTOR); + return -1; + } fb->actor = a; - fb->cur_msg = m; + fb->cur_msg = m.payload; + fb->msg_caller = m.caller; + fb->msg_caller_shard = m.caller_shard; a->active = fb; return 0; } @@ -714,6 +812,14 @@ int wo_vm_actor_send(wo_vm *vm, uint64_t addr, uint64_t msg_val, const char **ms *msg = "send: nil message"; return WO_T_BOUNDS; } + /* send-to-dead is a silent drop (spec'd v1): the message moved to the + * runtime, so the runtime discards it. The dead flag is written on the + * home thread; a racing remote read at worst enqueues an envelope the + * home drain then discards through its own dead check. */ + if (a->dead) { + actor_drop_payload(vm, msg_val); + return 0; + } /* iteration 24: the cap check happens SENDER-side on every path, so * the sender always learns — fail-fast backpressure, catchable. */ if (wo_mbox_reserve(a) != 0) { @@ -736,7 +842,8 @@ int wo_vm_actor_send(wo_vm *vm, uint64_t addr, uint64_t msg_val, const char **ms inbox_push_to(a->home, e); return 0; } - if (actor_push(a, msg_val) != 0) { + wo_msg m0 = { msg_val, NULL, 0 }; + if (actor_push(a, m0) != 0) { wo_mbox_release(a); *msg = "out of memory"; return WO_T_OOM; @@ -748,6 +855,75 @@ int wo_vm_actor_send(wo_vm *vm, uint64_t addr, uint64_t msg_val, const char **ms return 0; } +/* 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: + * call-to-dead traps immediately, callee-dies-mid-call error-unparks. */ +int wo_vm_actor_call(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { + uint32_t A = wo_ins_a(ins), B = wo_ins_b(ins); + wo_fiber *fb = vm->cur; + + if (fb->call_state == 2) { /* the reply: consume it */ + fb->call_state = 0; + R[A] = fb->call_reply; + return 0; + } + if (fb->call_state == 3) { /* the callee was/went dead */ + fb->call_state = 0; + *msg = "actor died during call"; + return WO_T_ACTOR; + } + + wo_actor *a = (wo_actor *)(uintptr_t)R[B]; + uint64_t msg_val = R[B + 1]; + if (!a) { + *msg = "call: nil actor address"; + return WO_T_BOUNDS; + } + if (!msg_val) { + *msg = "call: nil message"; + return WO_T_BOUNDS; + } + if (a->dead) { /* unlike send, the caller MUST learn */ + actor_drop_payload(vm, msg_val); + *msg = "actor died during call"; + return WO_T_ACTOR; + } + if (wo_mbox_reserve(a) != 0) { + *msg = "actor mailbox full"; + return WO_T_ACTOR; + } + if (a->home != vm->shard_id) { + wo_envelope *e = calloc(1, sizeof *e); + if (!e) { + wo_mbox_release(a); + *msg = "out of memory"; + return WO_T_OOM; + } + e->kind = 5; + e->actor = a; + e->payload = msg_val; + e->from_shard = vm->shard_id; + e->from_fiber = fb; + inbox_push_to(a->home, e); + } else { + wo_msg mc = { msg_val, fb, vm->shard_id }; + if (actor_push(a, mc) != 0) { + wo_mbox_release(a); + *msg = "out of memory"; + return WO_T_OOM; + } + if (!a->active && actor_activate(vm, a) != 0) { + *msg = "out of memory"; + return WO_T_OOM; + } + } + fb->call_state = 1; + fb->park_fd = WO_PARK_INBOX; + fb->park_done = 0; /* resume RE-EXECUTES the builtin: the consume path */ + return WO_SYS_PARKED; +} + /* 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) { @@ -830,7 +1006,7 @@ static void vm_gc_roots(wo_vm *vm) { for (const wo_actor *a = vm->actors; a; a = a->next_all) { if (a->instance) wo_gc_scan_root(&vm->rt, (wo_hdr *)(uintptr_t)a->instance); for (uint32_t i = 0; i < a->mlen; i++) { - uint64_t m = a->msgs[(a->mhead + i) % a->mcap]; + uint64_t m = a->msgs[(a->mhead + i) % a->mcap].payload; if (m) wo_gc_scan_root(&vm->rt, (wo_hdr *)(uintptr_t)m); } if (a->active && a->active->cur_msg) @@ -992,6 +1168,12 @@ static int vm_run(wo_vm *vm, uint64_t *ret, wo_err *err) { "wovm: fiber trap %d at %s:%d: %s\n", \ err->code, err->method, err->line, err->msg); \ wo_fiber *dead = vm->cur; \ + /* iteration 24: a receive trapping uncaught kills the ACTOR, \ + * not just the fiber — the dead flag, the in-flight caller, \ + * every queued caller, the state and the mailbox are all \ + * settled here (before: cur_msg leaked and a->active \ + * dangled — the actor's mailbox rotted forever). */ \ + if (dead->actor) actor_die(vm, dead->actor, dead); \ vm->nfibers--; \ free(dead); \ NEXT_RUNNABLE(); \ @@ -1332,10 +1514,15 @@ dispatch: wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)dead->cur_msg); \ dead->cur_msg = 0; \ } \ + /* iteration 24: the receive's return value IS the reply — \ + * ship it before this context is reused or freed */ \ + call_reply_to(vm, dead->msg_caller, dead->msg_caller_shard, \ + (rv), 0); \ + dead->msg_caller = NULL; \ if (a->mlen) { \ /* next message: REUSE this context, re-queued for \ * fairness (one message per turn, never a monopolist) */ \ - uint64_t m_ = actor_pop(a); \ + wo_msg m_ = actor_pop(a); \ const wo_methodrec *sme_ = &vm->mod->methods[a->method]; \ dead->depth = 1; \ dead->ncatch = 0; \ @@ -1343,9 +1530,11 @@ dispatch: dead->frames[0].pc = 0; \ dead->frames[0].base = 0; \ dead->regs[0] = a->instance; \ - dead->regs[1] = m_; \ + dead->regs[1] = m_.payload; \ memset(dead->regs + 2, 0, (size_t)(sme_->reg_cnt - 2) * 8u); \ - dead->cur_msg = m_; \ + dead->cur_msg = m_.payload; \ + dead->msg_caller = m_.caller; \ + dead->msg_caller_shard = m_.caller_shard; \ fib_enqueue(vm, dead); \ NEXT_RUNNABLE(); \ RELOAD(); \ diff --git a/runtime/src/vm.h b/runtime/src/vm.h index 13aca19..b86616a 100644 --- a/runtime/src/vm.h +++ b/runtime/src/vm.h @@ -82,6 +82,17 @@ typedef struct wo_fiber { /* arc stage 3: the in-flight DB request while parked on the DB actor's * reply (a wo_db_req*, opaque here; vm.c owns the protocol) */ void *dbreq; + /* iteration 24, caller side of call(): 0 = no call in flight, + * 1 = parked awaiting the reply, 2 = reply landed (call_reply is the + * scalar), 3 = the callee was/went dead (the re-executed builtin + * traps WO_T_ACTOR). Set on the caller's own thread or under its + * shard's inbox drain — never concurrently with the fiber running. */ + int call_state; + uint64_t call_reply; + /* iteration 24, delivery side: the CURRENT message's caller (NULL for + * a plain send) — where FIBER_DONE ships the receive's return value. */ + struct wo_fiber *msg_caller; + uint32_t msg_caller_shard; } wo_fiber; /* arc stage 3: park_fd sentinel — PARKED with NO plane wait; the wake is @@ -89,14 +100,26 @@ typedef struct wo_fiber { * scans, which key on park_fd == -1 exactly. */ #define WO_PARK_INBOX (-2) +/* iteration 24: one mailbox slot. A plain send has caller == NULL; a + * call carries the parked caller so the delivery's return value can + * route home as a kind-6 envelope (or a same-shard unpark). */ +typedef struct wo_msg { + uint64_t payload; + struct wo_fiber *caller; /* NULL = send */ + uint32_t caller_shard; +} wo_msg; + /* An actor: moved-in state, its receive method, a FIFO mailbox, and at * most one delivery fiber at a time (one message at a time — the actor - * guarantee). Actors live until program end (v1: no actor death). */ + * 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. */ 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) */ uint32_t home; /* the shard whose thread owns mailbox + delivery */ - uint64_t *msgs; /* FIFO ring, growable up to the cap */ + int dead; /* set on the home thread when a receive traps */ + wo_msg *msgs; /* FIFO ring, growable up to the cap */ uint32_t mhead, mlen, mcap; /* iteration 24: sent-but-not-delivered count, incremented by the * SENDER on any shard (the cap check), decremented by the home @@ -158,6 +181,10 @@ typedef struct wo_vm { int wo_vm_actor_spawn(wo_vm *vm, uint64_t instance, uint32_t method_idx, uint64_t *out_addr, const char **msg); int wo_vm_actor_send(wo_vm *vm, uint64_t addr, uint64_t msg_val, const char **msg); +/* iteration 24: send-that-waits. First entry enqueues the message with 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). */ +int wo_vm_actor_call(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg); /* ---- the shard engine (arc stage 2) ------------------------------------ * One pinned thread per shard, each a full wo_vm (own arena, GC, I/O @@ -178,9 +205,17 @@ typedef struct wo_envelope { int kind; /* 0 = SEND (actor, payload), 1 = SPAWN-ADOPT (actor), * 2 = FREE (payload = wo_hdr*), * 3 = DB_REQ (payload = wo_db_req*, to shard 0), - * 4 = DB_RESP (payload = wo_db_req*, back to the requester) */ + * 4 = DB_RESP (payload = wo_db_req*, back to the requester), + * 5 = CALL (iteration 24: actor, payload = moved message, + * 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) */ struct wo_actor *actor; uint64_t payload; + uint32_t from_shard; + struct wo_fiber *from_fiber; + int status; } wo_envelope; /* arc stage 3: the requester half of the transparent DB RPC (vm.c). Called diff --git a/runtime/src/wob.h b/runtime/src/wob.h index c18d1a9..36d44a0 100644 --- a/runtime/src/wob.h +++ b/runtime/src/wob.h @@ -462,9 +462,15 @@ enum { WO_B_SHA256 = 86, /* (bytes) -> fresh 32-byte Bytes */ WO_B_HMAC_SHA256 = 87, /* (key bytes, msg bytes) -> fresh 32-byte Bytes, * RFC 2104 (key > 64 bytes hashed first) */ + /* ---- iteration 24: actor lifecycle ---- */ + WO_B_CALL = 88, /* (addr, msg) -> R: send that WAITS — the + * caller's fiber parks until the receive's + * return value arrives. R is a SCALAR (v1, + * compiler-enforced WO-E226). Dead callee = + * WO_T_ACTOR, immediately or mid-call. */ }; -#define WO_B_MAX 87u +#define WO_B_MAX 88u /* 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/tests/corpus/compile-fail/call-reply-disagree/fixture.code b/tests/corpus/compile-fail/call-reply-disagree/fixture.code new file mode 100644 index 0000000..2c30994 --- /dev/null +++ b/tests/corpus/compile-fail/call-reply-disagree/fixture.code @@ -0,0 +1 @@ +WO-E226 \ No newline at end of file diff --git a/tests/corpus/compile-fail/call-reply-disagree/fixture.wo b/tests/corpus/compile-fail/call-reply-disagree/fixture.wo new file mode 100644 index 0000000..57cd582 --- /dev/null +++ b/tests/corpus/compile-fail/call-reply-disagree/fixture.wo @@ -0,0 +1,27 @@ +-- iteration 24: `actor M` erases the class, so `call`'s reply type is +-- only defined when every receive(msg: M) program-wide agrees on one +-- return type. Two classes disagreeing is WO-E226 naming both. + +class Q { + n: Int +} + +class IntAnswer { + pad: Int + fn receive(msg: Q) -> Int { + return msg.n; + } +} + +class BoolAnswer { + pad: Int + fn receive(msg: Q) -> Bool { + return msg.n > 0; + } +} + +fn main() -> Int { + let a: actor Q = spawn IntAnswer { pad: 0 }; + let r = call(a, Q { n: 1 }); + return r; +} diff --git a/tests/corpus/compile-fail/call-void-receive/fixture.code b/tests/corpus/compile-fail/call-void-receive/fixture.code new file mode 100644 index 0000000..2c30994 --- /dev/null +++ b/tests/corpus/compile-fail/call-void-receive/fixture.code @@ -0,0 +1 @@ +WO-E226 \ No newline at end of file diff --git a/tests/corpus/compile-fail/call-void-receive/fixture.wo b/tests/corpus/compile-fail/call-void-receive/fixture.wo new file mode 100644 index 0000000..b61bf07 --- /dev/null +++ b/tests/corpus/compile-fail/call-void-receive/fixture.wo @@ -0,0 +1,19 @@ +-- iteration 24: `call` needs a reply — a receive with no return type +-- serves `send` only (WO-E226). + +class Q { + n: Int +} + +class Sink { + pad: Int + fn receive(msg: Q) { + self.pad = msg.n; + } +} + +fn main() -> Int { + let s: actor Q = spawn Sink { pad: 0 }; + let r = call(s, Q { n: 1 }); + return r; +} diff --git a/tests/corpus/run/call-dead-trap/fixture.out b/tests/corpus/run/call-dead-trap/fixture.out new file mode 100644 index 0000000..0ba05e0 --- /dev/null +++ b/tests/corpus/run/call-dead-trap/fixture.out @@ -0,0 +1,2 @@ +mid-call: actor died during call +to-dead: actor died during call diff --git a/tests/corpus/run/call-dead-trap/fixture.wo b/tests/corpus/run/call-dead-trap/fixture.wo new file mode 100644 index 0000000..7b4a7c6 --- /dev/null +++ b/tests/corpus/run/call-dead-trap/fixture.wo @@ -0,0 +1,29 @@ +-- iteration 24: no hangs, ever. A receive that traps kills the ACTOR: +-- the parked caller is error-unparked into a catchable WO_T_ACTOR +-- (callee died mid-call), and a later call to the same address traps +-- immediately (call-to-dead). The runtime prints the fiber-trap report +-- to stderr; stdout carries only the caller's side. +class Q { + n: Int +} + +class Bomb { + pad: Int + fn receive(msg: Q) -> Int { + return msg.n / (msg.n - msg.n); -- DIV0: the actor dies mid-call + } +} + +fn call_one(b: actor Q, n: Int) -> Text { + let r = call(b, Q { n: n }); + return "replied ${r}"; +} + +fn main() -> Int { + let b: actor Q = spawn Bomb { pad: 0 }; + let first = try call_one(b, 7) catch (e) e.msg; + print("mid-call: ${first}"); + let second = try call_one(b, 8) catch (e) e.msg; + print("to-dead: ${second}"); + return 0; +} diff --git a/tests/corpus/run/call-echo/fixture.out b/tests/corpus/run/call-echo/fixture.out new file mode 100644 index 0000000..8d936d8 --- /dev/null +++ b/tests/corpus/run/call-echo/fixture.out @@ -0,0 +1,2 @@ +a=42 +b=84 diff --git a/tests/corpus/run/call-echo/fixture.wo b/tests/corpus/run/call-echo/fixture.wo new file mode 100644 index 0000000..c985149 --- /dev/null +++ b/tests/corpus/run/call-echo/fixture.wo @@ -0,0 +1,24 @@ +-- iteration 24: call — send that waits. The caller parks until the +-- receive's return value comes back; same-shard here (the runner pins +-- WO_SHARDS=1), so the reply path is the direct unpark. The actor's +-- state update across two calls proves the mailbox kept its order. +class Q { + n: Int +} + +class Doubler { + seen: Int + fn receive(msg: Q) -> Int { + self.seen = self.seen + 1; + return msg.n * 2; + } +} + +fn main() -> Int { + let d: actor Q = spawn Doubler { seen: 0 }; + let a = call(d, Q { n: 21 }); + print("a=${a}"); + let b = call(d, Q { n: a }); + print("b=${b}"); + return 0; +}