Skip to main content

shape_vm/executor/async_ops/
mod.rs

1//! Async operations for the VM executor.
2//!
3//! # Concurrency Model
4//!
5//! The Shape VM uses **cooperative, single-threaded concurrency**. All async
6//! operations execute on the thread that owns the `VirtualMachine` instance --
7//! there is no work-stealing or multi-threaded task execution within the VM
8//! itself. The VM is `!Sync` by design.
9//!
10//! ## Task Lifecycle
11//!
12//! 1. **Spawn** (`SpawnTask`): Pops a callable from the stack, assigns a
13//!    monotonic future ID, registers it with the `TaskScheduler`, and pushes
14//!    a `Future(id)` value onto the stack.
15//! 2. **Await** (`Await`): Pops a `Future(id)`, attempts synchronous inline
16//!    resolution via the `TaskScheduler`. If the task cannot be resolved
17//!    (e.g., it depends on an external I/O operation), execution suspends
18//!    with `VMError::Suspended` so the host runtime can schedule it.
19//! 3. **Join** (`JoinInit` + `JoinAwait`): Collects multiple futures into a
20//!    `TaskGroup` value, then resolves them according to a join strategy
21//!    (all, race, any, all-settled).
22//! 4. **Cancel** (`CancelTask`): Marks a task as cancelled in the scheduler.
23//!
24//! ## Structured Concurrency
25//!
26//! `AsyncScopeEnter` / `AsyncScopeExit` bracket a structured concurrency
27//! region. All tasks spawned within a scope are tracked; on scope exit, any
28//! still-pending tasks are cancelled in LIFO order. This guarantees that no
29//! task outlives its enclosing scope.
30//!
31//! ## Suspension Protocol
32//!
33//! When an operation cannot complete synchronously, it returns
34//! `AsyncExecutionResult::Suspended(SuspensionInfo)`. The dispatch layer in
35//! `dispatch.rs` converts this into `VMError::Suspended { future_id, resume_ip }`
36//! which propagates up to the host runtime. The host resolves the future and
37//! calls back into the VM to resume execution at `resume_ip`.
38//!
39//! ## Opcodes Handled
40//!
41//! `Yield`, `Suspend`, `Resume`, `Poll`, `AwaitBar`, `AwaitTick`,
42//! `EmitAlert`, `EmitEvent`, `Await`, `SpawnTask`, `JoinInit`, `JoinAwait`,
43//! `CancelTask`, `AsyncScopeEnter`, `AsyncScopeExit`.
44//!
45//! ## Wave 6.5 / E-async migration (ADR-006 §2.7.7 / Q9, §10 E-async row)
46//!
47//! Every push/pop in this file threads the kinded API
48//! (`push_kinded(bits, kind)` / `pop_kinded()`) per the playbook §2 / §3
49//! kind-sourcing rules. Future and TaskGroup payload kinds:
50//!
51//! - `Future(id)` ⇒ `NativeKind::Ptr(HeapKind::Future)` — inline scalar
52//!   payload (the future ID is stored directly in `bits`; no `Arc<T>`).
53//! - `TaskGroup(Arc<TaskGroupData>)` ⇒ `NativeKind::Ptr(HeapKind::TaskGroup)`
54//!   — `Arc<TaskGroupData>` payload per ADR-006 §2.3.
55//!
56//! ## Wave 8 W8-AS migration (ADR-006 §2.7.11/Q12, §2.7.4 Phase-2c boundary)
57//!
58//! The `task_scheduler::TaskScheduler` API was migrated to the kinded
59//! `(bits, kind)` carrier shape during Wave 6.5 R-async-time / E-async
60//! close, and `call_convention.rs::resolve_spawned_task` was filled by
61//! W7-cv-async (close `f3502b0`) per §2.7.11/Q12 — sync resolution of
62//! spawned closures + function-id callables routes through
63//! `call_closure_with_nb_args_keepalive` / `call_function_with_nb_args`.
64//! W8-AS lights up the per-await-site integration:
65//!
66//! - `op_await` resolves a `Future(id)`-kinded slot synchronously via
67//!   `vm.resolve_spawned_task(task_id)` and pushes the kinded result.
68//! - `op_spawn_task` allocates a fresh `future_id`, transfers the popped
69//!   callable share into `task_scheduler.register(id, bits, kind)`,
70//!   tracks the future id in the active async-scope stack, and pushes
71//!   the future id as `Ptr(HeapKind::Future)`.
72//! - `op_join_await` walks the carried `task_ids` and dispatches per-id
73//!   to `resolve_spawned_task` per the join strategy (All / Race / Any
74//!   / AllSettled), aggregating into an `Arc<TaskGroupData>` `Ptr(HeapKind::TaskGroup)`
75//!   carrier (Race/Any return the per-task result directly).
76//!
77//! ### §2.7.4 Phase-2c boundary
78//!
79//! Suspension state crossing a `resolve_spawned_task` frame boundary —
80//! a `VMError::Suspended` raised inside a spawned closure body — stays
81//! out of scope. The current sync-resolution path propagates the error
82//! upward; the task's cached entry remains `Pending` until a future
83//! Phase-2c rebuild lands the snapshot-tier resumption.
84
85use crate::{
86    bytecode::{Instruction, OpCode, Operand},
87    executor::VirtualMachine,
88    executor::vm_impl::stack::drop_with_kind,
89};
90use shape_value::{
91    NativeKind, VMError,
92    heap_value::{HeapKind, TaskGroupData},
93};
94use std::sync::Arc;
95
96/// Result of executing an async operation
97#[derive(Debug, Clone)]
98pub enum AsyncExecutionResult {
99    /// Continue normal execution
100    Continue,
101    /// Yield to event loop (cooperative scheduling)
102    Yielded,
103    /// Suspended waiting for external event
104    Suspended(SuspensionInfo),
105}
106
107/// Information about why execution was suspended
108#[derive(Debug, Clone)]
109pub struct SuspensionInfo {
110    /// What we're waiting for
111    pub wait_type: WaitType,
112    /// Instruction pointer to resume at
113    pub resume_ip: usize,
114}
115
116/// Type of wait condition
117#[derive(Debug, Clone)]
118pub enum WaitType {
119    /// Waiting for next data bar from source
120    NextBar { source: String },
121    /// Waiting for timer
122    Timer { id: u64 },
123    /// Waiting for any event
124    AnyEvent,
125    /// Waiting for a future to resolve (general-purpose await)
126    Future { id: u64 },
127    /// Waiting for a task group to resolve (join await)
128    TaskGroup { kind: u8, task_ids: Vec<u64> },
129}
130
131impl VirtualMachine {
132    /// Execute an async opcode
133    ///
134    /// Returns `AsyncExecutionResult` to indicate whether execution should
135    /// continue, yield, or suspend.
136    #[inline(always)]
137    pub(in crate::executor) fn exec_async_op(
138        &mut self,
139        instruction: &Instruction,
140    ) -> Result<AsyncExecutionResult, VMError> {
141        use OpCode::*;
142        match instruction.opcode {
143            Yield => self.op_yield(),
144            Suspend => self.op_suspend(instruction),
145            Resume => self.op_resume(instruction),
146            Poll => self.op_poll(),
147            AwaitBar => self.op_await_bar(instruction),
148            AwaitTick => self.op_await_tick(instruction),
149            EmitAlert => self.op_emit_alert(),
150            EmitEvent => self.op_emit_event(),
151            Await => self.op_await(),
152            SpawnTask => self.op_spawn_task(),
153            JoinInit => self.op_join_init(instruction),
154            JoinAwait => self.op_join_await(),
155            CancelTask => self.op_cancel_task(),
156            AsyncScopeEnter => self.op_async_scope_enter(),
157            AsyncScopeExit => self.op_async_scope_exit(),
158            _ => unreachable!(
159                "exec_async_op called with non-async opcode: {:?}",
160                instruction.opcode
161            ),
162        }
163    }
164
165    /// Yield to the event loop for cooperative scheduling
166    ///
167    /// This allows other tasks to run and prevents long-running
168    /// computations from blocking the event loop.
169    fn op_yield(&mut self) -> Result<AsyncExecutionResult, VMError> {
170        // Save current state - the IP is already pointing to next instruction
171        Ok(AsyncExecutionResult::Yielded)
172    }
173
174    /// Suspend execution until a condition is met
175    ///
176    /// The operand specifies the wait condition type.
177    fn op_suspend(&mut self, instruction: &Instruction) -> Result<AsyncExecutionResult, VMError> {
178        let wait_type = match &instruction.operand {
179            Some(Operand::Const(idx)) => {
180                // Get wait type from constant pool
181                // For now, default to waiting for any event
182                let _ = idx;
183                WaitType::AnyEvent
184            }
185            _ => WaitType::AnyEvent,
186        };
187
188        Ok(AsyncExecutionResult::Suspended(SuspensionInfo {
189            wait_type,
190            resume_ip: self.ip,
191        }))
192    }
193
194    /// Resume from suspension
195    ///
196    /// Called by the runtime when resuming suspended execution.
197    /// The resume value (if any) should be on the stack.
198    fn op_resume(&mut self, _instruction: &Instruction) -> Result<AsyncExecutionResult, VMError> {
199        // Resume is handled by the outer execution loop
200        // This opcode is a marker for where to resume
201        Ok(AsyncExecutionResult::Continue)
202    }
203
204    /// Poll the event queue
205    ///
206    /// Pushes the next event from the queue onto the stack,
207    /// or null if the queue is empty.
208    fn op_poll(&mut self) -> Result<AsyncExecutionResult, VMError> {
209        // In the VM, we don't have direct access to the event queue
210        // This is handled via the VMContext passed from the runtime.
211        // R5b-2-bool-null-sentinel-cluster (ADR-006 §2.7 + §2.7.7/Q9,
212        // 2026-05-19): no event available — push `NativeKind::Null` per
213        // post-disposition; pre-disposition `(0u64, NativeKind::Bool)`
214        // collided with legitimate `false` bool slots.
215        self.push_kinded(0u64, NativeKind::Null)?;
216        Ok(AsyncExecutionResult::Continue)
217    }
218
219    /// Await next data bar from a source
220    ///
221    /// Suspends execution until the next data point arrives
222    /// from the specified source.
223    fn op_await_bar(&mut self, instruction: &Instruction) -> Result<AsyncExecutionResult, VMError> {
224        let source = match &instruction.operand {
225            Some(Operand::Const(idx)) => {
226                // Get source name from constant pool
227                match self.program.constants.get(*idx as usize) {
228                    Some(crate::bytecode::Constant::String(s)) => s.clone(),
229                    _ => "default".to_string(),
230                }
231            }
232            _ => "default".to_string(),
233        };
234
235        Ok(AsyncExecutionResult::Suspended(SuspensionInfo {
236            wait_type: WaitType::NextBar { source },
237            resume_ip: self.ip,
238        }))
239    }
240
241    /// Await next timer tick
242    ///
243    /// Suspends execution until the specified timer fires.
244    fn op_await_tick(
245        &mut self,
246        instruction: &Instruction,
247    ) -> Result<AsyncExecutionResult, VMError> {
248        let timer_id = match &instruction.operand {
249            Some(Operand::Const(idx)) => {
250                // Get timer ID from constant pool
251                match self.program.constants.get(*idx as usize) {
252                    Some(crate::bytecode::Constant::Number(n)) => *n as u64,
253                    _ => 0,
254                }
255            }
256            _ => 0,
257        };
258
259        Ok(AsyncExecutionResult::Suspended(SuspensionInfo {
260            wait_type: WaitType::Timer { id: timer_id },
261            resume_ip: self.ip,
262        }))
263    }
264
265    /// Emit an alert to the alert pipeline
266    ///
267    /// Pops an alert object from the stack and sends it to
268    /// the alert router for processing.
269    fn op_emit_alert(&mut self) -> Result<AsyncExecutionResult, VMError> {
270        // Pop the alert payload and release its share — alert pipeline
271        // integration is deferred. Drop discipline (playbook §3): every
272        // `pop_kinded` either re-pushes or `drop_with_kind`s.
273        let (bits, kind) = self.pop_kinded()?;
274        drop_with_kind(bits, kind);
275        Ok(AsyncExecutionResult::Continue)
276    }
277
278    /// General-purpose await
279    ///
280    /// Pops a value from the stack. If it's a Future(id), attempts to resolve
281    /// the task inline from the task scheduler. If the task's callable is a
282    /// plain value (not a closure/function), it is used directly as the result.
283    /// Otherwise, suspends execution so the host runtime can schedule the task.
284    /// If the value is not a Future, pushes it back (sync shortcut).
285    fn op_await(&mut self) -> Result<AsyncExecutionResult, VMError> {
286        let sp_before = self.sp;
287        let (bits, kind) = self.pop_kinded()?;
288        match kind {
289            NativeKind::Ptr(HeapKind::Future) => {
290                // Future(id) is an inline scalar — `bits` IS the future ID
291                // (see TaskScheduler docstring + §2.7.11/Q12 Future row).
292                // No Arc share to drop on the popped slot; HeapKind::Future
293                // is a no-op in drop_with_kind / clone_with_kind.
294                //
295                // Sync-resolution path (§2.7.11/Q12 dispatch precedent,
296                // closed by W7-cv-async at `f3502b0`): hand the future id
297                // to `resolve_spawned_task`, which routes through the
298                // kinded `call_*_with_nb_args` family for closure /
299                // function-id callables and returns the result `KindedSlot`.
300                //
301                // Phase-2c boundary (ADR-006 §2.7.4): if the spawned
302                // body suspends mid-execution, `resolve_spawned_task`
303                // propagates `VMError::Suspended` up. The task's cached
304                // entry remains `Pending` until a future snapshot-tier
305                // rebuild lands. This is explicitly out-of-scope here.
306                let task_id = bits;
307                let result = self.resolve_spawned_task(task_id)?;
308
309                // Transfer the result share onto the stack via the
310                // canonical `push_kinded(raw, kind)` + `mem::forget`
311                // pattern (per §2.7.10/§2.7.11 dispatch-shell shape —
312                // `control_flow/mod.rs::dispatch_call_value_immediate`
313                // is the precedent). The carrier's Drop must not fire
314                // after the share moves to the stack slot.
315                self.push_kinded(result.raw(), result.kind())?;
316                std::mem::forget(result);
317                debug_assert_eq!(
318                    self.sp, sp_before,
319                    "op_await (Future): stack depth changed (before={}, after={})",
320                    sp_before, self.sp
321                );
322                Ok(AsyncExecutionResult::Continue)
323            }
324            _ => {
325                // Sync shortcut: value is already resolved, push it back.
326                // The popped share transfers directly back onto the stack —
327                // no `clone_with_kind` / `drop_with_kind` needed.
328                self.push_kinded(bits, kind)?;
329                debug_assert_eq!(
330                    self.sp, sp_before,
331                    "op_await (sync shortcut): stack depth changed (before={}, after={})",
332                    sp_before, self.sp
333                );
334                Ok(AsyncExecutionResult::Continue)
335            }
336        }
337    }
338
339    /// Await with a timeout.
340    ///
341    /// Spawn a task from a callable or pre-resolved value on the stack.
342    ///
343    /// Pops a top-of-stack kinded slot and creates a new async task identified
344    /// by a fresh future id. Pushes a `Ptr(HeapKind::Future)` onto the stack
345    /// representing the spawned task.
346    ///
347    /// Two callable-classification paths per the §2.7.11/Q12 value-call dispatch:
348    ///
349    /// - **Callable kinds** (`NativeKind::Ptr(HeapKind::Closure)` /
350    ///   `NativeKind::UInt64`): registered with the scheduler via `register`
351    ///   for later inline execution at `resolve_spawned_task` time. Same shape
352    ///   as the W7-cv-static / W7-cv-async dispatch classification at
353    ///   `call_convention.rs::resolve_spawned_task` (line 438).
354    ///
355    /// - **Non-callable kinds** (every other `NativeKind`, including scalars
356    ///   `Int64`/`Float64`/`Bool`/`String` and heap-bearing
357    ///   `Ptr(HeapKind::*)` for non-callable heap types): treated as already-
358    ///   resolved values per the §2.7.4 sync-shortcut semantic. The compiler
359    ///   surface emits `compile_expr + SpawnTask` for `async let x = 42`,
360    ///   `join all { 1+2, 3+4 }`, and other RHS / branch expressions that
361    ///   evaluate to plain values (not closure literals or function
362    ///   references). Without this path, `resolve_spawned_task` would
363    ///   surface `callable must be NativeKind::Ptr(HeapKind::Closure) or
364    ///   NativeKind::UInt64` for every non-callable RHS.
365    ///
366    ///   The non-callable path uses the scheduler's existing `complete(id,
367    ///   bits, kind)` API (same shape as the external-completion path used
368    ///   by remote calls); the share transfers from the stack slot into the
369    ///   scheduler's cached-result entry. `resolve_spawned_task` then hits
370    ///   its `TaskStatus::Completed` cached fast-path at line 407 and
371    ///   returns the value cleanly via `clone_with_kind`.
372    ///
373    /// If inside an async scope, the spawned future id is tracked for
374    /// cancellation. Cancellation of a non-callable / pre-completed task is
375    /// a no-op (the result is already stored — `cancel` only releases the
376    /// `callables` entry, which is empty for non-callable kinds).
377    fn op_spawn_task(&mut self) -> Result<AsyncExecutionResult, VMError> {
378        let sp_before = self.sp;
379        // Pop the top-of-stack kinded slot. The share transfers to the
380        // task_scheduler via either `register` (callable kinds) or
381        // `complete` (non-callable kinds) — same retain-on-store contract
382        // as the §2.7.7 stack and §2.7.8 cell-storage tracks (one
383        // strong-count share owned by the storage; released by
384        // `take_callable` / cached-result drop / `Drop`). No
385        // `drop_with_kind` here: the share is moved, not released.
386        let (slot_bits, slot_kind) = self.pop_kinded()?;
387
388        // Allocate a fresh future id. `next_future_id` is monotonic and
389        // single-threaded (the VM is `!Sync` per the module docstring's
390        // concurrency model section).
391        let task_id = self.next_future_id();
392
393        // Classify the popped slot per the §2.7.11/Q12 callable-vs-value
394        // dispatch shape (mirrors `call_convention.rs::resolve_spawned_task`
395        // line 438 + `call_value_immediate_nb` line 854). Callables route
396        // through `register` for later inline execution; non-callable
397        // values route through `complete` as pre-resolved results.
398        match slot_kind {
399            NativeKind::Ptr(HeapKind::Closure) | NativeKind::UInt64 => {
400                // Callable — register for later execution.
401                self.task_scheduler.register(task_id, slot_bits, slot_kind);
402            }
403            _ => {
404                // Non-callable pre-resolved value. Register with Pending
405                // status so the scope tracking + `is_resolved` checks
406                // behave uniformly with the callable path, then `complete`
407                // immediately so `resolve_spawned_task` hits the
408                // cached-result fast-path. The scheduler's `register`
409                // requires a `(bits, kind)` pair to slot into the
410                // `callables` map for refcount discipline; for the
411                // pre-resolved path we transfer the share directly to
412                // the `results` map via `complete` and ensure the
413                // `callables` entry is empty so the take-callable arm
414                // never fires.
415                //
416                // Refcount discipline: `complete(task_id, bits, kind)`
417                // takes one strong-count share into the cached-result
418                // entry; the share transferred to us from `pop_kinded`
419                // moves into the scheduler. `resolve_spawned_task`'s
420                // cached fast-path at `call_convention.rs:407` clones a
421                // fresh share via `clone_with_kind` for the returned
422                // `KindedSlot`, leaving the cached entry's share intact.
423                self.task_scheduler.complete(task_id, slot_bits, slot_kind);
424            }
425        }
426
427        // Track the spawned future id in the active async scope (if any)
428        // so `op_async_scope_exit` can cancel still-pending tasks in
429        // LIFO order (structured concurrency contract — see module
430        // docstring's "Structured Concurrency" section). Pre-completed
431        // tasks ignore `cancel` per `TaskScheduler::cancel` line 141
432        // ("Only cancel if still pending").
433        if let Some(scope) = self.async_scope_stack.last_mut() {
434            scope.push(task_id);
435        }
436
437        // Push the future id as `Ptr(HeapKind::Future)`. The Future kind
438        // is an inline-scalar payload — `bits` IS the future id, no Arc
439        // backing. Drop is a no-op in `drop_with_kind`.
440        self.push_kinded(task_id, NativeKind::Ptr(HeapKind::Future))?;
441        debug_assert_eq!(
442            self.sp, sp_before,
443            "op_spawn_task: stack depth changed (before={}, after={})",
444            sp_before, self.sp
445        );
446        Ok(AsyncExecutionResult::Continue)
447    }
448
449    /// Initialize a join group from futures on the stack
450    ///
451    /// Operand: Count(packed_u16) where high 2 bits = join kind, low 14 bits = arity.
452    /// Pops `arity` Future values from the stack (in reverse order).
453    /// Pushes a `Ptr(HeapKind::TaskGroup)`-kinded `Arc<TaskGroupData>` payload.
454    fn op_join_init(&mut self, instruction: &Instruction) -> Result<AsyncExecutionResult, VMError> {
455        let packed = match &instruction.operand {
456            Some(Operand::Count(n)) => *n,
457            _ => {
458                return Err(VMError::RuntimeError(
459                    "JoinInit requires Count operand".to_string(),
460                ));
461            }
462        };
463
464        let kind = ((packed >> 14) & 0x03) as u8;
465        let arity = (packed & 0x3FFF) as usize;
466
467        if self.sp < arity {
468            return Err(VMError::StackUnderflow);
469        }
470
471        let mut task_ids: Vec<u64> = Vec::with_capacity(arity);
472        for _ in 0..arity {
473            let (bits, slot_kind) = self.pop_kinded()?;
474            match slot_kind {
475                NativeKind::Ptr(HeapKind::Future) => {
476                    // Future is an inline scalar — bits IS the id. No share
477                    // to drop (HeapKind::Future is a no-op in drop_with_kind).
478                    task_ids.push(bits);
479                }
480                _ => {
481                    // Type mismatch — drop the popped share before surfacing
482                    // the error so refcount discipline holds (playbook §3).
483                    drop_with_kind(bits, slot_kind);
484                    return Err(VMError::RuntimeError(format!(
485                        "JoinInit expected Future, got {:?}",
486                        slot_kind
487                    )));
488                }
489            }
490        }
491        // Reverse so task_ids[0] corresponds to first branch
492        task_ids.reverse();
493
494        // Construct an Arc<TaskGroupData> and push as Ptr(HeapKind::TaskGroup).
495        // ADR-006 §2.3 / playbook §3 per-HeapKind push pattern: heap-bearing
496        // kinds push the `Arc::into_raw` pointer with the matching kind.
497        let arc: Arc<TaskGroupData> = Arc::new(TaskGroupData { kind, task_ids });
498        let bits = Arc::into_raw(arc) as u64;
499        self.push_kinded(bits, NativeKind::Ptr(HeapKind::TaskGroup))?;
500        Ok(AsyncExecutionResult::Continue)
501    }
502
503    /// Await a task group, resolving tasks inline
504    ///
505    /// Pops a `Ptr(HeapKind::TaskGroup)`-kinded slot from the stack.
506    /// Resolves all tasks inline using the task scheduler's `resolve_task_group`,
507    /// which executes each task's callable synchronously (same strategy as `op_await`).
508    /// Pushes the result value onto the stack according to the join strategy.
509    fn op_join_await(&mut self) -> Result<AsyncExecutionResult, VMError> {
510        let sp_before = self.sp;
511        let (bits, slot_kind) = self.pop_kinded()?;
512        match slot_kind {
513            NativeKind::Ptr(HeapKind::TaskGroup) => {
514                // Reclaim the `Arc<TaskGroupData>` share that `pop_kinded`
515                // transferred to us. Extract `kind` + `task_ids` for the
516                // join walk, then drop the Arc.
517                //
518                // SAFETY: the construction-side contract for
519                // `push_kinded(bits, Ptr(HeapKind::TaskGroup))` (see
520                // `op_join_init` above + ADR-006 §2.3) guarantees `bits`
521                // is the result of `Arc::into_raw::<TaskGroupData>` and
522                // we own exactly one strong-count share.
523                let arc: Arc<TaskGroupData> =
524                    unsafe { Arc::from_raw(bits as *const TaskGroupData) };
525                let join_kind = arc.kind;
526                let task_ids = arc.task_ids.clone();
527                drop(arc);
528
529                // Per-id sync resolution via the §2.7.11/Q12 dispatch
530                // entry-point. `resolve_spawned_task` consults the
531                // scheduler's cached-result fast-path first, then takes
532                // the callable share and routes through
533                // `call_*_with_nb_args` family. The borrow shape here
534                // mirrors `resolve_spawned_task`'s own `take_callable` /
535                // `complete` cycle — no per-call closure capture of
536                // `&mut self.task_scheduler` is needed because the
537                // scheduler is consulted/mutated point-wise inside
538                // `resolve_spawned_task` itself.
539                //
540                // Phase-2c boundary (ADR-006 §2.7.4): if a constituent
541                // task's body suspends (`VMError::Suspended`), the join
542                // surfaces the error upward. Snapshot-tier resumption
543                // of in-flight join groups stays out of scope per
544                // §2.7.11 out-of-scope clause.
545                match join_kind {
546                    // All: resolve every task, drop each per-task share
547                    // (the aggregate carrier is a `TaskGroupData` of
548                    // ids only, mirroring `TaskScheduler::resolve_task_group`'s
549                    // All-mode shape). Push a fresh `Arc<TaskGroupData>`
550                    // result carrier kinded `Ptr(HeapKind::TaskGroup)`.
551                    0 => {
552                        for &id in &task_ids {
553                            let result = self.resolve_spawned_task(id)?;
554                            drop_with_kind(result.raw(), result.kind());
555                            std::mem::forget(result);
556                        }
557                        let aggregate: Arc<TaskGroupData> = Arc::new(TaskGroupData {
558                            kind: 0,
559                            task_ids: task_ids.clone(),
560                        });
561                        let result_bits = Arc::into_raw(aggregate) as u64;
562                        self.push_kinded(
563                            result_bits,
564                            NativeKind::Ptr(HeapKind::TaskGroup),
565                        )?;
566                    }
567                    // Race: resolve all tasks; return the first result.
568                    // Matches `TaskScheduler::resolve_task_group`'s
569                    // race-mode semantics. Empty list → RuntimeError.
570                    1 => {
571                        let mut pushed = false;
572                        for (idx, &id) in task_ids.iter().enumerate() {
573                            let result = self.resolve_spawned_task(id)?;
574                            if idx == 0 {
575                                self.push_kinded(result.raw(), result.kind())?;
576                                std::mem::forget(result);
577                                pushed = true;
578                            } else {
579                                // Subsequent results: their shares aren't
580                                // returned to the user; release each.
581                                drop_with_kind(result.raw(), result.kind());
582                                std::mem::forget(result);
583                            }
584                        }
585                        if !pushed {
586                            return Err(VMError::RuntimeError(
587                                "Race join with empty task list".to_string(),
588                            ));
589                        }
590                    }
591                    // Any: return first success; on errors, keep the
592                    // last for the empty-success fallback. Matches
593                    // `TaskScheduler::resolve_task_group`'s any-mode.
594                    2 => {
595                        let mut last_err: Option<VMError> = None;
596                        let mut pushed = false;
597                        for &id in &task_ids {
598                            match self.resolve_spawned_task(id) {
599                                Ok(result) => {
600                                    self.push_kinded(result.raw(), result.kind())?;
601                                    std::mem::forget(result);
602                                    pushed = true;
603                                    break;
604                                }
605                                Err(e) => last_err = Some(e),
606                            }
607                        }
608                        if !pushed {
609                            return Err(last_err.unwrap_or_else(|| {
610                                VMError::RuntimeError(
611                                    "Any join with empty task list".to_string(),
612                                )
613                            }));
614                        }
615                    }
616                    // AllSettled: drive every task; per-task errors are
617                    // preserved in the scheduler's result map (caller
618                    // can inspect via `get_result`). Aggregate carrier
619                    // kind=3 mirrors `TaskScheduler::resolve_task_group`.
620                    // The {status, value/error} array view depends on
621                    // a kinded VMArray helper that's Phase-2c per
622                    // ADR-006 §2.7.4 — the TaskGroup carrier is the
623                    // minimum shape the await-time decoder can re-walk.
624                    3 => {
625                        for &id in &task_ids {
626                            if let Ok(result) = self.resolve_spawned_task(id) {
627                                drop_with_kind(result.raw(), result.kind());
628                                std::mem::forget(result);
629                            }
630                            // Errors per-task are preserved in the
631                            // scheduler's result map.
632                        }
633                        let aggregate: Arc<TaskGroupData> = Arc::new(TaskGroupData {
634                            kind: 3,
635                            task_ids: task_ids.clone(),
636                        });
637                        let result_bits = Arc::into_raw(aggregate) as u64;
638                        self.push_kinded(
639                            result_bits,
640                            NativeKind::Ptr(HeapKind::TaskGroup),
641                        )?;
642                    }
643                    other => {
644                        return Err(VMError::RuntimeError(format!(
645                            "Unknown join kind: {}",
646                            other
647                        )));
648                    }
649                }
650
651                debug_assert_eq!(
652                    self.sp, sp_before,
653                    "op_join_await: stack depth changed (before={}, after={})",
654                    sp_before, self.sp
655                );
656                Ok(AsyncExecutionResult::Continue)
657            }
658            _ => {
659                drop_with_kind(bits, slot_kind);
660                Err(VMError::RuntimeError(format!(
661                    "JoinAwait expected TaskGroup, got {:?}",
662                    slot_kind
663                )))
664            }
665        }
666    }
667
668    /// Cancel a task by its future ID
669    ///
670    /// Pops a Future(task_id) from the stack and signals cancellation.
671    /// The host runtime is responsible for actually cancelling the task.
672    fn op_cancel_task(&mut self) -> Result<AsyncExecutionResult, VMError> {
673        let (bits, slot_kind) = self.pop_kinded()?;
674        match slot_kind {
675            NativeKind::Ptr(HeapKind::Future) => {
676                // Future is an inline scalar — bits IS the id. No Arc share
677                // to drop (Future is a no-op in drop_with_kind).
678                let id = bits;
679                self.task_scheduler.cancel(id);
680                Ok(AsyncExecutionResult::Continue)
681            }
682            _ => {
683                drop_with_kind(bits, slot_kind);
684                Err(VMError::RuntimeError(format!(
685                    "CancelTask expected Future, got {:?}",
686                    slot_kind
687                )))
688            }
689        }
690    }
691
692    /// Enter a structured concurrency scope
693    ///
694    /// Pushes a new empty Vec onto the async_scope_stack.
695    /// All tasks spawned while this scope is active are tracked in that Vec.
696    fn op_async_scope_enter(&mut self) -> Result<AsyncExecutionResult, VMError> {
697        let depth_before = self.async_scope_stack.len();
698        self.async_scope_stack.push(Vec::new());
699        debug_assert_eq!(
700            self.async_scope_stack.len(),
701            depth_before + 1,
702            "op_async_scope_enter: scope stack depth not incremented"
703        );
704        Ok(AsyncExecutionResult::Continue)
705    }
706
707    /// Exit a structured concurrency scope
708    ///
709    /// Pops the current scope from the async_scope_stack and cancels
710    /// all tasks spawned within it that are still pending, in LIFO order.
711    /// The body's result value remains on top of the stack.
712    fn op_async_scope_exit(&mut self) -> Result<AsyncExecutionResult, VMError> {
713        debug_assert!(
714            !self.async_scope_stack.is_empty(),
715            "op_async_scope_exit: scope stack is empty (mismatched Enter/Exit)"
716        );
717        if let Some(mut scope_tasks) = self.async_scope_stack.pop() {
718            // Cancel in LIFO order (last spawned first)
719            scope_tasks.reverse();
720            for task_id in scope_tasks {
721                self.task_scheduler.cancel(task_id);
722            }
723        }
724        // Result value from the body is already on top of the stack
725        Ok(AsyncExecutionResult::Continue)
726    }
727
728    /// Emit a generic event to the event queue
729    ///
730    /// Pops an event object from the stack and pushes it to
731    /// the event queue for external consumers.
732    fn op_emit_event(&mut self) -> Result<AsyncExecutionResult, VMError> {
733        // Pop the event payload and release its share — event queue
734        // integration is deferred. Drop discipline (playbook §3): every
735        // `pop_kinded` either re-pushes or `drop_with_kind`s.
736        let (bits, kind) = self.pop_kinded()?;
737        drop_with_kind(bits, kind);
738        Ok(AsyncExecutionResult::Continue)
739    }
740}
741
742/// Check if an opcode is an async operation
743#[cfg(test)]
744pub fn is_async_opcode(opcode: OpCode) -> bool {
745    matches!(
746        opcode,
747        OpCode::Yield
748            | OpCode::Suspend
749            | OpCode::Resume
750            | OpCode::Poll
751            | OpCode::AwaitBar
752            | OpCode::AwaitTick
753            | OpCode::EmitAlert
754            | OpCode::EmitEvent
755            | OpCode::Await
756            | OpCode::SpawnTask
757            | OpCode::JoinInit
758            | OpCode::JoinAwait
759            | OpCode::CancelTask
760            | OpCode::AsyncScopeEnter
761            | OpCode::AsyncScopeExit
762    )
763}
764
765#[cfg(test)]
766mod tests {
767    use super::*;
768
769    #[test]
770    fn test_is_async_opcode() {
771        assert!(is_async_opcode(OpCode::Yield));
772        assert!(is_async_opcode(OpCode::Suspend));
773        assert!(is_async_opcode(OpCode::EmitAlert));
774        assert!(is_async_opcode(OpCode::AsyncScopeEnter));
775        assert!(is_async_opcode(OpCode::AsyncScopeExit));
776        assert!(!is_async_opcode(OpCode::AddInt));
777        assert!(!is_async_opcode(OpCode::Jump));
778    }
779
780    #[test]
781    fn test_is_async_opcode_all_variants() {
782        // Test all async opcodes
783        assert!(is_async_opcode(OpCode::Yield));
784        assert!(is_async_opcode(OpCode::Suspend));
785        assert!(is_async_opcode(OpCode::Resume));
786        assert!(is_async_opcode(OpCode::Poll));
787        assert!(is_async_opcode(OpCode::AwaitBar));
788        assert!(is_async_opcode(OpCode::AwaitTick));
789        assert!(is_async_opcode(OpCode::EmitAlert));
790        assert!(is_async_opcode(OpCode::EmitEvent));
791
792        // Test non-async opcodes
793        assert!(!is_async_opcode(OpCode::PushConst));
794        assert!(!is_async_opcode(OpCode::Return));
795        assert!(!is_async_opcode(OpCode::Call));
796        assert!(!is_async_opcode(OpCode::Nop));
797    }
798
799    #[test]
800    fn test_async_execution_result_variants() {
801        // Test Continue
802        let continue_result = AsyncExecutionResult::Continue;
803        assert!(matches!(continue_result, AsyncExecutionResult::Continue));
804
805        // Test Yielded
806        let yielded_result = AsyncExecutionResult::Yielded;
807        assert!(matches!(yielded_result, AsyncExecutionResult::Yielded));
808
809        // Test Suspended
810        let suspended_result = AsyncExecutionResult::Suspended(SuspensionInfo {
811            wait_type: WaitType::AnyEvent,
812            resume_ip: 42,
813        });
814        match suspended_result {
815            AsyncExecutionResult::Suspended(info) => {
816                assert_eq!(info.resume_ip, 42);
817                assert!(matches!(info.wait_type, WaitType::AnyEvent));
818            }
819            _ => panic!("Expected Suspended"),
820        }
821    }
822
823    #[test]
824    fn test_wait_type_variants() {
825        // NextBar
826        let next_bar = WaitType::NextBar {
827            source: "market_data".to_string(),
828        };
829        match next_bar {
830            WaitType::NextBar { source } => assert_eq!(source, "market_data"),
831            _ => panic!("Expected NextBar"),
832        }
833
834        // Timer
835        let timer = WaitType::Timer { id: 123 };
836        match timer {
837            WaitType::Timer { id } => assert_eq!(id, 123),
838            _ => panic!("Expected Timer"),
839        }
840
841        // AnyEvent
842        let any = WaitType::AnyEvent;
843        assert!(matches!(any, WaitType::AnyEvent));
844    }
845
846    #[test]
847    fn test_suspension_info_creation() {
848        let info = SuspensionInfo {
849            wait_type: WaitType::Timer { id: 999 },
850            resume_ip: 100,
851        };
852
853        assert_eq!(info.resume_ip, 100);
854        assert!(matches!(info.wait_type, WaitType::Timer { id: 999 }));
855    }
856
857    #[test]
858    fn test_is_async_opcode_await() {
859        assert!(is_async_opcode(OpCode::Await));
860    }
861
862    #[test]
863    fn test_wait_type_future() {
864        let future = WaitType::Future { id: 42 };
865        match future {
866            WaitType::Future { id } => assert_eq!(id, 42),
867            _ => panic!("Expected Future"),
868        }
869    }
870
871    #[test]
872    fn test_is_async_opcode_join_opcodes() {
873        assert!(is_async_opcode(OpCode::SpawnTask));
874        assert!(is_async_opcode(OpCode::JoinInit));
875        assert!(is_async_opcode(OpCode::JoinAwait));
876        assert!(is_async_opcode(OpCode::CancelTask));
877    }
878
879    #[test]
880    fn test_wait_type_task_group() {
881        let tg = WaitType::TaskGroup {
882            kind: 0,
883            task_ids: vec![1, 2, 3],
884        };
885        match tg {
886            WaitType::TaskGroup { kind, task_ids } => {
887                assert_eq!(kind, 0); // All
888                assert_eq!(task_ids.len(), 3);
889                assert_eq!(task_ids, vec![1, 2, 3]);
890            }
891            _ => panic!("Expected TaskGroup"),
892        }
893    }
894
895    #[test]
896    fn test_wait_type_task_group_race() {
897        let tg = WaitType::TaskGroup {
898            kind: 1,
899            task_ids: vec![10, 20],
900        };
901        match tg {
902            WaitType::TaskGroup { kind, task_ids } => {
903                assert_eq!(kind, 1); // Race
904                assert_eq!(task_ids, vec![10, 20]);
905            }
906            _ => panic!("Expected TaskGroup"),
907        }
908    }
909}