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}