Skip to main content

wasmtime/runtime/component/
concurrent.rs

1//! Runtime support for the Component Model Async ABI.
2//!
3//! This module and its submodules provide host runtime support for Component
4//! Model Async features such as async-lifted exports, async-lowered imports,
5//! streams, futures, and related intrinsics.  See [the Async
6//! Explainer](https://github.com/WebAssembly/component-model/blob/main/design/mvp/Concurrency.md)
7//! for a high-level overview.
8//!
9//! At the core of this support is an event loop which schedules and switches
10//! between guest tasks and any host tasks they create.  Each
11//! `Store` will have at most one event loop running at any given
12//! time, and that loop may be suspended and resumed by the host embedder using
13//! e.g. `StoreContextMut::run_concurrent`.  The `StoreContextMut::poll_until`
14//! function contains the loop itself, while the
15//! `StoreOpaque::concurrent_state` field holds its state.
16//!
17//! # Public API Overview
18//!
19//! ## Top-level API (e.g. kicking off host->guest calls and driving the event loop)
20//!
21//! - `[Typed]Func::call_concurrent`: Start a host->guest call to an
22//! async-lifted or sync-lifted import, creating a guest task.
23//!
24//! - `StoreContextMut::run_concurrent`: Run the event loop for the specified
25//! instance, allowing any and all tasks belonging to that instance to make
26//! progress.
27//!
28//! - `StoreContextMut::spawn`: Run a background task as part of the event loop
29//! for the specified instance.
30//!
31//! - `{Future,Stream}Reader::new`: Create a new Component Model `future` or
32//! `stream` which may be passed to the guest.  This takes a
33//! `{Future,Stream}Producer` implementation which will be polled for items when
34//! the consumer requests them.
35//!
36//! - `{Future,Stream}Reader::pipe`: Consume a `future` or `stream` by
37//! connecting it to a `{Future,Stream}Consumer` which will consume any items
38//! produced by the write end.
39//!
40//! ## Host Task API (e.g. implementing concurrent host functions and background tasks)
41//!
42//! - `LinkerInstance::func_wrap_concurrent`: Register a concurrent host
43//! function with the linker.  That function will take an `Accessor` as its
44//! first parameter, which provides access to the store between (but not across)
45//! await points.
46//!
47//! - `Accessor::with`: Access the store and its associated data.
48//!
49//! - `Accessor::spawn`: Run a background task as part of the event loop for the
50//! store.  This is equivalent to `StoreContextMut::spawn` but more convenient to use
51//! in host functions.
52
53use self::error_contexts::GlobalErrorContextRefCount;
54use crate::component::func::{Func, call_post_return};
55use crate::component::{
56    HasData, HasSelf, Instance, Resource, ResourceTable, ResourceTableError, RuntimeInstance,
57};
58use crate::fiber::{self, StoreFiber, StoreFiberYield};
59use crate::hash_set::HashSet;
60#[cfg(feature = "gc")]
61use crate::module::ModuleRegistry;
62use crate::prelude::*;
63use crate::store::{Store, StoreId, StoreInner, StoreOpaque, StoreToken};
64#[cfg(feature = "gc")]
65use crate::vm::GcRootsList;
66use crate::vm::component::{CallContext, ComponentInstance, CurrentScope, InstanceState, Scope};
67use crate::vm::{AlwaysMut, SendSyncPtr, VMFuncRef, VMLazyThread, VMMemoryDefinition, VMStore};
68use crate::{
69    AsContext, AsContextMut, FuncType, Result, StoreContext, StoreContextMut, ValRaw, ValType, bail,
70};
71use crate::{Instance as ModuleInstance, bail_bug};
72use alloc::borrow::ToOwned;
73use alloc::collections::{BTreeMap, BTreeSet, VecDeque};
74use core::any::Any;
75use core::cell::UnsafeCell;
76use core::fmt;
77use core::future;
78use core::future::Future;
79use core::marker::PhantomData;
80use core::mem::{self, ManuallyDrop, MaybeUninit};
81use core::ops::DerefMut;
82use core::pin::{Pin, pin};
83use core::ptr::{self, NonNull};
84use core::task::{Context, Poll, Waker};
85use futures::channel::oneshot;
86use futures::stream::{FuturesUnordered, StreamExt};
87use futures_and_streams::{FlatAbi, ReturnCode, TransmitHandle, TransmitIndex};
88use table::{TableDebug, TableId};
89use wasmtime_environ::component::{
90    CanonicalAbiInfo, CanonicalOptions, CanonicalOptionsDataModel, MAX_FLAT_PARAMS,
91    MAX_FLAT_RESULTS, OptionsIndex, PREPARE_ASYNC_NO_RESULT, PREPARE_ASYNC_WITH_RESULT,
92    RuntimeComponentInstanceIndex, RuntimeTableIndex, StringEncoding,
93    TypeComponentGlobalErrorContextTableIndex, TypeComponentLocalErrorContextTableIndex,
94    TypeFuncIndex, TypeFutureTableIndex, TypeStreamTableIndex, TypeTupleIndex,
95};
96use wasmtime_environ::packed_option::ReservedValue;
97use wasmtime_environ::{NUM_COMPONENT_CONTEXT_SLOTS, Trap};
98#[cfg(feature = "gc")]
99use wasmtime_unwinder::Unwind;
100
101pub use abort::JoinHandle;
102pub use func::{FuncCallConcurrent, TypedFuncCallConcurrent};
103pub use future_stream_any::{FutureAny, StreamAny};
104pub use futures_and_streams::{
105    Destination, DirectDestination, DirectSource, ErrorContext, FutureConsumer, FutureProducer,
106    FutureReader, GuardedFutureReader, GuardedStreamReader, ReadBuffer, Source, StreamConsumer,
107    StreamProducer, StreamReader, StreamResult, VecBuffer, WriteBuffer,
108};
109pub(crate) use futures_and_streams::{ResourcePair, lower_error_context_to_index};
110#[cfg(feature = "task-group-hook")]
111pub use task_group_hook::TaskGroupHook;
112pub use task_group_hook::TaskGroupId;
113
114mod abort;
115mod error_contexts;
116mod func;
117mod future_stream_any;
118mod futures_and_streams;
119pub(crate) mod table;
120#[cfg(feature = "task-group-hook")]
121mod task_group_hook;
122#[cfg(not(feature = "task-group-hook"))]
123mod task_group_hook_disabled;
124#[cfg(not(feature = "task-group-hook"))]
125use task_group_hook_disabled as task_group_hook;
126pub(crate) mod tls;
127
128/// Constant defined in the Component Model spec to indicate that the async
129/// intrinsic (e.g. `future.write`) has not yet completed.
130const BLOCKED: u32 = 0xffff_ffff;
131
132/// Corresponds to `CallState` in the upstream spec.
133#[derive(Clone, Copy, Eq, PartialEq, Debug)]
134pub enum Status {
135    Starting = 0,
136    Started = 1,
137    Returned = 2,
138    StartCancelled = 3,
139    ReturnCancelled = 4,
140}
141
142impl Status {
143    /// Packs this status and the optional `waitable` provided into a 32-bit
144    /// result that the canonical ABI requires.
145    ///
146    /// The low 4 bits are reserved for the status while the upper 28 bits are
147    /// the waitable, if present.
148    pub fn pack(self, waitable: Option<u32>) -> u32 {
149        assert!(matches!(self, Status::Returned) == waitable.is_none());
150        let waitable = waitable.unwrap_or(0);
151        assert!(waitable < (1 << 28));
152        (waitable << 4) | (self as u32)
153    }
154}
155
156/// Corresponds to `EventCode` in the Component Model spec, plus related payload
157/// data.
158#[derive(Clone, Copy, Debug)]
159enum Event {
160    None,
161    Subtask {
162        status: Status,
163    },
164    StreamRead {
165        code: ReturnCode,
166        pending: Option<(TypeStreamTableIndex, u32)>,
167    },
168    StreamWrite {
169        code: ReturnCode,
170        pending: Option<(TypeStreamTableIndex, u32)>,
171    },
172    FutureRead {
173        code: ReturnCode,
174        pending: Option<(TypeFutureTableIndex, u32)>,
175    },
176    FutureWrite {
177        code: ReturnCode,
178        pending: Option<(TypeFutureTableIndex, u32)>,
179    },
180    Cancelled,
181}
182
183impl Event {
184    /// Lower this event to core Wasm integers for delivery to the guest.
185    ///
186    /// Note that the waitable handle, if any, is assumed to be lowered
187    /// separately.
188    fn parts(self) -> (u32, u32) {
189        const EVENT_NONE: u32 = 0;
190        const EVENT_SUBTASK: u32 = 1;
191        const EVENT_STREAM_READ: u32 = 2;
192        const EVENT_STREAM_WRITE: u32 = 3;
193        const EVENT_FUTURE_READ: u32 = 4;
194        const EVENT_FUTURE_WRITE: u32 = 5;
195        const EVENT_CANCELLED: u32 = 6;
196        match self {
197            Event::None => (EVENT_NONE, 0),
198            Event::Cancelled => (EVENT_CANCELLED, 0),
199            Event::Subtask { status } => (EVENT_SUBTASK, status as u32),
200            Event::StreamRead { code, .. } => (EVENT_STREAM_READ, code.encode()),
201            Event::StreamWrite { code, .. } => (EVENT_STREAM_WRITE, code.encode()),
202            Event::FutureRead { code, .. } => (EVENT_FUTURE_READ, code.encode()),
203            Event::FutureWrite { code, .. } => (EVENT_FUTURE_WRITE, code.encode()),
204        }
205    }
206}
207
208/// Corresponds to `CallbackCode` in the spec.
209mod callback_code {
210    pub const EXIT: u32 = 0;
211    pub const YIELD: u32 = 1;
212    pub const WAIT: u32 = 2;
213}
214
215/// A flag indicating that the callee is an async-lowered export.
216///
217/// This may be passed to the `async-start` intrinsic from a fused adapter.
218const START_FLAG_ASYNC_CALLEE: u32 = wasmtime_environ::component::START_FLAG_ASYNC_CALLEE as u32;
219
220/// Provides access to either store data (via the `get` method) or the store
221/// itself (via [`AsContext`]/[`AsContextMut`]), as well as the component
222/// instance to which the current host task belongs.
223///
224/// See [`Accessor::with`] for details.
225pub struct Access<'a, T: 'static, D: HasData + ?Sized = HasSelf<T>> {
226    store: StoreContextMut<'a, T>,
227    get_data: fn(&mut T) -> D::Data<'_>,
228}
229
230impl<'a, T, D> Access<'a, T, D>
231where
232    D: HasData + ?Sized,
233    T: 'static,
234{
235    /// Creates a new [`Access`] from its component parts.
236    pub fn new(store: StoreContextMut<'a, T>, get_data: fn(&mut T) -> D::Data<'_>) -> Self {
237        Self { store, get_data }
238    }
239
240    /// Get mutable access to the store data.
241    pub fn data_mut(&mut self) -> &mut T {
242        self.store.data_mut()
243    }
244
245    /// Get mutable access to the store data.
246    pub fn get(&mut self) -> D::Data<'_> {
247        (self.get_data)(self.data_mut())
248    }
249
250    /// Spawn a background task.
251    ///
252    /// See [`Accessor::spawn`] for details.
253    pub fn spawn(&mut self, task: impl for<'fut> AccessorTask<'fut, T, D>) -> Result<JoinHandle>
254    where
255        T: 'static,
256    {
257        let accessor = Accessor {
258            get_data: self.get_data,
259            token: StoreToken::new(self.store.as_context_mut()),
260        };
261        self.store
262            .as_context_mut()
263            .spawn_with_accessor(accessor, task)
264    }
265
266    /// Returns the getter this accessor is using to project from `T` into
267    /// `D::Data`.
268    pub fn getter(&self) -> fn(&mut T) -> D::Data<'_> {
269        self.get_data
270    }
271}
272
273impl<'a, T, D> AsContext for Access<'a, T, D>
274where
275    D: HasData + ?Sized,
276    T: 'static,
277{
278    type Data = T;
279
280    fn as_context(&self) -> StoreContext<'_, T> {
281        self.store.as_context()
282    }
283}
284
285impl<'a, T, D> AsContextMut for Access<'a, T, D>
286where
287    D: HasData + ?Sized,
288    T: 'static,
289{
290    fn as_context_mut(&mut self) -> StoreContextMut<'_, T> {
291        self.store.as_context_mut()
292    }
293}
294
295/// Provides scoped mutable access to store data in the context of a concurrent
296/// host task future.
297///
298/// This allows multiple host task futures to execute concurrently and access
299/// the store between (but not across) `await` points.
300///
301/// # Rationale
302///
303/// This structure is sort of like `&mut T` plus a projection from `&mut T` to
304/// `D::Data<'_>`. The problem this is solving, however, is that it does not
305/// literally store these values. The basic problem is that when a concurrent
306/// host future is being polled it has access to `&mut T` (and the whole
307/// `Store`) but when it's not being polled it does not have access to these
308/// values. This reflects how the store is only ever polling one future at a
309/// time so the store is effectively being passed between futures.
310///
311/// Rust's `Future` trait, however, has no means of passing a `Store`
312/// temporarily between futures. The [`Context`](core::task::Context) type does
313/// not have the ability to attach arbitrary information to it at this time.
314/// This type, [`Accessor`], is used to bridge this expressivity gap.
315///
316/// The [`Accessor`] type here represents the ability to acquire, temporarily in
317/// a synchronous manner, the current store. The [`Accessor::with`] function
318/// yields an [`Access`] which can be used to access [`StoreContextMut`], `&mut
319/// T`, or `D::Data<'_>`. Note though that [`Accessor::with`] intentionally does
320/// not take an `async` closure as its argument, instead it's a synchronous
321/// closure which must complete during on run of `Future::poll`. This reflects
322/// how the store is temporarily made available while a host future is being
323/// polled.
324///
325/// # Implementation
326///
327/// This type does not actually store `&mut T` nor `StoreContextMut<T>`, and
328/// this type additionally doesn't even have a lifetime parameter. This is
329/// instead a representation of proof of the ability to acquire these while a
330/// future is being polled. Wasmtime will, when it polls a host future,
331/// configure ambient state such that the `Accessor` that a future closes over
332/// will work and be able to access the store.
333///
334/// This has a number of implications for users such as:
335///
336/// * It's intentional that `Accessor` cannot be cloned, it needs to stay within
337///   the lifetime of a single future.
338/// * A future is expected to, however, close over an `Accessor` and keep it
339///   alive probably for the duration of the entire future.
340/// * Different host futures will be given different `Accessor`s, and that's
341///   intentional.
342/// * The `Accessor` type is `Send` and `Sync` irrespective of `T` which
343///   alleviates some otherwise required bounds to be written down.
344///
345/// # Using `Accessor` in `Drop`
346///
347/// The methods on `Accessor` are only expected to work in the context of
348/// `Future::poll` and are not guaranteed to work in `Drop`. This is because a
349/// host future can be dropped at any time throughout the system and Wasmtime
350/// store context is not necessarily available at that time. It's recommended to
351/// not use `Accessor` methods in anything connected to a `Drop` implementation
352/// as they will panic and have unintended results. If you run into this though
353/// feel free to file an issue on the Wasmtime repository.
354pub struct Accessor<T: 'static, D = HasSelf<T>>
355where
356    D: HasData + ?Sized,
357{
358    token: StoreToken<T>,
359    get_data: fn(&mut T) -> D::Data<'_>,
360}
361
362/// A helper trait to take any type of accessor-with-data in functions.
363///
364/// This trait is similar to [`AsContextMut`] except that it's used when
365/// working with an [`Accessor`] instead of a [`StoreContextMut`]. The
366/// [`Accessor`] is the main type used in concurrent settings and is passed to
367/// functions such as [`Func::call_concurrent`].
368///
369/// This trait is implemented for [`Accessor`] and `&T` where `T` implements
370/// this trait. This effectively means that regardless of the `D` in
371/// `Accessor<T, D>` it can still be passed to a function which just needs a
372/// store accessor.
373///
374/// Acquiring an [`Accessor`] can be done through
375/// [`StoreContextMut::run_concurrent`] for example or in a host function
376/// through
377/// [`Linker::func_wrap_concurrent`](crate::component::LinkerInstance::func_wrap_concurrent).
378pub trait AsAccessor {
379    /// The `T` in `Store<T>` that this accessor refers to.
380    type Data: 'static;
381
382    /// The `D` in `Accessor<T, D>`, or the projection out of
383    /// `Self::Data`.
384    type AccessorData: HasData + ?Sized;
385
386    /// Returns the accessor that this is referring to.
387    fn as_accessor(&self) -> &Accessor<Self::Data, Self::AccessorData>;
388}
389
390impl<T: AsAccessor + ?Sized> AsAccessor for &T {
391    type Data = T::Data;
392    type AccessorData = T::AccessorData;
393
394    fn as_accessor(&self) -> &Accessor<Self::Data, Self::AccessorData> {
395        T::as_accessor(self)
396    }
397}
398
399impl<T, D: HasData + ?Sized> AsAccessor for Accessor<T, D> {
400    type Data = T;
401    type AccessorData = D;
402
403    fn as_accessor(&self) -> &Accessor<T, D> {
404        self
405    }
406}
407
408// Note that it is intentional at this time that `Accessor` does not actually
409// store `&mut T` or anything similar. This distinctly enables the `Accessor`
410// structure to be both `Send` and `Sync` regardless of what `T` is (or `D` for
411// that matter). This is used to ergonomically simplify bindings where the
412// majority of the time `Accessor` is closed over in a future which then needs
413// to be `Send` and `Sync`. To avoid needing to write `T: Send` everywhere (as
414// you already have to write `T: 'static`...) it helps to avoid this.
415//
416// Note as well that `Accessor` doesn't actually store its data at all. Instead
417// it's more of a "proof" of what can be accessed from TLS. API design around
418// `Accessor` and functions like `Linker::func_wrap_concurrent` are
419// intentionally made to ensure that `Accessor` is ideally only used in the
420// context that TLS variables are actually set. For example host functions are
421// given `&Accessor`, not `Accessor`, and this prevents them from persisting
422// the value outside of a future. Within the future the TLS variables are all
423// guaranteed to be set while the future is being polled.
424//
425// Finally though this is not an ironclad guarantee, but nor does it need to be.
426// The TLS APIs are designed to panic or otherwise model usage where they're
427// called recursively or similar. It's hoped that code cannot be constructed to
428// actually hit this at runtime but this is not a safety requirement at this
429// time.
430const _: () = {
431    const fn assert<T: Send + Sync>() {}
432    assert::<Accessor<UnsafeCell<u32>>>();
433};
434
435impl<T> Accessor<T> {
436    /// Creates a new `Accessor` backed by the specified functions.
437    ///
438    /// - `get`: used to retrieve the store
439    ///
440    /// - `get_data`: used to "project" from the store's associated data to
441    /// another type (e.g. a field of that data or a wrapper around it).
442    ///
443    /// - `spawn`: used to queue spawned background tasks to be run later
444    pub(crate) fn new(token: StoreToken<T>) -> Self {
445        Self {
446            token,
447            get_data: |x| x,
448        }
449    }
450}
451
452impl<T, D> Accessor<T, D>
453where
454    D: HasData + ?Sized,
455{
456    /// Run the specified closure, passing it mutable access to the store.
457    ///
458    /// This function is one of the main building blocks of the [`Accessor`]
459    /// type. This yields synchronous, blocking, access to the store via an
460    /// [`Access`]. The [`Access`] implements [`AsContextMut`] in addition to
461    /// providing the ability to access `D` via [`Access::get`]. Note that the
462    /// `fun` here is given only temporary access to the store and `T`/`D`
463    /// meaning that the return value `R` here is not allowed to capture borrows
464    /// into the two. If access is needed to data within `T` or `D` outside of
465    /// this closure then it must be `clone`d out, for example.
466    ///
467    /// # Panics
468    ///
469    /// This function will panic if it is call recursively with any other
470    /// accessor already in scope. For example if `with` is called within `fun`,
471    /// then this function will panic. It is up to the embedder to ensure that
472    /// this does not happen.
473    pub fn with<R>(&self, fun: impl FnOnce(Access<'_, T, D>) -> R) -> R {
474        tls::get(|vmstore| {
475            fun(Access {
476                store: self.token.as_context_mut(vmstore),
477                get_data: self.get_data,
478            })
479        })
480    }
481
482    /// Returns the getter this accessor is using to project from `T` into
483    /// `D::Data`.
484    pub fn getter(&self) -> fn(&mut T) -> D::Data<'_> {
485        self.get_data
486    }
487
488    /// Changes this accessor to access `D2` instead of the current type
489    /// parameter `D`.
490    ///
491    /// This changes the underlying data access from `T` to `D2::Data<'_>`.
492    ///
493    /// # Panics
494    ///
495    /// When using this API the returned value is disconnected from `&self` and
496    /// the lifetime binding the `self` argument. An `Accessor` only works
497    /// within the context of the closure or async closure that it was
498    /// originally given to, however. This means that due to the fact that the
499    /// returned value has no lifetime connection it's possible to use the
500    /// accessor outside of `&self`, the original accessor, and panic.
501    ///
502    /// The returned value should only be used within the scope of the original
503    /// `Accessor` that `self` refers to.
504    pub fn with_getter<D2: HasData>(
505        &self,
506        get_data: fn(&mut T) -> D2::Data<'_>,
507    ) -> Accessor<T, D2> {
508        Accessor {
509            token: self.token,
510            get_data,
511        }
512    }
513
514    /// Spawn a background task which will receive an `&Accessor<T, D>` and
515    /// run concurrently with any other tasks in progress for the current
516    /// store.
517    ///
518    /// This is particularly useful for host functions which return a `stream`
519    /// or `future` such that the code to write to the write end of that
520    /// `stream` or `future` must run after the function returns.
521    ///
522    /// The returned [`JoinHandle`] may be used to cancel the task.
523    ///
524    /// # Panics
525    ///
526    /// Panics if called within a closure provided to the [`Accessor::with`]
527    /// function. This can only be called outside an active invocation of
528    /// [`Accessor::with`].
529    pub fn spawn(&self, task: impl for<'fut> AccessorTask<'fut, T, D>) -> Result<JoinHandle>
530    where
531        T: 'static,
532    {
533        let accessor = self.clone_for_spawn();
534        self.with(|mut access| access.as_context_mut().spawn_with_accessor(accessor, task))
535    }
536
537    fn clone_for_spawn(&self) -> Self {
538        Self {
539            token: self.token,
540            get_data: self.get_data,
541        }
542    }
543
544    /// Polls to see if this store contains any "interesting" tasks still within
545    /// it.
546    ///
547    /// Returns `Poll::Ready(())` if there are no more interesting tasks, and
548    /// otherwise returns `Poll::Pending`. If pending is returned then whenever
549    /// the last remaining "interesting" task has exited the provided context's
550    /// waker will be notified. Note that only the waker passed to the last call
551    /// to `poll_no_interesting_tasks` for the store will be notified, so this
552    /// is only appropriate to use once-at-a-time per store.
553    ///
554    /// The component model specification, as of this current date, does not
555    /// have a distinction between "interesting" tasks and not. The current
556    /// intention is that in a future revision of the component model this will
557    /// be distinguished at the component ABI level where tasks will be able to
558    /// flag themselves as "interesting" optionally. Additionally extra work can
559    /// be opted-in to being "interesting".
560    ///
561    /// For now what this means is that all component model tasks within this
562    /// store are considered interesting. This specifically includes the entire
563    /// duration of a task, so even all of the time after a task has returned
564    /// but before it has exited. This means that this function is, today,
565    /// effectively a proxy for "are there any more tasks still running in this
566    /// store". This can be used by embedders to determine whether there's any
567    /// more work going on, even in the background, for a particular guest.
568    /// Hosts can use this as a signal that the guest wants to stay alive a
569    /// little longer, even after a task has returned.
570    ///
571    /// In the future this predicate won't include all tasks in this store. Some
572    /// tasks will be able to flag themselves as not interesting, meaning that
573    /// when this returns ready it'd be possible that there are still tasks
574    /// remaining in the store.
575    ///
576    /// Note that at this time spawned threads within a task are always
577    /// considered uninteresting. If this function returns ready, then spawned
578    /// threads may still be in the store.
579    pub fn poll_no_interesting_tasks(&self, cx: &mut Context<'_>) -> Poll<()> {
580        self.with(|mut access| {
581            let store = access.as_context_mut().0;
582            let state = store.concurrent_state_mut_without_forcing_current_thread();
583            if state.interesting_tasks == 0 {
584                Poll::Ready(())
585            } else {
586                state.interesting_tasks_empty_waker = Some(cx.waker().clone());
587                Poll::Pending
588            }
589        })
590    }
591
592    /// Poll to see if the component instance corresponding to the specified
593    /// function is ready to run a concurrent call without queuing it (i.e. does
594    /// not have backpressure enabled and does not have a sync call in
595    /// progress).
596    ///
597    /// Returns `Poll::Ready(())` if the component instance is ready to run a
598    /// concurrent call, and otherwise returns `Poll::Pending`.  If pending is
599    /// returned then whenever the instance becomes ready for a call the
600    /// provided context's waker will be notified.  Note that only the waker
601    /// passed to the last call to `poll_ready_for_concurrent_call` for the
602    /// store will be notified (regardless of whether the same or different
603    /// `Func` is specified relative to earlier calls), so this is only
604    /// appropriate to use once-at-a-time per store.  Also note that the waker
605    /// may be notified when _any_ instance becomes callable (i.e. not
606    /// necessarily the last one polled), so this function must be called again
607    /// to determine if the instance of interest is ready.
608    pub fn poll_ready_for_concurrent_call(&self, func: Func, cx: &mut Context<'_>) -> Poll<()> {
609        self.with(|mut access| {
610            let store = access.as_context_mut().0;
611            let (_, _, _, raw_options) = func.abi_info(store);
612            let instance = func.instance().runtime_instance(raw_options.instance);
613            let state = store.instance_state(instance).concurrent_state();
614            if state.backpressure == 0 {
615                Poll::Ready(())
616            } else {
617                store
618                    .concurrent_state_mut_without_forcing_current_thread()
619                    .ready_for_concurrent_call_waker = Some(cx.waker().clone());
620                Poll::Pending
621            }
622        })
623    }
624}
625
626/// Represents an async closure which may be provided to `Accessor::spawn`,
627/// `Accessor::forward`, or `StoreContextMut::spawn`.
628// TODO: Replace this with `core::ops::AsyncFnOnce` when we are able to put `Send`
629// bound on the unnamed `Future` directly.
630//
631// As of this writing, it's not possible to specify e.g. `Send` and `Sync`
632// bounds on the `Future` type returned by an `AsyncFnOnce`.  Also, using `F:
633// Future<Output = Result<()>> + Send + Sync, FN: FnOnce(&Accessor<T>) -> F +
634// Send + Sync + 'static` fails with a type mismatch error as we cannot describe
635// that `F` should have `&Accessor<T>`'s unnamed lifetime
636//
637// Instead, this trait is used as a workaround for this limitation, the bound on `Self`
638// implementing `AsyncFnOnce()` is required for Rust to automatically infer that an async
639// closure is to be provided wherever we are accepting `impl for<'fut> AccessorTask<'fut, T, D>`
640// as argument. Otherwise, users will have to fully qualify the closure types before it
641// is accepted as a valid value. This also means that it is not intended for a user
642// to manually implement this trait for any arbitrary type, since they would first
643// have to implement `AsyncFnOnce`, which is unstable.
644//
645// The blanket implementation for this trait will ensure that `Self` is an async closure
646// that returns a `Future` that is `Send`, and lives as long as `&Accessor<T, D>`
647pub trait AccessorTask<'fut, T, D = HasSelf<T>>:
648    AsyncFnOnce(&Accessor<T, D>) -> Result<()> + Send + 'static
649where
650    D: HasData + ?Sized,
651{
652    /// Run the task.
653    fn run(self, accessor: &'fut Accessor<T, D>) -> impl Future<Output = Result<()>> + Send + 'fut;
654}
655
656impl<'fut, F, Fut, T, D> AccessorTask<'fut, T, D> for F
657where
658    T: 'static,
659    F: AsyncFnOnce(&Accessor<T, D>) -> Result<()>,
660    F: FnOnce(&'fut Accessor<T, D>) -> Fut + Send + 'static,
661    Fut: Future<Output = Result<()>> + Send + 'fut,
662    D: HasData,
663{
664    fn run(self, accessor: &'fut Accessor<T, D>) -> impl Future<Output = Result<()>> + Send + 'fut {
665        (self)(accessor)
666    }
667}
668
669/// Represents parameter and result metadata for the caller side of a
670/// guest->guest call orchestrated by a fused adapter.
671enum CallerInfo {
672    /// Metadata for a call to an async-lowered import
673    Async {
674        params: Vec<ValRaw>,
675        has_result: bool,
676    },
677    /// Metadata for a call to an sync-lowered import
678    Sync {
679        params: Vec<ValRaw>,
680        result_count: u32,
681    },
682}
683
684/// Indicates how a guest task is waiting on a waitable set.
685enum WaitMode {
686    /// The guest task is waiting using `task.wait`
687    Fiber(StoreFiber<'static>),
688    /// The guest task is waiting via a callback declared as part of an
689    /// async-lifted export.
690    Callback(Instance),
691}
692
693impl fmt::Debug for WaitMode {
694    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
695        match self {
696            Self::Fiber(_) => f.debug_tuple("Fiber").finish(),
697            Self::Callback(instance) => f.debug_tuple("Callback").field(instance).finish(),
698        }
699    }
700}
701
702/// Represents the reason a fiber is suspending itself.
703#[derive(Debug)]
704enum SuspendReason {
705    /// The fiber is waiting for an event to be delivered to the specified
706    /// waitable set or task.
707    Waiting {
708        set: TableId<WaitableSet>,
709        thread: QualifiedThreadId,
710    },
711    /// The fiber is waiting for a subtask to suspend or exit, e.g. for a
712    /// guest-to-guest call or a `subtask.cancel`.
713    YieldingToSubtask { thread: QualifiedThreadId },
714    /// The fiber has finished handling its most recent work item and is waiting
715    /// for another (or to be dropped if it is no longer needed).
716    NeedWork,
717    /// The fiber is yielding and should be resumed once other tasks have had a
718    /// chance to run.
719    Yielding { thread: QualifiedThreadId },
720    /// The fiber was explicitly suspended with a call to `thread.suspend` or
721    /// `thread.switch-to`.
722    ExplicitlySuspending { thread: QualifiedThreadId },
723}
724
725/// Represents a pending call into guest code for a given guest task.
726enum GuestCallKind {
727    /// Indicates there's an event to deliver to the task, possibly related to a
728    /// waitable set the task has been waiting on or polling.
729    DeliverEvent {
730        /// The instance to which the task belongs.
731        instance: Instance,
732        /// The waitable set the event belongs to, if any.
733        ///
734        /// If this is `None` the event will be waiting in the
735        /// `GuestTask::event` field for the task.
736        set: Option<TableId<WaitableSet>>,
737    },
738    /// Indicates that a new guest task call is pending and may be executed
739    /// using the specified closure.
740    ///
741    /// If the closure returns `Ok(Some(call))`, the `call` should be run
742    /// immediately using `handle_guest_call`.
743    StartImplicit(Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send + Sync>),
744    StartExplicit(Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send + Sync>),
745}
746
747impl fmt::Debug for GuestCallKind {
748    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
749        match self {
750            Self::DeliverEvent { instance, set } => f
751                .debug_struct("DeliverEvent")
752                .field("instance", instance)
753                .field("set", set)
754                .finish(),
755            Self::StartImplicit(_) => f.debug_tuple("StartImplicit").finish(),
756            Self::StartExplicit(_) => f.debug_tuple("StartExplicit").finish(),
757        }
758    }
759}
760
761/// The target of a suspension intrinsic.
762#[derive(Copy, Clone, Debug)]
763pub enum SuspensionTarget {
764    Resume(u32),
765    Promote(u32),
766    None,
767}
768
769/// Behavior for `resume_thread`.
770#[derive(Copy, Clone, Debug)]
771pub enum ResumeThread {
772    Promote,
773    Resume,
774    ResumeLater,
775}
776
777/// Represents a pending call into guest code for a given guest thread.
778#[derive(Debug)]
779struct GuestCall {
780    thread: QualifiedThreadId,
781    kind: GuestCallKind,
782}
783
784impl GuestCall {
785    /// Returns whether or not the call is ready to run.
786    ///
787    /// A call will not be ready to run if either:
788    ///
789    /// - the (sub-)component instance to be called has already been entered and
790    /// cannot be reentered until an in-progress call completes
791    ///
792    /// - the call is for a not-yet started task and the (sub-)component
793    /// instance to be called has backpressure enabled
794    fn is_ready(&self, store: &mut StoreOpaque) -> Result<bool> {
795        let task = store.concurrent_state_mut()?.get_mut(self.thread.task)?;
796        let async_typed = task.async_typed;
797        let instance = task.instance;
798        let state = store.instance_state(instance).concurrent_state();
799
800        let ready = match &self.kind {
801            GuestCallKind::DeliverEvent { .. } => !state.do_not_enter,
802            GuestCallKind::StartImplicit(_) => {
803                !async_typed || !(state.do_not_enter || state.backpressure > 0)
804            }
805            GuestCallKind::StartExplicit(_) => true,
806        };
807        log::trace!(
808            "call {self:?} ready? {ready} (do_not_enter: {}; backpressure: {})",
809            state.do_not_enter,
810            state.backpressure
811        );
812        Ok(ready)
813    }
814}
815
816/// Job to be run on a worker fiber.
817enum WorkerItem {
818    GuestCall(GuestCall),
819    Function(AlwaysMut<Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send>>),
820}
821
822/// Represents a pending work item to be handled by the event loop for a given
823/// component instance.
824enum WorkItem {
825    /// A host task to be pushed to `ConcurrentState::futures`.
826    PushFuture(AlwaysMut<HostTaskFuture>),
827    /// A fiber to resume.
828    ResumeFiber {
829        instance: RuntimeInstance,
830        thread: QualifiedThreadId,
831        fiber: StoreFiber<'static>,
832    },
833    /// A thread to resume.
834    ResumeThread {
835        instance: RuntimeInstance,
836        thread: QualifiedThreadId,
837    },
838    /// A pending call into guest code for a given guest task.
839    GuestCall {
840        instance: RuntimeInstance,
841        call: GuestCall,
842    },
843    /// A job to run on a worker fiber.
844    WorkerFunction(AlwaysMut<Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send>>),
845}
846
847impl fmt::Debug for WorkItem {
848    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
849        match self {
850            Self::PushFuture(_) => f.debug_tuple("PushFuture").finish(),
851            Self::ResumeFiber {
852                instance, thread, ..
853            } => f
854                .debug_struct("ResumeFiber")
855                .field("instance", instance)
856                .field("thread", thread)
857                .finish(),
858            Self::ResumeThread { instance, thread } => f
859                .debug_struct("ResumeThread")
860                .field("instance", instance)
861                .field("thread", thread)
862                .finish(),
863            Self::GuestCall { instance, call } => f
864                .debug_struct("GuestCall")
865                .field("instance", instance)
866                .field("call", call)
867                .finish(),
868            Self::WorkerFunction(_) => f.debug_tuple("WorkerFunction").finish(),
869        }
870    }
871}
872
873/// Whether a suspension intrinsic was cancelled or completed
874#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
875pub(crate) enum WaitResult {
876    Cancelled,
877    Completed,
878}
879
880/// Poll the specified future until it completes on behalf of a guest->host call
881/// using a sync-lowered import.
882///
883/// This is similar to `Instance::first_poll` except it's for sync-lowered
884/// imports, meaning we don't need to handle cancellation and we can block the
885/// caller until the task completes, at which point the caller can handle
886/// lowering the result to the guest's stack and linear memory.
887pub(crate) fn poll_and_block<R: Send + Sync + 'static>(
888    store: &mut dyn VMStore,
889    host_task: EnteredHostTask,
890    future: impl Future<Output = Result<R>> + Send + 'static,
891) -> Result<R> {
892    // Poll the future once before creating a host task. The host task will be
893    // created lazily if it's needed during the poll and otherwise will be
894    // created if the future suspends.  We can use a dummy `Waker` here because
895    // we'll add the future to `ConcurrentState::futures` and poll it
896    // automatically from the event loop if it doesn't complete immediately
897    // here.
898    let mut future = Box::pin(future);
899    let poll = tls::set(store, || {
900        future
901            .as_mut()
902            .poll(&mut Context::from_waker(&Waker::noop()))
903    });
904
905    let caller = match host_task {
906        Some(caller) => caller,
907        None => bail_bug!("host task wasn't created but should have been"),
908    };
909
910    let task = match poll {
911        // It completed immediately, so no persistent host task is needed.
912        Poll::Ready(result) => return result,
913
914        // It did not complete immediately; create the host task and add it to
915        // `ConcurrentState::futures` so it will be polled via the event loop;
916        // then use `GuestThread::sync_call_set` to wait for the task to
917        // complete, suspending the current fiber until it does so.
918        Poll::Pending => {
919            let Some(task) = store.materialize_host_task_id()? else {
920                bail_bug!("current thread is not a host thread")
921            };
922
923            // Wrap the future in a closure which will stash its result in the
924            // host task and resume this fiber when it completes.
925            let future = Box::pin(async move {
926                let result = run_with_host_task_set(task, future).await??;
927                tls::get(move |store| {
928                    let state = store.concurrent_state_mut()?;
929                    let host_state = &mut state.get_mut(task)?.state;
930                    assert!(matches!(host_state, HostTaskState::CalleeStarted));
931                    *host_state = HostTaskState::CalleeFinished(Box::new(result));
932
933                    Waitable::Host(task).set_event(
934                        state,
935                        Some(Event::Subtask {
936                            status: Status::Returned,
937                        }),
938                    )?;
939
940                    Ok(())
941                })
942            }) as HostTaskFuture;
943
944            let caller_instance = store.concurrent_state_mut()?.get_mut(caller.task)?.instance;
945            store.switch_or_trap_if_may_not_suspend(caller_instance)?;
946
947            let state = store.concurrent_state_mut()?;
948            state.push_future(future);
949
950            let set = state.get_mut(caller.thread)?.sync_call_set;
951            Waitable::Host(task).join(state, Some(set))?;
952
953            store.suspend(SuspendReason::Waiting {
954                set,
955                thread: caller,
956            })?;
957
958            // Remove the `task` from the `sync_call_set` to ensure that when
959            // this function returns and the task is deleted that there are no
960            // more lingering references to this host task.
961            Waitable::Host(task).join(store.concurrent_state_mut()?, None)?;
962            task
963        }
964    };
965
966    // Retrieve and return the result.
967    let host_state = &mut store.concurrent_state_mut()?.get_mut(task)?.state;
968    match mem::replace(host_state, HostTaskState::CalleeDone { cancelled: false }) {
969        HostTaskState::CalleeFinished(result) => Ok(match result.downcast() {
970            Ok(result) => *result,
971            Err(_) => bail_bug!("host task finished with wrong type of result"),
972        }),
973        _ => bail_bug!("unexpected host task state after completion"),
974    }
975}
976
977/// Execute the specified guest call.
978fn handle_guest_call(store: &mut dyn VMStore, call: GuestCall) -> Result<()> {
979    match call.kind {
980        GuestCallKind::DeliverEvent { instance, set } => {
981            let (event, waitable) = match instance.get_event(store, call.thread.task, set, true)? {
982                Some(pair) => pair,
983                None => bail_bug!("delivering non-present event"),
984            };
985            let state = store.concurrent_state_mut()?;
986            let task = state.get_mut(call.thread.task)?;
987            let runtime_instance = task.instance;
988            let handle = waitable.map(|(_, v)| v).unwrap_or(0);
989
990            log::trace!(
991                "use callback to deliver event {event:?} to {:?} for {waitable:?}",
992                call.thread,
993            );
994
995            let old_thread = store.set_thread(call.thread)?;
996            log::trace!(
997                "GuestCallKind::DeliverEvent: replaced {old_thread:?} with {:?} as current thread",
998                call.thread
999            );
1000
1001            store.enter_instance(runtime_instance);
1002
1003            let Some(callback) = store
1004                .concurrent_state_mut()?
1005                .get_mut(call.thread.task)?
1006                .callback
1007                .take()
1008            else {
1009                bail_bug!("guest task callback field not present")
1010            };
1011
1012            let code = callback(store, event, handle)?;
1013
1014            store
1015                .concurrent_state_mut()?
1016                .get_mut(call.thread.task)?
1017                .callback = Some(callback);
1018
1019            store.exit_instance(runtime_instance)?;
1020
1021            store.set_thread(old_thread)?;
1022
1023            instance.handle_callback_code(store, call.thread, runtime_instance.index, code)?;
1024
1025            log::trace!("GuestCallKind::DeliverEvent: restored {old_thread:?} as current thread");
1026        }
1027        GuestCallKind::StartImplicit(fun) => {
1028            fun(store)?;
1029        }
1030        GuestCallKind::StartExplicit(fun) => {
1031            fun(store)?;
1032        }
1033    }
1034
1035    Ok(())
1036}
1037
1038impl<T> Store<T> {
1039    /// Convenience wrapper for [`StoreContextMut::run_concurrent`].
1040    pub async fn run_concurrent<R>(&mut self, fun: impl AsyncFnOnce(&Accessor<T>) -> R) -> Result<R>
1041    where
1042        T: Send + 'static,
1043    {
1044        ensure!(
1045            self.as_context().0.concurrency_support(),
1046            "cannot use `run_concurrent` when Config::concurrency_support disabled",
1047        );
1048        self.as_context_mut().run_concurrent(fun).await
1049    }
1050
1051    #[doc(hidden)]
1052    pub fn assert_concurrent_state_empty(&mut self) {
1053        self.as_context_mut().assert_concurrent_state_empty();
1054    }
1055
1056    #[doc(hidden)]
1057    pub fn concurrent_state_table_size(&mut self) -> usize {
1058        self.as_context_mut().concurrent_state_table_size()
1059    }
1060
1061    /// Convenience wrapper for [`StoreContextMut::spawn`].
1062    pub fn spawn(
1063        &mut self,
1064        task: impl for<'fut> AccessorTask<'fut, T, HasSelf<T>>,
1065    ) -> Result<JoinHandle>
1066    where
1067        T: 'static,
1068    {
1069        self.as_context_mut().spawn(task)
1070    }
1071}
1072
1073impl<T> StoreContextMut<'_, T> {
1074    /// Assert that all the relevant tables and queues in the concurrent state
1075    /// for this store are empty.
1076    ///
1077    /// This is for sanity checking in integration tests
1078    /// (e.g. `component-async-tests`) that the relevant state has been cleared
1079    /// after each test concludes.  This should help us catch leaks, e.g. guest
1080    /// tasks which haven't been deleted despite having completed and having
1081    /// been dropped by their supertasks.
1082    ///
1083    /// Only intended for use in Wasmtime's own testing.
1084    #[doc(hidden)]
1085    pub fn assert_concurrent_state_empty(self) {
1086        let store = self.0;
1087        store
1088            .store_data_mut()
1089            .components
1090            .assert_instance_states_empty();
1091        let state = store.concurrent_state_mut().unwrap();
1092        assert!(
1093            state.table.get_mut().is_empty(),
1094            "non-empty table: {:?}",
1095            state.table.get_mut()
1096        );
1097        assert!(state.switch_item.is_none());
1098        assert!(state.next_switch_item.is_none());
1099        assert!(state.high_priority.is_empty());
1100        assert!(state.low_priority.is_empty());
1101        assert!(state.unforced_current_thread.is_none());
1102        assert!(state.deferred_host_call_context.is_none());
1103        assert!(state.futures_mut().unwrap().is_empty());
1104        assert!(state.global_error_context_ref_counts.is_empty());
1105    }
1106
1107    /// Helper function to perform tests over the size of the concurrent state
1108    /// table which can be useful for detecting leaks.
1109    ///
1110    /// Only intended for use in Wasmtime's own testing.
1111    #[doc(hidden)]
1112    pub fn concurrent_state_table_size(&mut self) -> usize {
1113        self.0
1114            .concurrent_state_mut()
1115            .unwrap()
1116            .table
1117            .get_mut()
1118            .iter_mut()
1119            .count()
1120    }
1121
1122    /// Spawn a background task to run as part of this instance's event loop.
1123    ///
1124    /// The task will receive an `&Accessor<U>` and run concurrently with
1125    /// any other tasks in progress for the instance.
1126    ///
1127    /// Note that the task will only make progress if and when the event loop
1128    /// for this instance is run.
1129    ///
1130    /// The returned [`JoinHandle`] may be used to cancel the task.
1131    pub fn spawn(mut self, task: impl for<'fut> AccessorTask<'fut, T>) -> Result<JoinHandle>
1132    where
1133        T: 'static,
1134    {
1135        let accessor = Accessor::new(StoreToken::new(self.as_context_mut()));
1136        self.spawn_with_accessor(accessor, task)
1137    }
1138
1139    /// Internal implementation of `spawn` functions where a `store` is
1140    /// available along with an `Accessor`.
1141    fn spawn_with_accessor<D>(
1142        self,
1143        accessor: Accessor<T, D>,
1144        task: impl for<'fut> AccessorTask<'fut, T, D>,
1145    ) -> Result<JoinHandle>
1146    where
1147        T: 'static,
1148        D: HasData + ?Sized,
1149    {
1150        // Create an "abortable future" here where internally the future will
1151        // hook calls to poll and possibly spawn more background tasks on each
1152        // iteration.
1153        let (handle, future) = JoinHandle::run(async move { task.run(&accessor).await });
1154        self.0
1155            .concurrent_state_mut()?
1156            .push_future(Box::pin(async move { future.await.unwrap_or(Ok(())) }));
1157        Ok(handle)
1158    }
1159
1160    /// Run the specified closure `fun` to completion as part of this store's
1161    /// event loop.
1162    ///
1163    /// This will run `fun` as part of this store's event loop until it
1164    /// yields a result.  `fun` is provided an [`Accessor`], which provides
1165    /// controlled access to the store and its data.
1166    ///
1167    /// This function can be used to invoke [`Func::call_concurrent`] for
1168    /// example within the async closure provided here.
1169    ///
1170    /// This function will unconditionally return an error if
1171    /// [`Config::concurrency_support`] is disabled.
1172    ///
1173    /// [`Config::concurrency_support`]: crate::Config::concurrency_support
1174    ///
1175    /// # Store-blocking behavior
1176    ///
1177    /// At this time there are certain situations in which the `Future` returned
1178    /// by the `AsyncFnOnce` passed to this function will not be polled for an
1179    /// extended period of time, despite one or more `Waker::wake` events having
1180    /// occurred for the task to which it belongs.  This can manifest as the
1181    /// `Future` seeming to be "blocked" or "locked up", but is actually due to
1182    /// the `Store` being held by e.g. a blocking host function, preventing the
1183    /// `Future` from being polled. A canonical example of this is when the
1184    /// `fun` provided to this function attempts to set a timeout for an
1185    /// invocation of a wasm function. In this situation the async closure is
1186    /// waiting both on (a) the wasm computation to finish, and (b) the timeout
1187    /// to elapse. At this time this setup will not always work and the timeout
1188    /// may not reliably fire.
1189    ///
1190    /// This function will not block the current thread and as such is always
1191    /// suitable to run in an `async` context, but the current implementation of
1192    /// Wasmtime can lead to situations where a certain wasm computation is
1193    /// required to make progress the closure to make progress. This is an
1194    /// artifact of Wasmtime's historical implementation of `async` functions
1195    /// and is the topic of [#11869] and [#11870]. In the timeout example from
1196    /// above it means that Wasmtime can get "wedged" for a bit where (a) must
1197    /// progress for a readiness notification of (b) to get delivered.
1198    ///
1199    /// This effectively means that it's not possible to reliably perform a
1200    /// "select" operation within the `fun` closure, which timeouts for example
1201    /// are based on. Fixing this requires some relatively major refactoring
1202    /// work within Wasmtime itself. This is a known pitfall otherwise and one
1203    /// that is intended to be fixed one day. In the meantime it's recommended
1204    /// to apply timeouts or such to the entire `run_concurrent` call itself
1205    /// rather than internally.
1206    ///
1207    /// [#11869]: https://github.com/bytecodealliance/wasmtime/issues/11869
1208    /// [#11870]: https://github.com/bytecodealliance/wasmtime/issues/11870
1209    ///
1210    /// # Example
1211    ///
1212    /// ```
1213    /// # use {
1214    /// #   wasmtime::{
1215    /// #     error::{Result},
1216    /// #     component::{ Component, Linker, Resource, ResourceTable},
1217    /// #     Config, Engine, Store
1218    /// #   },
1219    /// # };
1220    /// #
1221    /// # struct MyResource(u32);
1222    /// # struct Ctx { table: ResourceTable }
1223    /// #
1224    /// # async fn foo() -> Result<()> {
1225    /// # let mut config = Config::new();
1226    /// # let engine = Engine::new(&config)?;
1227    /// # let mut store = Store::new(&engine, Ctx { table: ResourceTable::new() });
1228    /// # let mut linker = Linker::new(&engine);
1229    /// # let component = Component::new(&engine, "")?;
1230    /// # let instance = linker.instantiate_async(&mut store, &component).await?;
1231    /// # let foo = instance.get_typed_func::<(Resource<MyResource>,), (Resource<MyResource>,)>(&mut store, "foo")?;
1232    /// # let bar = instance.get_typed_func::<(u32,), ()>(&mut store, "bar")?;
1233    /// store.run_concurrent(async |accessor| -> wasmtime::Result<_> {
1234    ///    let resource = accessor.with(|mut access| access.get().table.push(MyResource(42)))?;
1235    ///    let (another_resource,) = foo.call_concurrent(accessor, (resource,)).await?;
1236    ///    let value = accessor.with(|mut access| access.get().table.delete(another_resource))?;
1237    ///    bar.call_concurrent(accessor, (value.0,)).await?;
1238    ///    Ok(())
1239    /// }).await??;
1240    /// # Ok(())
1241    /// # }
1242    /// ```
1243    pub async fn run_concurrent<R>(self, fun: impl AsyncFnOnce(&Accessor<T>) -> R) -> Result<R>
1244    where
1245        T: Send + 'static,
1246    {
1247        ensure!(
1248            self.0.concurrency_support(),
1249            "cannot use `run_concurrent` when Config::concurrency_support disabled",
1250        );
1251        self.do_run_concurrent(fun, false).await
1252    }
1253
1254    pub(super) async fn run_concurrent_trap_on_idle<R>(
1255        self,
1256        fun: impl AsyncFnOnce(&Accessor<T>) -> R,
1257    ) -> Result<R> {
1258        self.do_run_concurrent(fun, true).await
1259    }
1260
1261    async fn do_run_concurrent<R>(
1262        mut self,
1263        fun: impl AsyncFnOnce(&Accessor<T>) -> R,
1264        trap_on_idle: bool,
1265    ) -> Result<R> {
1266        debug_assert!(self.0.concurrency_support());
1267        let already_running = self
1268            .0
1269            .concurrent_state_mut_already_forced_current_thread()
1270            .event_loop_running;
1271        if already_running {
1272            bail!("Recursive `StoreContextMut::run_concurrent` calls not supported")
1273        }
1274        let token = StoreToken::new(self.as_context_mut());
1275
1276        struct Dropper<'a, T: 'static, V> {
1277            store: StoreContextMut<'a, T>,
1278            value: ManuallyDrop<V>,
1279        }
1280
1281        impl<'a, T, V> Drop for Dropper<'a, T, V> {
1282            fn drop(&mut self) {
1283                self.store
1284                    .0
1285                    .concurrent_state_mut_already_forced_current_thread()
1286                    .event_loop_running = false;
1287
1288                tls::set(self.store.0, || {
1289                    // SAFETY: Here we drop the value without moving it for the
1290                    // first and only time -- per the contract for `Drop::drop`,
1291                    // this code won't run again, and the `value` field will no
1292                    // longer be accessible.
1293                    unsafe { ManuallyDrop::drop(&mut self.value) }
1294                });
1295            }
1296        }
1297
1298        let accessor = &Accessor::new(token);
1299        self.0
1300            .concurrent_state_mut_already_forced_current_thread()
1301            .event_loop_running = true;
1302        let dropper = &mut Dropper {
1303            store: self,
1304            value: ManuallyDrop::new(fun(accessor)),
1305        };
1306        // SAFETY: We never move `dropper` nor its `value` field.
1307        let future = unsafe { Pin::new_unchecked(dropper.value.deref_mut()) };
1308
1309        let result = dropper
1310            .store
1311            .as_context_mut()
1312            .poll_until(future, trap_on_idle)
1313            .await;
1314
1315        if result.is_err() {
1316            dropper.store.0.set_trapped();
1317        }
1318
1319        result
1320    }
1321
1322    /// Run this store's event loop.
1323    ///
1324    /// The returned future will resolve when the specified future completes or,
1325    /// if `trap_on_idle` is true, when the event loop can't make further
1326    /// progress.
1327    async fn poll_until<R>(
1328        mut self,
1329        mut future: Pin<&mut impl Future<Output = R>>,
1330        trap_on_idle: bool,
1331    ) -> Result<R> {
1332        struct Reset<'a, T: 'static> {
1333            store: StoreContextMut<'a, T>,
1334            futures: Option<FuturesUnordered<HostTaskFuture>>,
1335        }
1336
1337        impl<'a, T> Drop for Reset<'a, T> {
1338            fn drop(&mut self) {
1339                if let Some(futures) = self.futures.take() {
1340                    *self
1341                        .store
1342                        .0
1343                        .concurrent_state_mut_already_forced_current_thread()
1344                        .futures
1345                        .get_mut() = Some(futures);
1346                }
1347            }
1348        }
1349
1350        loop {
1351            // Take `ConcurrentState::futures` out of the store so we can poll
1352            // it while also safely giving any of the futures inside access to
1353            // `self`.
1354            let futures = self.0.concurrent_state_mut()?.futures.get_mut().take();
1355            let mut reset = Reset {
1356                store: self.as_context_mut(),
1357                futures,
1358            };
1359            let mut next = match reset.futures.as_mut() {
1360                Some(f) => pin!(f.next()),
1361                None => bail_bug!("concurrent state missing futures field"),
1362            };
1363
1364            enum PollResult<R> {
1365                Complete(R),
1366                ProcessWork {
1367                    ready: Option<WorkItem>,
1368                    low_priority: bool,
1369                },
1370            }
1371
1372            let result = future::poll_fn(|cx| {
1373                // First, poll the future we were passed as an argument and
1374                // return immediately if it's ready.
1375                if let Poll::Ready(value) = tls::set(reset.store.0, || future.as_mut().poll(cx)) {
1376                    return Poll::Ready(Ok(PollResult::Complete(value)));
1377                }
1378
1379                // Next, poll `ConcurrentState::futures` (which includes any
1380                // pending host tasks and/or background tasks), returning
1381                // immediately if one of them fails.
1382                let next = match tls::set(reset.store.0, || next.as_mut().poll(cx)) {
1383                    Poll::Ready(Some(output)) => {
1384                        match output {
1385                            Err(e) => return Poll::Ready(Err(e)),
1386                            Ok(()) => {}
1387                        }
1388                        Poll::Ready(true)
1389                    }
1390                    Poll::Ready(None) => Poll::Ready(false),
1391                    Poll::Pending => Poll::Pending,
1392                };
1393
1394                // Next, identify the next work item to process, if any, using
1395                // the following priority order:
1396                //
1397                // - `switch_item`: Represents the guest thread we _must_ switch
1398                // to before any other thread runs per the determinism
1399                // requirements in the Component Model spec.
1400                //
1401                // - `high_priority`: "Urgent" work items, e.g. async calls have
1402                // become freshly unblocked due to backpressure clearing or
1403                // similar, stream or future state updates, etc.
1404                //
1405                // - `low_priority`: Work items such as resuming a fiber after
1406                // it yields, in which case the point is to let other items run
1407                // first.
1408                let state = reset.store.0.concurrent_state_mut()?;
1409                let mut ready = state.switch_item.take();
1410                let mut low_priority = false;
1411                if ready.is_none() {
1412                    ready = state.high_priority.pop_back();
1413                    if ready.is_none() {
1414                        ready = state.low_priority.pop_back();
1415                        low_priority = true;
1416                    }
1417                }
1418                if ready.is_some() {
1419                    return Poll::Ready(Ok(PollResult::ProcessWork {
1420                        ready,
1421                        low_priority,
1422                    }));
1423                }
1424
1425                // Finally, if we have nothing else to do right now, determine what to do
1426                // based on whether there are any pending futures in
1427                // `ConcurrentState::futures`.
1428                return match next {
1429                    Poll::Ready(true) => {
1430                        // In this case, one of the futures in
1431                        // `ConcurrentState::futures` completed
1432                        // successfully, so we return now and continue
1433                        // the outer loop in case there is another one
1434                        // ready to complete.
1435                        Poll::Ready(Ok(PollResult::ProcessWork {
1436                            ready: None,
1437                            low_priority: false,
1438                        }))
1439                    }
1440                    Poll::Ready(false) => {
1441                        // Poll the future we were passed one last time
1442                        // in case one of `ConcurrentState::futures` had
1443                        // the side effect of unblocking it.
1444                        if let Poll::Ready(value) =
1445                            tls::set(reset.store.0, || future.as_mut().poll(cx))
1446                        {
1447                            Poll::Ready(Ok(PollResult::Complete(value)))
1448                        } else {
1449                            // In this case, there are no more pending
1450                            // futures in `ConcurrentState::futures`,
1451                            // there are no remaining work items, _and_
1452                            // the future we were passed as an argument
1453                            // still hasn't completed.
1454                            if trap_on_idle {
1455                                // `trap_on_idle` is true, so we exit
1456                                // immediately.
1457
1458                                // If there are any tasks belonging to
1459                                // an instance which may not suspend, trap with
1460                                // `CannotBlockSyncTask`:
1461                                Poll::Ready(Err(if reset.store.0.any_may_not_suspend()? {
1462                                    Trap::CannotBlockSyncTask.into()
1463                                } else {
1464                                    // Otherwise, trap with `AsyncDeadlock`:
1465                                    Trap::AsyncDeadlock.into()
1466                                }))
1467                            } else {
1468                                // `trap_on_idle` is false, so we assume
1469                                // that future will wake up and give us
1470                                // more work to do when it's ready to.
1471                                Poll::Pending
1472                            }
1473                        }
1474                    }
1475                    // There is at least one pending future in
1476                    // `ConcurrentState::futures` and we have nothing
1477                    // else to do but wait for now, so we return
1478                    // `Pending`.
1479                    Poll::Pending => Poll::Pending,
1480                };
1481            })
1482            .await;
1483
1484            // Put the `ConcurrentState::futures` back into the store before we
1485            // return or handle any work items since one or more of those items
1486            // might append more futures.
1487            drop(reset);
1488
1489            match result? {
1490                // The future we were passed as an argument completed, so we
1491                // return the result.
1492                PollResult::Complete(value) => break Ok(value),
1493                // The future we were passed has not yet completed, so handle
1494                // any work items and then loop again.
1495                PollResult::ProcessWork {
1496                    ready,
1497                    low_priority,
1498                } => {
1499                    struct Dispose<'a, T: 'static> {
1500                        store: StoreContextMut<'a, T>,
1501                        ready: Option<WorkItem>,
1502                    }
1503
1504                    impl<'a, T> Drop for Dispose<'a, T> {
1505                        fn drop(&mut self) {
1506                            if let Some(item) = self.ready.take() {
1507                                match item {
1508                                    WorkItem::ResumeFiber { mut fiber, .. } => {
1509                                        fiber.dispose(self.store.0)
1510                                    }
1511                                    WorkItem::PushFuture(future) => {
1512                                        tls::set(self.store.0, move || drop(future))
1513                                    }
1514                                    _ => {}
1515                                }
1516                            }
1517                        }
1518                    }
1519
1520                    let mut dispose = Dispose {
1521                        store: self.as_context_mut(),
1522                        ready,
1523                    };
1524
1525                    // If we're about to run a low-priority task, first yield to
1526                    // the executor.  This ensures that it won't be starved of
1527                    // the ability to e.g. update the readiness of sockets,
1528                    // etc. which the guest may be using `thread.yield` along
1529                    // with `waitable-set.poll` to monitor in a CPU-heavy loop.
1530                    //
1531                    // This works for e.g. `thread.yield` and callbacks which
1532                    // return `CALLBACK_CODE_YIELD` because we queue a low
1533                    // priority item to resume the task (i.e. resume the thread
1534                    // or call the callback, respectively) just prior to
1535                    // suspending it.  Indeed, as of this writing those are the
1536                    // _only_ situations we queue low-priority tasks.
1537                    // Therefore, we interpret the guest's request to yield as
1538                    // meaning "yield to other guest tasks _and_/_or_ host
1539                    // operations such as updating socket readiness", the latter
1540                    // being the async runtime's responsibility.
1541                    //
1542                    // In the future, if this ends up causing measurable
1543                    // performance issues, this could be optimized such that we
1544                    // only yield periodically (e.g. for batches of low priority
1545                    // items) and not for each and every idividual item.
1546                    if low_priority {
1547                        dispose.store.0.yield_now().await
1548                    }
1549
1550                    if let Some(item) = dispose.ready.take() {
1551                        dispose
1552                            .store
1553                            .as_context_mut()
1554                            .handle_work_item(item)
1555                            .await?;
1556                    }
1557                }
1558            }
1559        }
1560    }
1561
1562    /// Handle the specified work item, possibly resuming a fiber if applicable.
1563    async fn handle_work_item(self, item: WorkItem) -> Result<()> {
1564        log::trace!("handle work item {item:?}");
1565        match item {
1566            WorkItem::PushFuture(future) => {
1567                self.0
1568                    .concurrent_state_mut()?
1569                    .futures_mut()?
1570                    .push(future.into_inner());
1571            }
1572            WorkItem::ResumeFiber { fiber, .. } => {
1573                self.0.resume_fiber(fiber).await?;
1574            }
1575            WorkItem::ResumeThread { thread, .. } => {
1576                if let GuestThreadState::Ready { fiber, .. } = mem::replace(
1577                    &mut self.0.concurrent_state_mut()?.get_mut(thread.thread)?.state,
1578                    GuestThreadState::Running,
1579                ) {
1580                    self.0.resume_fiber(fiber).await?;
1581                } else {
1582                    bail_bug!("cannot resume non-pending thread {thread:?}");
1583                }
1584            }
1585            WorkItem::GuestCall { call, .. } => {
1586                if call.is_ready(self.0)? {
1587                    self.0
1588                        .concurrent_state_mut()?
1589                        .get_mut(call.thread.thread)?
1590                        .wake_on_cancel = WakeOnCancel::None;
1591                    self.run_on_worker(WorkerItem::GuestCall(call)).await?;
1592                } else {
1593                    let state = self.0.concurrent_state_mut()?;
1594                    let task = state.get_mut(call.thread.task)?;
1595                    if !task.starting_sent {
1596                        task.starting_sent = true;
1597                        if let GuestCallKind::StartImplicit(_) = &call.kind {
1598                            Waitable::Guest(call.thread.task).set_event(
1599                                state,
1600                                Some(Event::Subtask {
1601                                    status: Status::Starting,
1602                                }),
1603                            )?;
1604                        }
1605                    }
1606
1607                    let instance = state.get_mut(call.thread.task)?.instance;
1608                    self.0
1609                        .instance_state(instance)
1610                        .concurrent_state()
1611                        .pending
1612                        .insert(call.thread, call.kind);
1613
1614                    // Switch back to the caller (or canceller) immediately if
1615                    // applicable since we aren't yet able to run the subtask it
1616                    // yielded to.
1617                    self.0.concurrent_state_mut()?.take_next_switch_item()?;
1618                }
1619            }
1620            WorkItem::WorkerFunction(fun) => {
1621                self.run_on_worker(WorkerItem::Function(fun)).await?;
1622            }
1623        }
1624
1625        Ok(())
1626    }
1627
1628    /// Execute the specified guest call on a worker fiber.
1629    async fn run_on_worker(self, item: WorkerItem) -> Result<()> {
1630        let worker = if let Some(fiber) = self.0.concurrent_state_mut()?.worker.take() {
1631            fiber
1632        } else {
1633            // SAFETY: the `make_fiber_unchecked` function is unsafe because the
1634            // returned fiber is unconditionally `Send` as opposed to being
1635            // conditionally send depending on the argument (in this case
1636            // `self.0`). This `async` function, however, is conditionally
1637            // `Send` depending on `self`, in this case `StoreContextMut<T>`,
1638            // which is already going to be conditionally `Send` depending on
1639            // `T`.
1640            //
1641            // The returned fiber is possibly stored within the `Store<T>` as
1642            // well. If `T: Send` then that's fine and everything's dandy. If
1643            // `T: !Send`, however, then the store is already not-`Send` meaning
1644            // that putting more actually-not-`Send` things inside of it isn't
1645            // an issue.
1646            //
1647            // The main issue here is that the returned fiber effectively can't
1648            // get transferred outside the context of the store. That's an
1649            // implementation detail we'll have to rely on, but is currently
1650            // true.
1651            unsafe {
1652                fiber::make_fiber_unchecked(self.0, move |store| {
1653                    loop {
1654                        let Some(item) = store.concurrent_state_mut()?.worker_item.take() else {
1655                            bail_bug!("worker_item not present when resuming fiber")
1656                        };
1657                        match item {
1658                            WorkerItem::GuestCall(call) => handle_guest_call(store, call)?,
1659                            WorkerItem::Function(fun) => fun.into_inner()(store)?,
1660                        }
1661
1662                        store.suspend(SuspendReason::NeedWork)?;
1663                    }
1664                })?
1665            }
1666        };
1667
1668        let worker_item = &mut self.0.concurrent_state_mut()?.worker_item;
1669        assert!(worker_item.is_none());
1670        *worker_item = Some(item);
1671
1672        self.0.resume_fiber(worker).await
1673    }
1674
1675    /// Wrap the specified host function in a future which will call it, passing
1676    /// it an `&Accessor<T>`.
1677    ///
1678    /// See the `Accessor` documentation for details.
1679    pub(crate) fn wrap_call<F, R>(self, closure: F) -> impl Future<Output = Result<R>> + 'static
1680    where
1681        T: 'static,
1682        F: FnOnce(&Accessor<T>) -> Pin<Box<dyn Future<Output = Result<R>> + Send + '_>>
1683            + Send
1684            + Sync
1685            + 'static,
1686        R: Send + Sync + 'static,
1687    {
1688        let token = StoreToken::new(self);
1689        async move {
1690            let mut accessor = Accessor::new(token);
1691            closure(&mut accessor).await
1692        }
1693    }
1694
1695    pub(crate) async fn start_instance(
1696        &mut self,
1697        instance: ModuleInstance,
1698    ) -> Result<ModuleInstance> {
1699        let (tx, rx) = oneshot::channel();
1700        let token = StoreToken::new(self.as_context_mut());
1701        self.0.queue_task(move |store| {
1702            _ = tx.send(
1703                instance
1704                    .start_raw(&mut token.as_context_mut(store))
1705                    .map(|()| instance),
1706            );
1707            Ok(())
1708        })?;
1709        self.as_context_mut()
1710            .run_concurrent_trap_on_idle(async |_| {
1711                rx.await
1712                    .map_err(|_| format_err!("oneshot channel canceled"))
1713            })
1714            .await??
1715    }
1716}
1717
1718/// Return value of [`StoreOpaque::host_task_create`].
1719///
1720/// This is an `Option` to handle the dynamic `store.concurrency_support()`
1721/// property. When present this records the guest thread to restore when the
1722/// host call exits. The corresponding [`HostTask`] will need to be lazily
1723/// created if needed via [`StoreOpaque::materialize_host_task_id`].
1724pub type EnteredHostTask = Option<QualifiedThreadId>;
1725
1726impl StoreOpaque {
1727    /// Returns the currently-running thread, promoting any deferred lazy guest
1728    /// thread into a fully-materialized `CurrentThread`. Deferred [`HostTask`]s
1729    /// are not materialized.
1730    #[inline]
1731    pub(crate) fn current_thread(&mut self) -> Result<CurrentThread> {
1732        // Without concurrency support there is nothing to force.
1733        if !self.concurrency_support() {
1734            return Ok(CurrentThread::None);
1735        }
1736
1737        // If the JIT-visible current thread isn't a deferred thread then
1738        // `ConcurrentState` is already up to date.
1739        if !self
1740            .vm_store_context_mut()
1741            .current_thread_mut()
1742            .is_deferred()
1743        {
1744            return Ok(self
1745                .concurrent_state_mut_already_forced_current_thread()
1746                .unforced_current_thread);
1747        }
1748
1749        self.force_deferred_current_thread()
1750    }
1751
1752    /// Slow path of [`Self::current_thread`]: promote the deferred lazy
1753    /// thread into a fully-materialized `CurrentThread`.
1754    #[cold]
1755    fn force_deferred_current_thread(&mut self) -> Result<CurrentThread> {
1756        // The component instance whose adapters pushed the deferred frames; all
1757        // frames in a guest-to-guest, sync-to-sync call chain of fused adapters
1758        // live within a single `wasmtime::component::Instance` (because
1759        // cross-`wasmtime::component::Instance` calls don't go through fused
1760        // adapters), and guest code only ever runs as a guest thread, so the
1761        // chain's base thread is already materialized in `ConcurrentState` and
1762        // we can get the `ComponentInstanceId` shared by the whole chain from
1763        // here.
1764        let state = self.concurrent_state_mut_without_forcing_current_thread();
1765        let id = match state.unforced_current_thread.guest_task() {
1766            Some(task) => state.get_mut(task)?.instance.instance,
1767            None => bail_bug!("deferred component-model thread with non-guest base"),
1768        };
1769
1770        // Collect the deferred frames pushed inline by fused adapters, walking
1771        // the `parent` chain from innermost to the base.
1772        let mut frames = Vec::new();
1773        let mut cur = *self.vm_store_context_mut().current_thread_mut();
1774        while let Some(ptr) = cur.as_deferred() {
1775            // SAFETY: `ptr` points at a `VMDeferredThread` living in a fused
1776            // adapter's stack frame that is suspended below us on the stack
1777            // (mid-call, waiting for this nested call to return), so the
1778            // referent is still valid and exclusively ours to read.
1779            let deferred = unsafe { ptr.as_non_null().as_ref() };
1780            frames.push((
1781                deferred.callee_async != 0,
1782                deferred.callee_instance,
1783                deferred.saved_context,
1784            ));
1785            cur = deferred.parent;
1786        }
1787
1788        // Mark the current thread forced *before* replaying so that any
1789        // reentrant `force_current_thread` call short-circuits via the
1790        // non-deferred path above.
1791        *self.vm_store_context_mut().current_thread_mut() = VMLazyThread::forced();
1792
1793        // Save the current context, as we need to overwrite it while replaying
1794        // below.
1795        let current_context = *self.vm_store_context_mut().component_context_mut();
1796
1797        // Replay the deferred `enter_guest_sync_call`s outermost-first so that
1798        // the resulting `ConcurrentState` matches what the non-deferred path
1799        // would have otherwise produced.
1800        for (callee_async, callee_instance, saved_context) in frames.into_iter().rev() {
1801            // Restore the caller's context slots so that we save the correct
1802            // values into the caller's thread, exactly as the non-deferred path
1803            // would have on entry.
1804            *self.vm_store_context_mut().component_context_mut() = saved_context;
1805            let callee = RuntimeInstance {
1806                instance: id,
1807                index: RuntimeComponentInstanceIndex::from_u32(callee_instance),
1808            };
1809            self.enter_guest_sync_call(callee_async, callee)?;
1810        }
1811
1812        // Replaying done; restore the current context.
1813        *self.vm_store_context_mut().component_context_mut() = current_context;
1814
1815        Ok(self
1816            .concurrent_state_mut_without_forcing_current_thread()
1817            .unforced_current_thread)
1818    }
1819
1820    fn current_guest_thread(&mut self) -> Result<QualifiedThreadId> {
1821        match self.current_thread()?.guest() {
1822            Some(id) => Ok(*id),
1823            None => bail_bug!("current thread is not a guest thread"),
1824        }
1825    }
1826
1827    // A result of `None` may indicate that this is either the top-level event
1828    // loop, a deferred host task, or concurrency support is disabled. In all
1829    // cases we don't have an ID for the task.
1830    pub(crate) fn current_materialized_host_task(&mut self) -> Result<Option<TableId<HostTask>>> {
1831        match self.current_thread()? {
1832            CurrentThread::Host(id) => Ok(Some(id)),
1833            CurrentThread::DeferredHost(_) | CurrentThread::None => Ok(None),
1834            _ => bail_bug!("current thread is not a host thread"),
1835        }
1836    }
1837
1838    /// Returns the current host task ID, materializing a deferred host task if
1839    /// one is active. `None` represents a call from the top-level host.
1840    fn materialize_host_task_id(&mut self) -> Result<Option<TableId<HostTask>>> {
1841        Ok(self
1842            .concurrent_state_mut()?
1843            .materialize_current_host_task_id()?)
1844    }
1845
1846    fn enter_sync_call(&mut self, callee: RuntimeInstance) -> Result<()> {
1847        log::trace!("enter sync-typed call {callee:?}");
1848        let state = self.instance_state(callee).concurrent_state();
1849        let old_do_not_suspend = state.do_not_suspend;
1850        state.do_not_suspend = true;
1851
1852        let thread = self.current_guest_thread()?;
1853        let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
1854        if thread.old_do_not_suspend.is_some() {
1855            bail_bug!("current thread already has `old_do_not_suspend` value");
1856        }
1857
1858        thread.old_do_not_suspend = Some(old_do_not_suspend);
1859
1860        Ok(())
1861    }
1862
1863    fn exit_sync_call(&mut self, callee: RuntimeInstance) -> Result<()> {
1864        log::trace!("exit sync-typed call {callee:?}");
1865        let thread = self.current_guest_thread()?;
1866        let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
1867        let Some(old_do_not_suspend) = thread.old_do_not_suspend.take() else {
1868            bail_bug!("current thread missing `old_do_not_suspend` value");
1869        };
1870        let state = self.instance_state(callee).concurrent_state();
1871        state.do_not_suspend = old_do_not_suspend;
1872        Ok(())
1873    }
1874
1875    /// Push a `GuestTask` onto the task stack for either a sync-to-sync,
1876    /// guest-to-guest call or a sync host-to-guest call.
1877    ///
1878    /// This task will only be used for the purpose of handling calls to
1879    /// intrinsic functions; both parameter lowering and result lifting are
1880    /// assumed to be taken care of elsewhere.
1881    ///
1882    /// NB: for sync-to-sync, guest-to-guest calls we delay task construction in
1883    /// fused adapters, see `StoreOpaque::current_thread`, `VMDeferredThread`,
1884    /// and `lower_fact_enter_sync_call`. Make sure all this stuff stays in
1885    /// sync!
1886    pub(crate) fn enter_guest_sync_call(
1887        &mut self,
1888        callee_async_typed: bool,
1889        callee: RuntimeInstance,
1890    ) -> Result<()> {
1891        log::trace!("enter sync-lifted call {callee:?}");
1892        if !self.concurrency_support() {
1893            return self.enter_call_not_concurrent();
1894        }
1895
1896        let thread = self.current_thread()?;
1897        let caller = if let Some(thread) = thread.guest() {
1898            Caller::Guest { thread: *thread }
1899        } else {
1900            Caller::Host {
1901                tx: None,
1902                host_future_present: false,
1903                caller: self.materialize_host_task_id()?,
1904            }
1905        };
1906
1907        let state = self.concurrent_state_mut()?;
1908        let guest_thread = GuestTask::new(
1909            state,
1910            Box::new(move |_, _| bail_bug!("cannot lower params in sync call")),
1911            LiftResult {
1912                lift: Box::new(move |_, _| bail_bug!("cannot lift result in sync call")),
1913                ty: TypeTupleIndex::reserved_value(),
1914                memory: None,
1915                string_encoding: StringEncoding::Utf8,
1916            },
1917            caller,
1918            None,
1919            callee,
1920            callee_async_typed,
1921            true,
1922        )?;
1923
1924        Instance::from_wasmtime(self, callee.instance).add_guest_thread_to_instance_table(
1925            guest_thread.thread,
1926            self,
1927            callee.index,
1928        )?;
1929        self.set_thread(guest_thread)?;
1930
1931        if !callee_async_typed {
1932            self.enter_sync_call(callee)?;
1933        }
1934
1935        Ok(())
1936    }
1937
1938    /// Pop a `GuestTask` previously pushed using `enter_guest_sync_call`.
1939    ///
1940    /// NB: for sync-to-sync, guest-to-guest calls we delay task construction in
1941    /// fused adapters and then when the call returns we check to see if the
1942    /// task's contruction was forced and if not avoid calling out of the JIT
1943    /// code to this function. See `lower_fact_exit_sync_call`. Make sure all
1944    /// this stuff stays in sync!
1945    pub(crate) fn exit_guest_sync_call(&mut self) -> Result<()> {
1946        if !self.concurrency_support() {
1947            return Ok(self.exit_call_not_concurrent());
1948        }
1949
1950        let thread = match self.current_thread()?.guest() {
1951            Some(t) => *t,
1952            None => bail_bug!("expected task when exiting"),
1953        };
1954        let task = self.concurrent_state_mut()?.get_mut(thread.task)?;
1955        let instance = task.instance;
1956
1957        let caller = match &task.caller {
1958            &Caller::Guest { thread } => thread.into(),
1959            &Caller::Host { caller, .. } => caller
1960                .map(CurrentThread::Host)
1961                .unwrap_or(CurrentThread::None),
1962        };
1963        task.lift_result = None;
1964        task.exited = true;
1965        let async_typed = task.async_typed;
1966
1967        if !async_typed {
1968            self.exit_sync_call(instance)?;
1969        }
1970
1971        self.set_thread(caller)?;
1972
1973        log::trace!("exit sync-lifted call {instance:?}");
1974
1975        if async_typed {
1976            // If we're async-typed, returning control to our caller won't help
1977            // resolve any outstanding sync-typed call which might be in
1978            // progress, so we may need to switch or trap before exiting this
1979            // thread:
1980            self.switch_or_trap_if_may_not_suspend(instance)?;
1981        }
1982
1983        self.cleanup_thread(thread, instance, CleanupTask::Yes)?;
1984
1985        Ok(())
1986    }
1987
1988    /// Similar to `enter_guest_sync_call` except for when the guest makes a
1989    /// transition to the host.
1990    ///
1991    /// This initially records a deferred host call. A full [`HostTask`] should
1992    /// be allocated later if needed via
1993    /// [`StoreOpaque::materialize_host_task_id`].
1994    pub(crate) fn host_task_create(&mut self) -> Result<EnteredHostTask> {
1995        if !self.concurrency_support() {
1996            self.enter_call_not_concurrent()?;
1997            return Ok(None);
1998        }
1999        let caller = self.current_guest_thread()?;
2000        log::trace!("new deferred host task with caller {caller:?}");
2001
2002        self.set_thread(CurrentThread::DeferredHost(caller))?;
2003        let state = self.concurrent_state_mut()?;
2004        debug_assert!(state.deferred_host_call_context.is_none());
2005        state.deferred_host_call_context = Some(CallContext::default());
2006        state.debug_assert_deferred_host_invariant();
2007        Ok(Some(caller))
2008    }
2009
2010    /// Dual of `host_task_create` and signifies that the host has finished and
2011    /// will be cleaned up.
2012    ///
2013    /// Note that this isn't invoked when the host is invoked asynchronously and
2014    /// the host isn't complete yet. In that situation the host task persists
2015    /// and will be cleaned up separately in `subtask_drop`
2016    pub(crate) fn host_task_delete(
2017        &mut self,
2018        original_task: EnteredHostTask,
2019        materialized_task: Option<TableId<HostTask>>,
2020    ) -> Result<()> {
2021        match original_task {
2022            Some(caller) => {
2023                self.set_thread(caller)?;
2024                if materialized_task.is_none() {
2025                    let state = self.concurrent_state_mut()?;
2026                    let context = state
2027                        .deferred_host_call_context
2028                        .take()
2029                        .expect("deferred host call context should be present");
2030                    debug_assert!(context.is_empty());
2031                    state.debug_assert_deferred_host_invariant();
2032                }
2033                log::trace!(
2034                    "delete host task with caller {original_task:?} and materialized as {materialized_task:?}"
2035                );
2036                if let Some(task) = materialized_task {
2037                    Waitable::Host(task).delete_from(self)?;
2038                }
2039            }
2040            None => {
2041                debug_assert!(materialized_task.is_none());
2042                self.exit_call_not_concurrent();
2043            }
2044        }
2045        Ok(())
2046    }
2047
2048    /// Helper function to retrieve the `InstanceState` for the
2049    /// specified instance.
2050    fn instance_state(&mut self, instance: RuntimeInstance) -> &mut InstanceState {
2051        self.component_instance_mut(instance.instance)
2052            .instance_state(instance.index)
2053    }
2054
2055    /// Configure the currently running `thread`.
2056    ///
2057    /// This will save off any state necessary for the previous thread, if
2058    /// applicable, and then it'll additionally update state for `thread` if
2059    /// needed too.
2060    pub(crate) fn set_thread(&mut self, thread: impl Into<CurrentThread>) -> Result<CurrentThread> {
2061        let thread = thread.into();
2062        let state = self.concurrent_state_mut()?;
2063        state.debug_assert_deferred_host_invariant();
2064        let old_thread = mem::replace(&mut state.unforced_current_thread, thread);
2065
2066        state.handle_thread_switch(old_thread, thread)?;
2067
2068        // First thing to do after swapping threads is updating the context
2069        // slots for this thread within the store. This restores the behavior of
2070        // `context.{get,set}`. This involves taking the old state out of the
2071        // store, saving it in the thread that's being swapped from, and doing
2072        // the inverse for the new thread. When debug assertions are enabled
2073        // this also leaves behind sentinel values to try to uncover bugs where
2074        // this may be forgotten.
2075        if let Some(old_thread) = old_thread.guest() {
2076            let old_context = *self.vm_store_context_mut().component_context_mut();
2077            self.concurrent_state_mut()?
2078                .get_mut(old_thread.thread)?
2079                .context = old_context;
2080        }
2081        if cfg!(debug_assertions) {
2082            *self.vm_store_context_mut().component_context_mut() =
2083                [u32::MAX; NUM_COMPONENT_CONTEXT_SLOTS];
2084        }
2085        if let Some(thread) = thread.guest() {
2086            let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
2087            let context = thread.context;
2088            if cfg!(debug_assertions) {
2089                thread.context = [u32::MAX; NUM_COMPONENT_CONTEXT_SLOTS];
2090            }
2091            *self.vm_store_context_mut().component_context_mut() = context;
2092        }
2093
2094        // Keep the JIT-visible current-thread pointer in sync.
2095        *self.vm_store_context_mut().current_thread_mut() = if thread.is_none() {
2096            VMLazyThread::none()
2097        } else {
2098            VMLazyThread::forced()
2099        };
2100
2101        Ok(old_thread)
2102    }
2103
2104    /// Call `switch_if_may_not_suspend` and trap if it returns `false`.
2105    fn switch_or_trap_if_may_not_suspend(&mut self, instance: RuntimeInstance) -> Result<()> {
2106        if self.switch_if_may_not_suspend(instance)? {
2107            Ok(())
2108        } else {
2109            Err(Trap::CannotBlockSyncTask.into())
2110        }
2111    }
2112
2113    /// Check if the specified instance has a sync-typed call in progress; if so
2114    /// attempt to switch to another ready thread for that instance, and if no
2115    /// such thread exists, return false.
2116    fn switch_if_may_not_suspend(&mut self, instance: RuntimeInstance) -> Result<bool> {
2117        // Call this for the side effect of forcing any deferred task creation,
2118        // which may influence the value of `ConcurrentState::do_not_suspend`
2119        // below:
2120        self.concurrent_state_mut()?;
2121
2122        Ok(!self.concurrency_support()
2123            || !self
2124                .instance_state(instance)
2125                .concurrent_state()
2126                .do_not_suspend
2127            || self
2128                .concurrent_state_mut()?
2129                .promote_instance_local_thread_work_item(instance)?)
2130    }
2131
2132    /// Record that we're about to enter a (sub-)component instance which does
2133    /// not support more than one concurrent, stackful activation, meaning it
2134    /// cannot be entered again until the next call returns.
2135    fn enter_instance(&mut self, instance: RuntimeInstance) {
2136        log::trace!("enter {instance:?}");
2137        self.instance_state(instance)
2138            .concurrent_state()
2139            .do_not_enter = true;
2140    }
2141
2142    /// Record that we've exited a (sub-)component instance previously entered
2143    /// with `Self::enter_instance` and then calls `Self::partition_pending`.
2144    /// See the documentation for the latter for details.
2145    fn exit_instance(&mut self, instance: RuntimeInstance) -> Result<()> {
2146        log::trace!("exit {instance:?}");
2147        self.instance_state(instance)
2148            .concurrent_state()
2149            .do_not_enter = false;
2150        self.partition_pending(instance)
2151    }
2152
2153    /// Iterate over `InstanceState::pending`, moving any ready items into the
2154    /// "high priority" work item queue.
2155    ///
2156    /// Also, notify `ConcurrentState::ready_for_concurrent_call_waker` if
2157    /// present.
2158    ///
2159    /// See `GuestCall::is_ready` for details.
2160    fn partition_pending(&mut self, instance: RuntimeInstance) -> Result<()> {
2161        for (thread, kind) in
2162            mem::take(&mut self.instance_state(instance).concurrent_state().pending).into_iter()
2163        {
2164            let call = GuestCall { thread, kind };
2165            if call.is_ready(self)? {
2166                self.concurrent_state_mut()?
2167                    .push_high_priority(WorkItem::GuestCall { instance, call });
2168            } else {
2169                self.instance_state(instance)
2170                    .concurrent_state()
2171                    .pending
2172                    .insert(call.thread, call.kind);
2173            }
2174        }
2175
2176        if let Some(waker) = self
2177            .concurrent_state_mut()?
2178            .ready_for_concurrent_call_waker
2179            .take()
2180        {
2181            waker.wake();
2182        }
2183
2184        Ok(())
2185    }
2186
2187    /// Implements the `backpressure.{inc,dec}` intrinsics.
2188    pub(crate) fn backpressure_modify(
2189        &mut self,
2190        caller_instance: RuntimeInstance,
2191        modify: impl FnOnce(u16) -> Option<u16>,
2192    ) -> Result<()> {
2193        let state = self.instance_state(caller_instance).concurrent_state();
2194        let old = state.backpressure;
2195        let new = modify(old).ok_or_else(|| Trap::BackpressureOverflow)?;
2196        state.backpressure = new;
2197
2198        if old > 0 && new == 0 {
2199            // Backpressure was previously enabled and is now disabled; move any
2200            // newly-eligible guest calls to the "high priority" queue.
2201            self.partition_pending(caller_instance)?;
2202        }
2203
2204        Ok(())
2205    }
2206
2207    /// Resume the specified fiber, giving it exclusive access to the specified
2208    /// store.
2209    async fn resume_fiber(&mut self, fiber: StoreFiber<'static>) -> Result<()> {
2210        let old_thread = self.current_thread()?;
2211        log::trace!("resume_fiber: save current thread {old_thread:?}");
2212
2213        let fiber = fiber::resolve_or_release(self, fiber).await?;
2214
2215        self.set_thread(old_thread)?;
2216
2217        let state = self.concurrent_state_mut()?;
2218
2219        if let Some(ot) = old_thread.guest() {
2220            state.get_mut(ot.thread)?.state = GuestThreadState::Running;
2221        }
2222        log::trace!("resume_fiber: restore current thread {old_thread:?}");
2223
2224        if let Some(mut fiber) = fiber {
2225            log::trace!("resume_fiber: suspend reason {:?}", &state.suspend_reason);
2226            // See the `SuspendReason` documentation for what each case means.
2227            let reason = match state.suspend_reason.take() {
2228                Some(r) => r,
2229                None => bail_bug!("suspend reason missing when resuming fiber"),
2230            };
2231            match reason {
2232                SuspendReason::NeedWork => {
2233                    if state.worker.is_none() {
2234                        state.worker = Some(fiber);
2235                    } else {
2236                        fiber.dispose(self);
2237                    }
2238                }
2239                SuspendReason::Yielding { thread } => {
2240                    state.get_mut(thread.thread)?.state = GuestThreadState::Ready { fiber };
2241                    let instance = state.get_mut(thread.task)?.instance;
2242                    state.push_low_priority(WorkItem::ResumeThread { instance, thread });
2243                }
2244                SuspendReason::ExplicitlySuspending { thread } => {
2245                    state.get_mut(thread.thread)?.state = GuestThreadState::Suspended(fiber);
2246                }
2247                SuspendReason::Waiting { set, thread } => {
2248                    let old = state
2249                        .get_mut(set)?
2250                        .waiting
2251                        .insert(thread, WaitMode::Fiber(fiber));
2252                    assert!(old.is_none());
2253                }
2254                SuspendReason::YieldingToSubtask { thread } => {
2255                    // In this case, the thread has either invoked or sent a
2256                    // cancel request to a subtask, and is now yielding to that
2257                    // subtask.  According to the CM spec, that subtask may only
2258                    // yield back to the original thread the first time it
2259                    // suspends or exits (or a thread that it has resumed
2260                    // suspends or exits, etc.), which we ensure by setting
2261                    // `ConcurrentState::next_switch_item` here.
2262
2263                    let item = WorkItem::ResumeFiber {
2264                        instance: state.get_mut(thread.task)?.instance,
2265                        thread,
2266                        fiber,
2267                    };
2268
2269                    if state.next_switch_item.replace(item).is_some() {
2270                        // This should be unreachable per the save/restore code
2271                        // in `Self::suspend`.
2272                        bail_bug!(
2273                            "`ConcurrentState::next_switch_item` was already `Some(_)` when \
2274                             a thread wanted to wait on a subtask"
2275                        );
2276                    }
2277                }
2278            };
2279        } else {
2280            log::trace!("resume_fiber: fiber has exited");
2281        }
2282
2283        Ok(())
2284    }
2285
2286    /// Suspend the current fiber, storing the reason in
2287    /// `ConcurrentState::suspend_reason` to indicate the conditions under which
2288    /// it should be resumed.
2289    ///
2290    /// See the `SuspendReason` documentation for details.
2291    fn suspend(&mut self, reason: SuspendReason) -> Result<()> {
2292        log::trace!("suspend fiber: {reason:?}");
2293
2294        let state = self.concurrent_state_mut()?;
2295
2296        // If we're yielding or waiting on behalf of a guest thread, we'll be
2297        // overwriting the current thread, so save it now and restore it once
2298        // we've resumed.
2299        //
2300        // Also, if we're yielding to a subtask, we're about to overwrite
2301        // `ConcurrentState::next_switch_item`, so also save and restore that.
2302        let (save_and_restore_thread, save_and_restore_next_switch_item) = match &reason {
2303            SuspendReason::Yielding { .. }
2304            | SuspendReason::Waiting { .. }
2305            | SuspendReason::ExplicitlySuspending { .. } => {
2306                // If there's a thread waiting for this subtask to suspend, this
2307                // is a good time to switch back to it.
2308                if state.switch_item.is_none() {
2309                    state.take_next_switch_item()?;
2310                }
2311
2312                (true, false)
2313            }
2314            SuspendReason::YieldingToSubtask { .. } => (true, true),
2315            SuspendReason::NeedWork => (false, false),
2316        };
2317
2318        let old_next_switch_item = if save_and_restore_next_switch_item {
2319            let item = state.next_switch_item.take();
2320            // Note that we store it in the table here rather than directly in a
2321            // local variable to ensure the fiber is disposed of properly if we
2322            // end up trapping or panicking.
2323            Some(state.push(item)?)
2324        } else {
2325            None
2326        };
2327
2328        let old_guest_thread = if save_and_restore_thread {
2329            self.current_thread()?
2330        } else {
2331            CurrentThread::None
2332        };
2333
2334        let suspend_reason = &mut self.concurrent_state_mut()?.suspend_reason;
2335        assert!(suspend_reason.is_none());
2336        *suspend_reason = Some(reason);
2337
2338        // We'll panic if we call `Self::with_blocking` when the fiber is being
2339        // disposed, so check for that and exit ASAP if appropriate.
2340        if !self.fiber_async_state_mut().can_block() {
2341            return Err(format_err!("future dropped"));
2342        }
2343
2344        self.with_blocking(|_, cx| cx.suspend(StoreFiberYield::ReleaseStore))?;
2345
2346        if save_and_restore_thread {
2347            self.set_thread(old_guest_thread)?;
2348        }
2349
2350        if let Some(item) = old_next_switch_item {
2351            let state = self.concurrent_state_mut()?;
2352            state.next_switch_item = state.delete(item)?;
2353        }
2354
2355        Ok(())
2356    }
2357
2358    fn wait_for_event(
2359        &mut self,
2360        caller_instance: RuntimeInstance,
2361        waitable: Waitable,
2362    ) -> Result<()> {
2363        let caller = self.current_guest_thread()?;
2364        let state = self.concurrent_state_mut()?;
2365
2366        waitable.trap_if_in_waitable_set(state)?;
2367
2368        let set = state.get_mut(caller.thread)?.sync_call_set;
2369        waitable.join(state, Some(set))?;
2370
2371        self.switch_or_trap_if_may_not_suspend(caller_instance)?;
2372
2373        self.suspend(SuspendReason::Waiting {
2374            set,
2375            thread: caller,
2376        })?;
2377        let state = self.concurrent_state_mut()?;
2378
2379        waitable.join(state, None)
2380    }
2381
2382    /// Cleans up the data structures backing the `guest_thread` specified,
2383    /// removing it from the internal tables of `runtime_instance` as well.
2384    ///
2385    /// This function is used whenever a guest thread has fully exited and
2386    /// completed. This'll clean up the associated `GuestThread` structure and
2387    /// related resources it contains.
2388    ///
2389    /// Other functionality that this implements is:
2390    ///
2391    /// * This will perform conditional cleanup of the `GuestTask` that owns
2392    ///   this thread if `cleanup_task` is `CleanupTask::Yes`.
2393    /// * If there are no more threads in the `GuestTask` that this thread is
2394    ///   associated with, and if the task hasn't produced a result (e.g. it's not
2395    ///   returned or cancelled), then a trap will be raised that a result
2396    ///   wasn't ever produced.
2397    /// * If this task is finished, meaning the top-level thread exited and
2398    ///   additionally it's been returned or cancelled, then this will handle
2399    ///   management of the store's "active interesting tasks" counter.
2400    ///
2401    /// Effectively this is intended to be a "narrow waist" through which many
2402    /// destruction operations are funneled through.
2403    fn cleanup_thread(
2404        &mut self,
2405        guest_thread: QualifiedThreadId,
2406        runtime_instance: RuntimeInstance,
2407        cleanup_task: CleanupTask,
2408    ) -> Result<()> {
2409        let state = self.concurrent_state_mut()?;
2410        // If we never suspended, we never had a chance to deliver a subtask
2411        // status update, if any, to our caller, so we do that here:
2412        state.take_next_switch_item()?;
2413        let thread_data = state.get_mut(guest_thread.thread)?;
2414        let sync_call_set = thread_data.sync_call_set;
2415        if let Some(guest_id) = thread_data.instance_rep {
2416            self.instance_state(runtime_instance)
2417                .thread_handle_table()
2418                .guest_thread_remove(guest_id)?;
2419        }
2420        let state = self.concurrent_state_mut()?;
2421
2422        // Clean up any pending subtasks in the sync_call_set
2423        for waitable in mem::take(&mut state.get_mut(sync_call_set)?.ready) {
2424            if let Some(Event::Subtask {
2425                status: Status::Returned | Status::ReturnCancelled,
2426            }) = waitable.common(self.concurrent_state_mut()?)?.event
2427            {
2428                waitable.delete_from(self)?;
2429            }
2430        }
2431
2432        let state = self.concurrent_state_mut()?;
2433        state.delete(guest_thread.thread)?;
2434        state.delete(sync_call_set)?;
2435        let task = state.get_mut(guest_thread.task)?;
2436        task.threads.remove(&guest_thread.thread);
2437
2438        if task.threads.is_empty() && !task.returned_or_cancelled() {
2439            bail!(Trap::NoAsyncResult);
2440        }
2441        let ready_to_delete = task.ready_to_delete();
2442
2443        if !task.decremented_interesting_task_count && task.exited && task.returned_or_cancelled() {
2444            task.decremented_interesting_task_count = true;
2445
2446            debug_assert!(state.interesting_tasks > 0);
2447            state.interesting_tasks -= 1;
2448            if state.interesting_tasks == 0
2449                && let Some(waker) = state.interesting_tasks_empty_waker.take()
2450            {
2451                waker.wake();
2452            }
2453        }
2454
2455        match cleanup_task {
2456            CleanupTask::Yes => {
2457                if ready_to_delete {
2458                    Waitable::Guest(guest_thread.task).delete_from(self)?;
2459                }
2460            }
2461            CleanupTask::No => {}
2462        }
2463
2464        Ok(())
2465    }
2466
2467    /// Performs cancellation of the `guest_task` specified with the
2468    /// precondition that the task hasn't lowered its parameters.
2469    ///
2470    /// In this situation the task hasn't ever been started meaning it hasn't
2471    /// actually run any wasm code yet. This requires cleaning up metadata such
2472    /// as thread information attached to the task.
2473    ///
2474    /// The main two entrypoints for this function are:
2475    ///
2476    /// * Task cancellation via `subtask.cancel`, the intrinsic.
2477    /// * Dropping a host `call_async` future which needs to cancel the task
2478    ///   because it cannot reference its parameters any more.
2479    fn cancel_guest_subtask_without_lowered_parameters(
2480        &mut self,
2481        caller_instance: RuntimeInstance,
2482        guest_task: TableId<GuestTask>,
2483    ) -> Result<()> {
2484        let concurrent_state = self.concurrent_state_mut()?;
2485        let task = concurrent_state.get_mut(guest_task)?;
2486        assert!(!task.already_lowered_parameters());
2487        // The task is in a `starting` state, meaning it hasn't run at
2488        // all yet.  Here we update its fields to indicate that it is
2489        // ready to delete immediately once `subtask.drop` is called.
2490        task.lower_params = None;
2491        task.lift_result = None;
2492        task.exited = true;
2493        let instance = task.instance;
2494
2495        // Clean up the thread within this task as it's now never going
2496        // to run.
2497        assert_eq!(1, task.threads.len());
2498        let thread = *task.threads.iter().next().unwrap();
2499        self.cleanup_thread(
2500            QualifiedThreadId {
2501                task: guest_task,
2502                thread,
2503            },
2504            caller_instance,
2505            CleanupTask::No,
2506        )?;
2507
2508        // Not yet started; cancel and remove from pending
2509        let pending = &mut self.instance_state(instance).concurrent_state().pending;
2510        let pending_count = pending.len();
2511        pending.retain(|thread, _| thread.task != guest_task);
2512        // If there were no pending threads for this task, we're in an error state
2513        if pending.len() == pending_count {
2514            bail!(Trap::SubtaskCancelAfterTerminal);
2515        }
2516        Ok(())
2517    }
2518
2519    /// Used by `ResourceTables` to record the scope of a borrow to get undone
2520    /// in the future.
2521    pub(crate) fn current_scope(&mut self) -> Result<Option<CurrentScope>> {
2522        if !self.concurrency_support() {
2523            return Ok(self
2524                .current_scope_id_not_concurrent()?
2525                .map(|id| CurrentScope::Id(Scope::Id(id))));
2526        }
2527
2528        Ok(match self.current_thread()? {
2529            CurrentThread::Guest(id) => Some(CurrentScope::Id(Scope::Id(id.task.rep()))),
2530            CurrentThread::Host(id) => Some(CurrentScope::Id(Scope::HostId(id.rep()))),
2531            CurrentThread::DeferredHost(_) => Some(CurrentScope::DeferredHost),
2532            CurrentThread::None => return Ok(None),
2533        })
2534    }
2535
2536    pub(crate) fn queue_task(
2537        &mut self,
2538        task: impl FnOnce(&mut dyn VMStore) -> Result<()> + Send + 'static,
2539    ) -> Result<()> {
2540        self.concurrent_state_mut()?
2541            .push_high_priority(WorkItem::WorkerFunction(AlwaysMut::new(Box::new(task))));
2542        Ok(())
2543    }
2544
2545    /// Used in `poll_until` just prior to trapping due to a "deadlock"
2546    /// condition.
2547    ///
2548    /// This helps us distinguish between a simple deadlock condition (where no
2549    /// work is available for the event loop to do, nor is there any way for new
2550    /// work to be added) and a "cannot block sync task" condition where at
2551    /// least once instance has an outstanding sync-typed task running, in which
2552    /// case we'll trap with a different error message.
2553    fn any_may_not_suspend(&mut self) -> Result<bool> {
2554        // Note that this currently requires a linear search across the whole
2555        // `ConcurrentState::table`.  We _could_ optimize that, but since (1)
2556        // this function is only used when trapping, (2) the only thing you can
2557        // really do with a store that's been poisoned by a trap is drop it, and
2558        // (3) we must do a linear search through the table when dropping a
2559        // store anyway to dispose of fibers, it's reasonable for us to also do
2560        // a linear search here.
2561        Ok(self
2562            .concurrent_state_mut()?
2563            .table
2564            .get_mut()
2565            .iter_mut()
2566            .filter_map(|(_, entry)| {
2567                if let Some(task) = entry.downcast_ref::<GuestTask>() {
2568                    Some(task.instance)
2569                } else {
2570                    None
2571                }
2572            })
2573            .collect::<Vec<_>>()
2574            .into_iter()
2575            .any(|instance| {
2576                self.instance_state(instance)
2577                    .concurrent_state()
2578                    .do_not_suspend
2579            }))
2580    }
2581}
2582
2583enum CleanupTask {
2584    Yes,
2585    No,
2586}
2587
2588impl Instance {
2589    /// Get the next pending event for the specified task and (optional)
2590    /// waitable set, along with the waitable handle if applicable.
2591    fn get_event(
2592        self,
2593        store: &mut StoreOpaque,
2594        guest_task: TableId<GuestTask>,
2595        set: Option<TableId<WaitableSet>>,
2596        cancellable: bool,
2597    ) -> Result<Option<(Event, Option<(Waitable, u32)>)>> {
2598        let state = store.concurrent_state_mut()?;
2599
2600        let task = state.get_mut(guest_task)?;
2601        let event = &mut task.event;
2602        if let Some(ev) = event
2603            && (cancellable || !matches!(ev, Event::Cancelled))
2604        {
2605            log::trace!("deliver event {ev:?} to {guest_task:?}");
2606
2607            if matches!(ev, Event::Cancelled) {
2608                task.cancel_request_delivered = true;
2609            }
2610
2611            let ev = *ev;
2612            *event = None;
2613            return Ok(Some((ev, None)));
2614        }
2615
2616        let set = match set {
2617            Some(set) => set,
2618            None => return Ok(None),
2619        };
2620        let waitable = match state.get_mut(set)?.ready.pop_first() {
2621            Some(v) => v,
2622            None => return Ok(None),
2623        };
2624
2625        let common = waitable.common(state)?;
2626        let handle = match common.handle {
2627            Some(h) => h,
2628            None => bail_bug!("handle not set when delivering event"),
2629        };
2630        let event = match common.event.take() {
2631            Some(e) => e,
2632            None => bail_bug!("event not set when delivering event"),
2633        };
2634
2635        log::trace!(
2636            "deliver event {event:?} to {guest_task:?} for {waitable:?} (handle {handle}); set {set:?}"
2637        );
2638
2639        waitable.on_delivery(store, self, event)?;
2640
2641        Ok(Some((event, Some((waitable, handle)))))
2642    }
2643
2644    /// Handle the `CallbackCode` returned from an async-lifted export or its
2645    /// callback.
2646    ///
2647    /// If this returns `Ok(Some(call))`, then `call` should be run immediately
2648    /// using `handle_guest_call`.
2649    fn handle_callback_code(
2650        self,
2651        store: &mut StoreOpaque,
2652        guest_thread: QualifiedThreadId,
2653        runtime_instance: RuntimeComponentInstanceIndex,
2654        code: u32,
2655    ) -> Result<()> {
2656        let (code, set) = unpack_callback_code(code);
2657
2658        log::trace!("received callback code from {guest_thread:?}: {code} (set: {set})");
2659
2660        let state = store.concurrent_state_mut()?;
2661
2662        state.take_next_switch_item()?;
2663
2664        let get_set = |store: &mut StoreOpaque, handle| -> Result<_> {
2665            let set = store
2666                .instance_state(self.runtime_instance(runtime_instance))
2667                .handle_table()
2668                .waitable_set_rep(handle)?;
2669
2670            Ok(TableId::<WaitableSet>::new(set))
2671        };
2672
2673        match code {
2674            callback_code::EXIT => {
2675                log::trace!("implicit thread {guest_thread:?} completed");
2676                let task = store.concurrent_state_mut()?.get_mut(guest_thread.task)?;
2677                task.exited = true;
2678                task.callback = None;
2679
2680                let runtime_instance = self.runtime_instance(runtime_instance);
2681
2682                // Since we're async-typed, returning control to our caller
2683                // won't help resolve any outstanding sync-typed call which
2684                // might be in progress, so we may need to switch or trap before
2685                // exiting this thread:
2686                store.switch_or_trap_if_may_not_suspend(runtime_instance)?;
2687
2688                store.cleanup_thread(guest_thread, runtime_instance, CleanupTask::Yes)?;
2689            }
2690            callback_code::YIELD => {
2691                // Set `GuestTask::wake_on_cancel` to allow `subtask.cancel` to
2692                // promote this thread if appropriate.
2693                let old = state
2694                    .get_mut(guest_thread.thread)?
2695                    .wake_on_cancel
2696                    .replace(WakeOnCancel::Yielding);
2697                if !old.is_none() {
2698                    bail_bug!("thread unexpectedly had wake_on_cancel set");
2699                }
2700
2701                let task = state.get_mut(guest_thread.task)?;
2702                // If an `Event::Cancelled` is pending, we'll deliver that;
2703                // otherwise, we'll deliver `Event::None`.  Note that
2704                // `GuestTask::event` is only ever set to one of those two
2705                // `Event` variants.
2706                if let Some(event) = task.event {
2707                    assert!(matches!(event, Event::None | Event::Cancelled));
2708                } else {
2709                    task.event = Some(Event::None);
2710                }
2711                let call = GuestCall {
2712                    thread: guest_thread,
2713                    kind: GuestCallKind::DeliverEvent {
2714                        instance: self,
2715                        set: None,
2716                    },
2717                };
2718                // Push this thread onto the "low priority" queue so it runs
2719                // after any other threads have had a chance to run.
2720                state.push_low_priority(WorkItem::GuestCall {
2721                    instance: self.runtime_instance(runtime_instance),
2722                    call,
2723                });
2724            }
2725            callback_code::WAIT => {
2726                let set = get_set(store, set)?;
2727                let state = store.concurrent_state_mut()?;
2728
2729                if state.get_mut(guest_thread.task)?.event.is_some()
2730                    || !state.get_mut(set)?.ready.is_empty()
2731                {
2732                    // An event is immediately available; deliver it ASAP.
2733                    state.push_high_priority(WorkItem::GuestCall {
2734                        instance: self.runtime_instance(runtime_instance),
2735                        call: GuestCall {
2736                            thread: guest_thread,
2737                            kind: GuestCallKind::DeliverEvent {
2738                                instance: self,
2739                                set: Some(set),
2740                            },
2741                        },
2742                    });
2743                } else {
2744                    // No event is immediately available.
2745                    //
2746                    // We're waiting, so register to be woken up when an event
2747                    // is published for this waitable set.
2748                    //
2749                    // Here we also set `GuestTask::wake_on_cancel` which allows
2750                    // `subtask.cancel` to interrupt the wait.
2751                    let old = state
2752                        .get_mut(guest_thread.thread)?
2753                        .wake_on_cancel
2754                        .replace(WakeOnCancel::Waiting(set));
2755                    if !old.is_none() {
2756                        bail_bug!("thread unexpectedly had wake_on_cancel set");
2757                    }
2758                    let old = state
2759                        .get_mut(set)?
2760                        .waiting
2761                        .insert(guest_thread, WaitMode::Callback(self));
2762                    if !old.is_none() {
2763                        bail_bug!("set's waiting set already had this thread registered");
2764                    }
2765                }
2766            }
2767            _ => bail!(Trap::UnsupportedCallbackCode),
2768        }
2769
2770        Ok(())
2771    }
2772
2773    /// Stage the specified guest call as the work item to run next, to be
2774    /// started as soon as backpressure and/or reentrance rules allow.
2775    ///
2776    /// SAFETY: The raw pointer arguments must be valid references to guest
2777    /// functions (with the appropriate signatures) when the closures staged by
2778    /// this function are called.
2779    unsafe fn stage_call<T: 'static>(
2780        self,
2781        mut store: StoreContextMut<T>,
2782        guest_thread: QualifiedThreadId,
2783        callee: SendSyncPtr<VMFuncRef>,
2784        param_count: usize,
2785        result_count: usize,
2786        async_: bool,
2787        callback: Option<SendSyncPtr<VMFuncRef>>,
2788        post_return: Option<SendSyncPtr<VMFuncRef>>,
2789        host_caller: bool,
2790    ) -> Result<()> {
2791        /// Return a closure which will call the specified function in the scope
2792        /// of the specified task.
2793        ///
2794        /// This will use `GuestTask::lower_params` to lower the parameters, but
2795        /// will not lift the result; instead, it returns a
2796        /// `[MaybeUninit<ValRaw>; MAX_FLAT_PARAMS]` from which the result, if
2797        /// any, may be lifted.  Note that an async-lifted export will have
2798        /// returned its result using the `task.return` intrinsic (or not
2799        /// returned a result at all, in the case of `task.cancel`), in which
2800        /// case the "result" of this call will either be a callback code or
2801        /// nothing.
2802        ///
2803        /// SAFETY: `callee` must be a valid `*mut VMFuncRef` at the time when
2804        /// the returned closure is called.
2805        unsafe fn make_call<T: 'static>(
2806            store: StoreContextMut<T>,
2807            guest_thread: QualifiedThreadId,
2808            callee: SendSyncPtr<VMFuncRef>,
2809            param_count: usize,
2810            result_count: usize,
2811        ) -> impl FnOnce(&mut dyn VMStore) -> Result<[MaybeUninit<ValRaw>; MAX_FLAT_PARAMS]>
2812        + Send
2813        + Sync
2814        + 'static
2815        + use<T> {
2816            let token = StoreToken::new(store);
2817            move |store: &mut dyn VMStore| {
2818                let mut storage = [MaybeUninit::uninit(); MAX_FLAT_PARAMS];
2819
2820                store
2821                    .concurrent_state_mut()?
2822                    .get_mut(guest_thread.thread)?
2823                    .state = GuestThreadState::Running;
2824                let task = store.concurrent_state_mut()?.get_mut(guest_thread.task)?;
2825                let lower = match task.lower_params.take() {
2826                    Some(l) => l,
2827                    None => bail_bug!("lower_params missing"),
2828                };
2829
2830                lower(store, &mut storage[..param_count])?;
2831
2832                let mut store = token.as_context_mut(store);
2833
2834                // SAFETY: Per the contract documented in `make_call's`
2835                // documentation, `callee` must be a valid pointer.
2836                unsafe {
2837                    crate::Func::call_unchecked_raw(
2838                        &mut store,
2839                        callee.as_non_null(),
2840                        NonNull::new(
2841                            &mut storage[..param_count.max(result_count)]
2842                                as *mut [MaybeUninit<ValRaw>] as _,
2843                        )
2844                        .unwrap(),
2845                    )?;
2846                }
2847
2848                Ok(storage)
2849            }
2850        }
2851
2852        // SAFETY: Per the contract described in this function documentation,
2853        // the `callee` pointer which `call` closes over must be valid when
2854        // called by the closure we queue below.
2855        let call = unsafe {
2856            make_call(
2857                store.as_context_mut(),
2858                guest_thread,
2859                callee,
2860                param_count,
2861                result_count,
2862            )
2863        };
2864
2865        let callee_instance = store
2866            .0
2867            .concurrent_state_mut()?
2868            .get_mut(guest_thread.task)?
2869            .instance;
2870
2871        let fun = if callback.is_some() {
2872            assert!(async_);
2873
2874            Box::new(move |store: &mut dyn VMStore| {
2875                self.add_guest_thread_to_instance_table(
2876                    guest_thread.thread,
2877                    store,
2878                    callee_instance.index,
2879                )?;
2880                let old_thread = store.set_thread(guest_thread)?;
2881                log::trace!(
2882                    "stackless call: replaced {old_thread:?} with {guest_thread:?} as current thread"
2883                );
2884
2885                store.enter_instance(callee_instance);
2886
2887                // SAFETY: See the documentation for `make_call` to review the
2888                // contract we must uphold for `call` here.
2889                //
2890                // Per the contract described in the `stage_call`
2891                // documentation, the `callee` pointer which `call` closes
2892                // over must be valid.
2893                let storage = call(store)?;
2894
2895                store.exit_instance(callee_instance)?;
2896
2897                store.set_thread(old_thread)?;
2898                let state = store.concurrent_state_mut()?;
2899                if let Some(t) = old_thread.guest() {
2900                    state.get_mut(t.thread)?.state = GuestThreadState::Running;
2901                }
2902                log::trace!("stackless call: restored {old_thread:?} as current thread");
2903
2904                // SAFETY: `wasmparser` will have validated that the callback
2905                // function returns a `i32` result.
2906                let code = unsafe { storage[0].assume_init() }.get_i32() as u32;
2907
2908                self.handle_callback_code(store, guest_thread, callee_instance.index, code)
2909            }) as Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send + Sync>
2910        } else {
2911            let token = StoreToken::new(store.as_context_mut());
2912            Box::new(move |store: &mut dyn VMStore| {
2913                self.add_guest_thread_to_instance_table(
2914                    guest_thread.thread,
2915                    store,
2916                    callee_instance.index,
2917                )?;
2918                let old_thread = store.set_thread(guest_thread)?;
2919                log::trace!(
2920                    "sync/async-stackful call: replaced {old_thread:?} with {guest_thread:?} as current thread",
2921                );
2922                let flags = self.id().get(store).instance_flags(callee_instance.index);
2923
2924                let callee_async_typed = store
2925                    .concurrent_state_mut()?
2926                    .get_mut(guest_thread.task)?
2927                    .async_typed;
2928
2929                // Unless this is a callback-less (i.e. stackful) async-lifted
2930                // or sync-typed export, we need to record that the instance
2931                // cannot be entered until the call returns.
2932                if !async_ && callee_async_typed {
2933                    store.enter_instance(callee_instance);
2934                }
2935
2936                if !callee_async_typed {
2937                    store.enter_sync_call(callee_instance)?;
2938                }
2939
2940                // SAFETY: See the documentation for `make_call` to review the
2941                // contract we must uphold for `call` here.
2942                //
2943                // Per the contract described in the `stage_call`
2944                // documentation, the `callee` pointer which `call` closes
2945                // over must be valid.
2946                let storage = call(store)?;
2947
2948                if !callee_async_typed {
2949                    store.exit_sync_call(callee_instance)?;
2950                }
2951
2952                if !async_ {
2953                    // This is a sync-lifted export, so now is when we lift the
2954                    // result, optionally call the post-return function, if any,
2955                    // and finally notify any current or future waiters that the
2956                    // subtask has returned.
2957
2958                    if callee_async_typed {
2959                        store.exit_instance(callee_instance)?;
2960                    }
2961
2962                    let lift = {
2963                        let state = store.concurrent_state_mut()?;
2964                        if !state.get_mut(guest_thread.task)?.result.is_none() {
2965                            bail_bug!("task has already produced a result");
2966                        }
2967
2968                        match state.get_mut(guest_thread.task)?.lift_result.take() {
2969                            Some(lift) => lift,
2970                            None => bail_bug!("lift_result field is missing"),
2971                        }
2972                    };
2973
2974                    // SAFETY: `result_count` represents the number of core Wasm
2975                    // results returned, per `wasmparser`.
2976                    let result = (lift.lift)(store, unsafe {
2977                        mem::transmute::<&[MaybeUninit<ValRaw>], &[ValRaw]>(
2978                            &storage[..result_count],
2979                        )
2980                    })?;
2981
2982                    let post_return_arg = match result_count {
2983                        0 => ValRaw::i32(0),
2984                        // SAFETY: `result_count` represents the number of
2985                        // core Wasm results returned, per `wasmparser`.
2986                        1 => unsafe { storage[0].assume_init() },
2987                        _ => unreachable!(),
2988                    };
2989
2990                    unsafe {
2991                        call_post_return(
2992                            token.as_context_mut(store),
2993                            post_return.map(|v| v.as_non_null()),
2994                            post_return_arg,
2995                            flags,
2996                        )?;
2997                    }
2998
2999                    self.task_complete(store, guest_thread.task, result, Status::Returned)?;
3000                }
3001
3002                store.set_thread(old_thread)?;
3003
3004                store
3005                    .concurrent_state_mut()?
3006                    .get_mut(guest_thread.task)?
3007                    .exited = true;
3008
3009                log::trace!(
3010                    "clean up thread; async lifted? {async_} async typed? {callee_async_typed}"
3011                );
3012
3013                if callee_async_typed {
3014                    // If we're async-typed, returning control to our caller
3015                    // won't help resolve any outstanding sync-typed call which
3016                    // might be in progress, so we may need to switch or trap
3017                    // before exiting this thread:
3018                    store.switch_or_trap_if_may_not_suspend(callee_instance)?;
3019                }
3020
3021                // This is a callback-less call, so the implicit thread has now completed
3022                store.cleanup_thread(guest_thread, callee_instance, CleanupTask::Yes)?;
3023                Ok(())
3024            })
3025        };
3026
3027        store.0.concurrent_state_mut()?.push_work_item(
3028            WorkItem::GuestCall {
3029                instance: callee_instance,
3030                call: GuestCall {
3031                    thread: guest_thread,
3032                    kind: GuestCallKind::StartImplicit(fun),
3033                },
3034            },
3035            if host_caller {
3036                Priority::High
3037            } else {
3038                Priority::Switch
3039            },
3040        )?;
3041
3042        Ok(())
3043    }
3044
3045    /// Prepare (but do not start) a guest->guest call.
3046    ///
3047    /// This is called from fused adapter code generated in
3048    /// `wasmtime_environ::fact::trampoline::Compiler`.  `start` and `return_`
3049    /// are synthesized Wasm functions which move the parameters from the caller
3050    /// to the callee and the result from the callee to the caller,
3051    /// respectively.  The adapter will call `Self::start_call` immediately
3052    /// after calling this function.
3053    ///
3054    /// SAFETY: All the pointer arguments must be valid pointers to guest
3055    /// entities (and with the expected signatures for the function references
3056    /// -- see `wasmtime_environ::fact::trampoline::Compiler` for details).
3057    unsafe fn prepare_call<T: 'static>(
3058        self,
3059        mut store: StoreContextMut<T>,
3060        start: NonNull<VMFuncRef>,
3061        return_: NonNull<VMFuncRef>,
3062        caller_instance: RuntimeComponentInstanceIndex,
3063        callee_instance: RuntimeComponentInstanceIndex,
3064        task_return_type: TypeTupleIndex,
3065        callee_async_typed: bool,
3066        memory: *mut VMMemoryDefinition,
3067        string_encoding: StringEncoding,
3068        caller_info: CallerInfo,
3069    ) -> Result<()> {
3070        enum ResultInfo {
3071            Heap { results: u32 },
3072            Stack { result_count: u32 },
3073        }
3074
3075        let result_info = match &caller_info {
3076            CallerInfo::Async {
3077                has_result: true,
3078                params,
3079            } => ResultInfo::Heap {
3080                results: match params.last() {
3081                    Some(r) => r.get_u32(),
3082                    None => bail_bug!("retptr missing"),
3083                },
3084            },
3085            CallerInfo::Async {
3086                has_result: false, ..
3087            } => ResultInfo::Stack { result_count: 0 },
3088            CallerInfo::Sync {
3089                result_count,
3090                params,
3091            } if *result_count > u32::try_from(MAX_FLAT_RESULTS)? => ResultInfo::Heap {
3092                results: match params.last() {
3093                    Some(r) => r.get_u32(),
3094                    None => bail_bug!("arg ptr missing"),
3095                },
3096            },
3097            CallerInfo::Sync { result_count, .. } => ResultInfo::Stack {
3098                result_count: *result_count,
3099            },
3100        };
3101
3102        let sync_caller = matches!(caller_info, CallerInfo::Sync { .. });
3103
3104        // Create a new guest task for the call, closing over the `start` and
3105        // `return_` functions to lift the parameters and lower the result,
3106        // respectively.
3107        let start = SendSyncPtr::new(start);
3108        let return_ = SendSyncPtr::new(return_);
3109        let token = StoreToken::new(store.as_context_mut());
3110        let old_thread = store.0.current_guest_thread()?;
3111
3112        let state = store.0.concurrent_state_mut()?;
3113
3114        debug_assert_eq!(
3115            state.get_mut(old_thread.task)?.instance,
3116            self.runtime_instance(caller_instance)
3117        );
3118
3119        let guest_thread = GuestTask::new(
3120            state,
3121            Box::new(move |store, dst| {
3122                let mut store = token.as_context_mut(store);
3123                assert!(dst.len() <= MAX_FLAT_PARAMS);
3124                // The `+ 1` here accounts for the return pointer, if any:
3125                let mut src = [MaybeUninit::uninit(); MAX_FLAT_PARAMS + 1];
3126                let count = match caller_info {
3127                    // Async callers, if they have a result, use the last
3128                    // parameter as a return pointer so chop that off if
3129                    // relevant here.
3130                    CallerInfo::Async { params, has_result } => {
3131                        let params = &params[..params.len() - usize::from(has_result)];
3132                        for (param, src) in params.iter().zip(&mut src) {
3133                            src.write(*param);
3134                        }
3135                        params.len()
3136                    }
3137
3138                    // Sync callers forward everything directly.
3139                    CallerInfo::Sync { params, .. } => {
3140                        for (param, src) in params.iter().zip(&mut src) {
3141                            src.write(*param);
3142                        }
3143                        params.len()
3144                    }
3145                };
3146                // SAFETY: `start` is a valid `*mut VMFuncRef` from
3147                // `wasmtime-cranelift`-generated fused adapter code.  Based on
3148                // how it was constructed (see
3149                // `wasmtime_environ::fact::trampoline::Compiler::compile_async_start_adapter`
3150                // for details) we know it takes count parameters and returns
3151                // `dst.len()` results.
3152                unsafe {
3153                    crate::Func::call_unchecked_raw(
3154                        &mut store,
3155                        start.as_non_null(),
3156                        NonNull::new(
3157                            &mut src[..count.max(dst.len())] as *mut [MaybeUninit<ValRaw>] as _,
3158                        )
3159                        .unwrap(),
3160                    )?;
3161                }
3162                dst.copy_from_slice(&src[..dst.len()]);
3163                let task = store.0.current_guest_thread()?.task;
3164                let state = store.0.concurrent_state_mut()?;
3165                Waitable::Guest(task).set_event(
3166                    state,
3167                    Some(Event::Subtask {
3168                        status: Status::Started,
3169                    }),
3170                )?;
3171                Ok(())
3172            }),
3173            LiftResult {
3174                lift: Box::new(move |store, src| {
3175                    // SAFETY: See comment in closure passed as `lower_params`
3176                    // parameter above.
3177                    let mut store = token.as_context_mut(store);
3178                    let mut my_src = src.to_owned(); // TODO: use stack to avoid allocation?
3179                    if let ResultInfo::Heap { results } = &result_info {
3180                        my_src.push(ValRaw::u32(*results));
3181                    }
3182
3183                    // SAFETY: `return_` is a valid `*mut VMFuncRef` from
3184                    // `wasmtime-cranelift`-generated fused adapter code.  Based
3185                    // on how it was constructed (see
3186                    // `wasmtime_environ::fact::trampoline::Compiler::compile_async_return_adapter`
3187                    // for details) we know it takes `src.len()` parameters and
3188                    // returns up to 1 result.
3189                    unsafe {
3190                        crate::Func::call_unchecked_raw(
3191                            &mut store,
3192                            return_.as_non_null(),
3193                            my_src.as_mut_slice().into(),
3194                        )?;
3195                    }
3196
3197                    let thread = store.0.current_guest_thread()?;
3198                    let state = store.0.concurrent_state_mut()?;
3199                    if sync_caller {
3200                        state.get_mut(thread.task)?.sync_result = SyncResult::Produced(
3201                            if let ResultInfo::Stack { result_count } = &result_info {
3202                                match result_count {
3203                                    0 => None,
3204                                    1 => Some(my_src[0]),
3205                                    _ => unreachable!(),
3206                                }
3207                            } else {
3208                                None
3209                            },
3210                        );
3211                    }
3212                    Ok(Box::new(DummyResult) as Box<dyn Any + Send + Sync>)
3213                }),
3214                ty: task_return_type,
3215                memory: NonNull::new(memory).map(SendSyncPtr::new),
3216                string_encoding,
3217            },
3218            Caller::Guest { thread: old_thread },
3219            None,
3220            self.runtime_instance(callee_instance),
3221            callee_async_typed,
3222            // We don't know whether the callee export was lifted sync or async
3223            // yet, but we'll update this in `start_call`:
3224            false,
3225        )?;
3226
3227        // Make the new thread the current one so that `Self::start_call` knows
3228        // which one to start.
3229        store.0.set_thread(guest_thread)?;
3230        log::trace!("pushed {guest_thread:?} as current thread; old thread was {old_thread:?}");
3231
3232        Ok(())
3233    }
3234
3235    /// Call the specified callback function for an async-lifted export.
3236    ///
3237    /// SAFETY: `function` must be a valid reference to a guest function of the
3238    /// correct signature for a callback.
3239    unsafe fn call_callback<T>(
3240        self,
3241        mut store: StoreContextMut<T>,
3242        function: SendSyncPtr<VMFuncRef>,
3243        event: Event,
3244        handle: u32,
3245    ) -> Result<u32> {
3246        let (ordinal, result) = event.parts();
3247        let params = &mut [
3248            ValRaw::u32(ordinal),
3249            ValRaw::u32(handle),
3250            ValRaw::u32(result),
3251        ];
3252        // SAFETY: `func` is a valid `*mut VMFuncRef` from either
3253        // `wasmtime-cranelift`-generated fused adapter code or
3254        // `component::Options`.  Per `wasmparser` callback signature
3255        // validation, we know it takes three parameters and returns one.
3256        unsafe {
3257            crate::Func::call_unchecked_raw(
3258                &mut store,
3259                function.as_non_null(),
3260                params.as_mut_slice().into(),
3261            )?;
3262        }
3263        Ok(params[0].get_u32())
3264    }
3265
3266    /// Start a guest->guest call previously prepared using
3267    /// `Self::prepare_call`.
3268    ///
3269    /// This is called from fused adapter code generated in
3270    /// `wasmtime_environ::fact::trampoline::Compiler`.  The adapter will call
3271    /// this function immediately after calling `Self::prepare_call`.
3272    ///
3273    /// SAFETY: The `*mut VMFuncRef` arguments must be valid pointers to guest
3274    /// functions with the appropriate signatures for the current guest task.
3275    /// If this is a call to an async-lowered import, the actual call may be
3276    /// deferred and run after this function returns, in which case the pointer
3277    /// arguments must also be valid when the call happens.
3278    unsafe fn start_call<T: 'static>(
3279        self,
3280        mut store: StoreContextMut<T>,
3281        callback: *mut VMFuncRef,
3282        post_return: *mut VMFuncRef,
3283        callee: NonNull<VMFuncRef>,
3284        param_count: u32,
3285        result_count: u32,
3286        flags: u32,
3287        storage: Option<&mut [MaybeUninit<ValRaw>]>,
3288    ) -> Result<u32> {
3289        let token = StoreToken::new(store.as_context_mut());
3290        let async_caller = storage.is_none();
3291        let guest_thread = store.0.current_guest_thread()?;
3292        let state = store.0.concurrent_state_mut()?;
3293
3294        if !state.event_loop_running {
3295            bail_bug!("Instance::start_call called without a running event loop");
3296        }
3297
3298        let callee = SendSyncPtr::new(callee);
3299        let param_count = usize::try_from(param_count)?;
3300        assert!(param_count <= MAX_FLAT_PARAMS);
3301        let result_count = usize::try_from(result_count)?;
3302        assert!(result_count <= MAX_FLAT_RESULTS);
3303
3304        let task = state.get_mut(guest_thread.task)?;
3305        let callee_async_typed = task.async_typed;
3306        let callee_instance = task.instance;
3307
3308        task.async_lifted = (flags & START_FLAG_ASYNC_CALLEE) != 0;
3309
3310        if let Some(callback) = NonNull::new(callback) {
3311            // We're calling an async-lifted export with a callback, so store
3312            // the callback and related context as part of the task so we can
3313            // call it later when needed.
3314            let callback = SendSyncPtr::new(callback);
3315            task.callback = Some(Box::new(move |store, event, handle| {
3316                let store = token.as_context_mut(store);
3317                unsafe { self.call_callback::<T>(store, callback, event, handle) }
3318            }));
3319        }
3320
3321        let Caller::Guest { thread: caller } = &task.caller else {
3322            // As of this writing, `start_call` is only used for guest->guest
3323            // calls.
3324            bail_bug!("start_call unexpectedly invoked for host->guest call");
3325        };
3326        let caller = *caller;
3327        let caller_instance = state.get_mut(caller.task)?.instance;
3328
3329        // Stage the call as the work item to run next, modulo backpressure, etc.
3330        unsafe {
3331            self.stage_call(
3332                store.as_context_mut(),
3333                guest_thread,
3334                callee,
3335                param_count,
3336                result_count,
3337                (flags & START_FLAG_ASYNC_CALLEE) != 0,
3338                NonNull::new(callback).map(SendSyncPtr::new),
3339                NonNull::new(post_return).map(SendSyncPtr::new),
3340                false,
3341            )?;
3342        }
3343
3344        let old_do_not_suspend = if callee_async_typed {
3345            // If we're starting an async-typed call, it is permitted to suspend
3346            // since that will just return control back to the caller, which may
3347            // be able to avoid blocking if needed.
3348            //
3349            // We'll restore the old value after the call either returns or
3350            // suspends.
3351            let state = store.0.instance_state(callee_instance).concurrent_state();
3352            let old_do_not_suspend = state.do_not_suspend;
3353            state.do_not_suspend = false;
3354            Some(old_do_not_suspend)
3355        } else {
3356            None
3357        };
3358
3359        let state = store.0.concurrent_state_mut()?;
3360
3361        // Use the caller's `GuestThread::sync_call_set` to register interest in
3362        // the subtask...
3363        let guest_waitable = Waitable::Guest(guest_thread.task);
3364        let old_set = guest_waitable.common(state)?.set;
3365        let set = state.get_mut(caller.thread)?.sync_call_set;
3366        guest_waitable.join(state, Some(set))?;
3367
3368        store.0.set_thread(CurrentThread::None)?;
3369
3370        // ... and suspend this fiber temporarily while we wait for it to start.
3371        //
3372        // Note that we _could_ call the callee directly using the current fiber
3373        // rather than suspend this one, but that would make reasoning about the
3374        // event loop more complicated and is probably only worth doing if
3375        // there's a measurable performance benefit.  In addition, it would mean
3376        // blocking the caller if the callee calls a blocking sync-lowered
3377        // import, and as of this writing the spec says we must not do that.
3378        //
3379        // Alternatively, the fused adapter code could be modified to call the
3380        // callee directly without calling a host-provided intrinsic at all (in
3381        // which case it would need to do its own, inline backpressure checks,
3382        // etc.).  Again, we'd want to see a measurable performance benefit
3383        // before committing to such an optimization.  And again, we'd need to
3384        // update the spec to allow that.
3385        let mut yielded = false;
3386        let (status, waitable) = loop {
3387            store.0.suspend(if yielded {
3388                SuspendReason::Waiting {
3389                    set,
3390                    thread: caller,
3391                }
3392            } else {
3393                yielded = true;
3394                SuspendReason::YieldingToSubtask { thread: caller }
3395            })?;
3396
3397            if let Some(old_do_not_suspend) = old_do_not_suspend {
3398                store
3399                    .0
3400                    .instance_state(callee_instance)
3401                    .concurrent_state()
3402                    .do_not_suspend = old_do_not_suspend;
3403            }
3404
3405            let state = store.0.concurrent_state_mut()?;
3406
3407            log::trace!("taking event for {:?}", guest_thread.task);
3408            let event = guest_waitable.take_event(state)?;
3409            let Some(Event::Subtask { status }) = event else {
3410                bail_bug!("subtasks should only get subtask events, got {event:?}")
3411            };
3412
3413            log::trace!("status {status:?} for {:?}", guest_thread.task);
3414
3415            if status == Status::Returned {
3416                // It returned, so we can stop waiting.
3417                break (status, None);
3418            } else if async_caller {
3419                // It hasn't returned yet, but the caller is calling via an
3420                // async-lowered import, so we generate a handle for the task
3421                // waitable and return the status.
3422                let handle = store
3423                    .0
3424                    .instance_state(caller_instance)
3425                    .handle_table()
3426                    .subtask_insert_guest(guest_thread.task.rep())?;
3427                store
3428                    .0
3429                    .concurrent_state_mut()?
3430                    .get_mut(guest_thread.task)?
3431                    .common
3432                    .handle = Some(handle);
3433                break (status, Some(handle));
3434            } else {
3435                // The callee hasn't returned yet, and the caller is calling via
3436                // a sync-lowered import, so we loop and keep waiting until the
3437                // callee returns.
3438                store.0.switch_or_trap_if_may_not_suspend(caller_instance)?;
3439            }
3440        };
3441
3442        guest_waitable.join(store.0.concurrent_state_mut()?, old_set)?;
3443
3444        // Reset the current thread to point to the caller as it resumes control.
3445        store.0.set_thread(caller)?;
3446        store
3447            .0
3448            .concurrent_state_mut()?
3449            .get_mut(caller.thread)?
3450            .state = GuestThreadState::Running;
3451        log::trace!("popped current thread {guest_thread:?}; new thread is {caller:?}");
3452
3453        if let Some(storage) = storage {
3454            // The caller used a sync-lowered import to call an async-lifted
3455            // export, in which case the result, if any, has been stashed in
3456            // `GuestTask::sync_result`.
3457            let state = store.0.concurrent_state_mut()?;
3458            let task = state.get_mut(guest_thread.task)?;
3459            if let Some(result) = task.sync_result.take()? {
3460                if let Some(result) = result {
3461                    storage[0] = MaybeUninit::new(result);
3462                }
3463
3464                if task.exited && task.ready_to_delete() {
3465                    Waitable::Guest(guest_thread.task).delete_from(store.0)?;
3466                }
3467            }
3468        }
3469
3470        Ok(status.pack(waitable))
3471    }
3472
3473    /// Poll the specified future once on behalf of a guest->host call using an
3474    /// async-lowered import.
3475    ///
3476    /// If it returns `Ready`, return `Ok(None)`.  Otherwise, if it returns
3477    /// `Pending`, add it to the set of futures to be polled as part of this
3478    /// instance's event loop until it completes, and then return
3479    /// `Ok(Some(handle))` where `handle` is the waitable handle to return.
3480    ///
3481    /// Whether the future returns `Ready` immediately or later, the `lower`
3482    /// function will be used to lower the result, if any, into the guest caller's
3483    /// stack and linear memory. The `lower` function is invoked with the
3484    /// `Option<R>` param `None` if the future is cancelled. The
3485    /// `Option<TableId<HostTask>>` is passed as `Some` if the host task was
3486    /// materialized during execution and allows `lower` to delete the task if
3487    /// needed.
3488    pub(crate) fn first_poll<T: 'static, R: Send + 'static>(
3489        self,
3490        mut store: StoreContextMut<'_, T>,
3491        host_task: EnteredHostTask,
3492        future: impl Future<Output = Result<R>> + Send + 'static,
3493        lower: impl FnOnce(StoreContextMut<T>, Option<R>, bool, Option<TableId<HostTask>>) -> Result<()>
3494        + Send
3495        + 'static,
3496    ) -> Result<u32> {
3497        let token = StoreToken::new(store.as_context_mut());
3498
3499        // Create an abortable future which hooks calls to poll and manages call
3500        // context state for the future.
3501        let (join_handle, future) = JoinHandle::run(future);
3502        let mut future = Box::pin(future);
3503
3504        // Finally, poll the future.  We can use a dummy `Waker` here because
3505        // we'll add the future to `ConcurrentState::futures` and poll it
3506        // automatically from the event loop if it doesn't complete immediately
3507        // here.
3508        let poll = tls::set(store.0, || {
3509            future
3510                .as_mut()
3511                .poll(&mut Context::from_waker(&Waker::noop()))
3512        });
3513
3514        match poll {
3515            // It finished immediately; lower the result and delete the task.
3516            Poll::Ready(result) => {
3517                let result = result.transpose()?;
3518                // Check if the host task was materialized so that it can be
3519                // deleted in `lower`.
3520                let task = store.0.current_materialized_host_task()?;
3521                lower(store.as_context_mut(), result, true, task)?;
3522                return Ok(Status::Returned.pack(None));
3523            }
3524
3525            // Future isn't ready yet, so fall through.
3526            Poll::Pending => {}
3527        }
3528
3529        // The future will outlive this call frame, so materialize the deferred
3530        // host task and attach its cancellation handle before publishing it to
3531        // the event loop.
3532        let Some(task) = store.0.materialize_host_task_id()? else {
3533            bail_bug!("current thread is not a host thread")
3534        };
3535        {
3536            let state = &mut store.0.concurrent_state_mut()?.get_mut(task)?.state;
3537            assert!(matches!(state, HostTaskState::CalleeStarted));
3538            *state = HostTaskState::CalleeRunning(join_handle);
3539        }
3540
3541        // It hasn't finished yet; add the future to
3542        // `ConcurrentState::futures` so it will be polled by the event
3543        // loop and allocate a waitable handle to return to the guest.
3544
3545        // Wrap the future in a closure responsible for lowering the result into
3546        // the guest's stack and memory, as well as notifying any waiters that
3547        // the task returned.
3548        let future = Box::pin(async move {
3549            let result = match run_with_host_task_set(task, future).await? {
3550                Some(result) => Some(result?),
3551                None => None,
3552            };
3553            let on_complete = move |store: &mut dyn VMStore| {
3554                // Restore the `current_thread` to be the host so `lower` knows
3555                // how to manipulate borrows and knows which scope of borrows
3556                // to check.
3557                let mut store = token.as_context_mut(store);
3558                let old = store.0.set_thread(task)?;
3559
3560                let status = if result.is_some() {
3561                    Status::Returned
3562                } else {
3563                    Status::ReturnCancelled
3564                };
3565
3566                lower(store.as_context_mut(), result, false, Some(task))?;
3567                let state = store.0.concurrent_state_mut()?;
3568                match &mut state.get_mut(task)?.state {
3569                    // The task is already flagged as finished because it was
3570                    // cancelled. No need to transition further.
3571                    HostTaskState::CalleeDone { .. } => {}
3572
3573                    // Otherwise transition this task to the done state.
3574                    other => *other = HostTaskState::CalleeDone { cancelled: false },
3575                }
3576                Waitable::Host(task).set_event(state, Some(Event::Subtask { status }))?;
3577
3578                store.0.set_thread(old)?;
3579                Ok(())
3580            };
3581
3582            // Here we schedule a task to run on a worker fiber to do the
3583            // lowering since it may involve a call to the guest's realloc
3584            // function. This is necessary because calling the guest while
3585            // there are host embedder frames on the stack is unsound.
3586            tls::get(move |store| {
3587                store
3588                    .concurrent_state_mut()?
3589                    .push_high_priority(WorkItem::WorkerFunction(AlwaysMut::new(Box::new(
3590                        on_complete,
3591                    ))));
3592                Ok(())
3593            })
3594        });
3595
3596        // Make this task visible to the guest and then record what it
3597        // was made visible as.
3598        let caller = match host_task {
3599            Some(caller) => caller,
3600            None => bail_bug!("host task wasn't created but should have been"),
3601        };
3602        let state = store.0.concurrent_state_mut()?;
3603        state.push_future(future);
3604        let instance = state.get_mut(caller.task)?.instance;
3605        let handle = store
3606            .0
3607            .instance_state(instance)
3608            .handle_table()
3609            .subtask_insert_host(task.rep())?;
3610        store.0.concurrent_state_mut()?.get_mut(task)?.common.handle = Some(handle);
3611        log::trace!("assign {task:?} handle {handle} for {caller:?} instance {instance:?}");
3612
3613        // Restore the currently running thread to this host task's
3614        // caller. Note that the host task isn't deallocated as it's
3615        // within the store and will get deallocated later.
3616        store.0.set_thread(caller)?;
3617        Ok(Status::Started.pack(Some(handle)))
3618    }
3619
3620    /// Implements the `task.return` intrinsic, lifting the result for the
3621    /// current guest task.
3622    pub(crate) fn task_return(
3623        self,
3624        store: &mut dyn VMStore,
3625        ty: TypeTupleIndex,
3626        options: OptionsIndex,
3627        storage: &[ValRaw],
3628    ) -> Result<()> {
3629        let guest_thread = store.current_guest_thread()?;
3630        let state = store.concurrent_state_mut()?;
3631        let lift = state
3632            .get_mut(guest_thread.task)?
3633            .lift_result
3634            .take()
3635            .ok_or_else(|| Trap::TaskCancelOrReturnTwice)?;
3636        if !state.get_mut(guest_thread.task)?.result.is_none() {
3637            bail_bug!("task result unexpectedly already set");
3638        }
3639
3640        let CanonicalOptions {
3641            string_encoding,
3642            data_model,
3643            ..
3644        } = &self.id().get(store).component().env_component().options[options];
3645
3646        let invalid = ty != lift.ty
3647            || string_encoding != &lift.string_encoding
3648            || match data_model {
3649                CanonicalOptionsDataModel::LinearMemory(opts) => match opts.memory {
3650                    Some(memory) => {
3651                        let expected = lift.memory.map(|v| v.as_ptr()).unwrap_or(ptr::null_mut());
3652                        let actual = self.id().get(store).runtime_memory(memory);
3653                        expected != actual.as_ptr()
3654                    }
3655                    // Memory not specified, meaning it didn't need to be
3656                    // specified per validation, so not invalid.
3657                    None => false,
3658                },
3659                // Always invalid as this isn't supported.
3660                CanonicalOptionsDataModel::Gc { .. } => true,
3661            };
3662
3663        if invalid {
3664            bail!(Trap::TaskReturnInvalid);
3665        }
3666
3667        log::trace!("task.return for {guest_thread:?}");
3668
3669        let result = (lift.lift)(store, storage)?;
3670        self.task_complete(store, guest_thread.task, result, Status::Returned)
3671    }
3672
3673    /// Implements the `task.cancel` intrinsic.
3674    pub(crate) fn task_cancel(self, store: &mut StoreOpaque) -> Result<()> {
3675        let guest_thread = store.current_guest_thread()?;
3676        let state = store.concurrent_state_mut()?;
3677        let task = state.get_mut(guest_thread.task)?;
3678        if !task.cancel_request_delivered {
3679            bail!(Trap::TaskCancelNotCancelled);
3680        }
3681        _ = task
3682            .lift_result
3683            .take()
3684            .ok_or_else(|| Trap::TaskCancelOrReturnTwice)?;
3685
3686        if !task.result.is_none() {
3687            bail_bug!("task result should not bet set yet");
3688        }
3689
3690        log::trace!("task.cancel for {guest_thread:?}");
3691
3692        self.task_complete(
3693            store,
3694            guest_thread.task,
3695            Box::new(DummyResult),
3696            Status::ReturnCancelled,
3697        )
3698    }
3699
3700    /// Complete the specified guest task (i.e. indicate that it has either
3701    /// returned a (possibly empty) result or cancelled itself).
3702    ///
3703    /// This will return any resource borrows and notify any current or future
3704    /// waiters that the task has completed.
3705    fn task_complete(
3706        self,
3707        store: &mut StoreOpaque,
3708        guest_task: TableId<GuestTask>,
3709        result: Box<dyn Any + Send + Sync>,
3710        status: Status,
3711    ) -> Result<()> {
3712        store
3713            .component_resource_tables(Some(self))?
3714            .validate_scope_exit()?;
3715
3716        let state = store.concurrent_state_mut()?;
3717        let task = state.get_mut(guest_task)?;
3718
3719        if let Caller::Host { tx, .. } = &mut task.caller {
3720            if let Some(tx) = tx.take() {
3721                _ = tx.send(result);
3722            }
3723        } else {
3724            task.result = Some(result);
3725            Waitable::Guest(guest_task).set_event(state, Some(Event::Subtask { status }))?;
3726        }
3727
3728        Ok(())
3729    }
3730
3731    /// Implements the `waitable-set.new` intrinsic.
3732    pub(crate) fn waitable_set_new(
3733        self,
3734        store: &mut StoreOpaque,
3735        caller_instance: RuntimeComponentInstanceIndex,
3736    ) -> Result<u32> {
3737        let set = store.concurrent_state_mut()?.push(WaitableSet::default())?;
3738        let handle = store
3739            .instance_state(self.runtime_instance(caller_instance))
3740            .handle_table()
3741            .waitable_set_insert(set.rep())?;
3742        log::trace!("new waitable set {set:?} (handle {handle})");
3743        Ok(handle)
3744    }
3745
3746    /// Implements the `waitable-set.drop` intrinsic.
3747    pub(crate) fn waitable_set_drop(
3748        self,
3749        store: &mut StoreOpaque,
3750        caller_instance: RuntimeComponentInstanceIndex,
3751        set: u32,
3752    ) -> Result<()> {
3753        let rep = store
3754            .instance_state(self.runtime_instance(caller_instance))
3755            .handle_table()
3756            .waitable_set_remove(set)?;
3757
3758        log::trace!("drop waitable set {rep} (handle {set})");
3759
3760        // Note that we're careful to check for waiters _before_ deleting the
3761        // set to avoid dropping any waiters in `WaitMode::Fiber(_)`, which
3762        // would panic.  See `drop-waitable-set-with-waiters.wast` for details.
3763        if !store
3764            .concurrent_state_mut()?
3765            .get_mut(TableId::<WaitableSet>::new(rep))?
3766            .waiting
3767            .is_empty()
3768        {
3769            bail!(Trap::WaitableSetDropHasWaiters);
3770        }
3771
3772        store
3773            .concurrent_state_mut()?
3774            .delete(TableId::<WaitableSet>::new(rep))?;
3775
3776        Ok(())
3777    }
3778
3779    /// Implements the `waitable.join` intrinsic.
3780    pub(crate) fn waitable_join(
3781        self,
3782        store: &mut StoreOpaque,
3783        caller_instance: RuntimeComponentInstanceIndex,
3784        waitable_handle: u32,
3785        set_handle: u32,
3786    ) -> Result<()> {
3787        let mut instance = self.id().get_mut(store);
3788        let waitable =
3789            Waitable::from_instance(instance.as_mut(), caller_instance, waitable_handle)?;
3790
3791        let set = if set_handle == 0 {
3792            None
3793        } else {
3794            let set = instance.instance_states().0[caller_instance]
3795                .handle_table()
3796                .waitable_set_rep(set_handle)?;
3797
3798            let state = store.concurrent_state_mut()?;
3799            if let Some(old) = waitable.common(state)?.set
3800                && state.get_mut(old)?.is_sync_call_set
3801            {
3802                bail!(Trap::WaitableSyncAndAsync);
3803            }
3804
3805            Some(TableId::<WaitableSet>::new(set))
3806        };
3807
3808        log::trace!(
3809            "waitable {waitable:?} (handle {waitable_handle}) join set {set:?} (handle {set_handle})",
3810        );
3811
3812        waitable.join(store.concurrent_state_mut()?, set)
3813    }
3814
3815    /// Implements the `subtask.drop` intrinsic.
3816    pub(crate) fn subtask_drop(
3817        self,
3818        store: &mut StoreOpaque,
3819        caller_instance: RuntimeComponentInstanceIndex,
3820        task_id: u32,
3821    ) -> Result<()> {
3822        self.waitable_join(store, caller_instance, task_id, 0)?;
3823
3824        let (rep, is_host) = store
3825            .instance_state(self.runtime_instance(caller_instance))
3826            .handle_table()
3827            .subtask_remove(task_id)?;
3828
3829        let concurrent_state = store.concurrent_state_mut()?;
3830        let (waitable, delete) = if is_host {
3831            let id = TableId::<HostTask>::new(rep);
3832            let task = concurrent_state.get_mut(id)?;
3833            match &task.state {
3834                HostTaskState::CalleeRunning(_) => bail!(Trap::SubtaskDropNotResolved),
3835                HostTaskState::CalleeDone { .. } => {}
3836                HostTaskState::CalleeStarted | HostTaskState::CalleeFinished(_) => {
3837                    bail_bug!("invalid state for callee in `subtask.drop`")
3838                }
3839            }
3840
3841            (Waitable::Host(id), true)
3842        } else {
3843            let id = TableId::<GuestTask>::new(rep);
3844            let task = concurrent_state.get_mut(id)?;
3845            if task.lift_result.is_some() {
3846                bail!(Trap::SubtaskDropNotResolved);
3847            }
3848            (
3849                Waitable::Guest(id),
3850                concurrent_state.get_mut(id)?.ready_to_delete(),
3851            )
3852        };
3853
3854        waitable.common(concurrent_state)?.handle = None;
3855
3856        // If this subtask has an event that means that the terminal status of
3857        // this subtask wasn't yet received so it can't be dropped yet.
3858        if waitable.take_event(concurrent_state)?.is_some() {
3859            bail!(Trap::SubtaskDropNotResolved);
3860        }
3861
3862        if delete {
3863            waitable.delete_from(store)?;
3864        }
3865
3866        log::trace!("subtask_drop {waitable:?} (handle {task_id})");
3867        Ok(())
3868    }
3869
3870    /// Implements the `waitable-set.wait` intrinsic.
3871    pub(crate) fn waitable_set_wait(
3872        self,
3873        store: &mut StoreOpaque,
3874        options: OptionsIndex,
3875        set: u32,
3876        payload: u32,
3877    ) -> Result<u32> {
3878        let &CanonicalOptions {
3879            instance: caller_instance,
3880            ..
3881        } = &self.id().get(store).component().env_component().options[options];
3882        let caller = self.runtime_instance(caller_instance);
3883        let rep = store
3884            .instance_state(self.runtime_instance(caller_instance))
3885            .handle_table()
3886            .waitable_set_rep(set)?;
3887
3888        self.waitable_check(
3889            store,
3890            caller,
3891            WaitableCheck::Wait,
3892            WaitableCheckParams {
3893                set: TableId::new(rep),
3894                options,
3895                payload,
3896            },
3897        )
3898    }
3899
3900    /// Implements the `waitable-set.poll` intrinsic.
3901    pub(crate) fn waitable_set_poll(
3902        self,
3903        store: &mut StoreOpaque,
3904        options: OptionsIndex,
3905        set: u32,
3906        payload: u32,
3907    ) -> Result<u32> {
3908        let &CanonicalOptions {
3909            instance: caller_instance,
3910            ..
3911        } = &self.id().get(store).component().env_component().options[options];
3912        let caller = self.runtime_instance(caller_instance);
3913        let rep = store
3914            .instance_state(caller)
3915            .handle_table()
3916            .waitable_set_rep(set)?;
3917
3918        self.waitable_check(
3919            store,
3920            caller,
3921            WaitableCheck::Poll,
3922            WaitableCheckParams {
3923                set: TableId::new(rep),
3924                options,
3925                payload,
3926            },
3927        )
3928    }
3929
3930    /// Implements the `thread.index` intrinsic.
3931    pub(crate) fn thread_index(&self, store: &mut dyn VMStore) -> Result<u32> {
3932        let thread_id = store.current_guest_thread()?.thread;
3933        match store
3934            .concurrent_state_mut()?
3935            .get_mut(thread_id)?
3936            .instance_rep
3937        {
3938            Some(r) => Ok(r),
3939            None => bail_bug!("thread should have instance_rep by now"),
3940        }
3941    }
3942
3943    /// Implements the `thread.new-indirect` intrinsic.
3944    pub(crate) fn thread_new_indirect<T: 'static>(
3945        self,
3946        mut store: StoreContextMut<T>,
3947        runtime_instance: RuntimeComponentInstanceIndex,
3948        _func_ty_idx: TypeFuncIndex, // currently unused
3949        start_func_table_idx: RuntimeTableIndex,
3950        start_func_idx: u32,
3951        context: i32,
3952    ) -> Result<u32> {
3953        log::trace!("creating new thread");
3954
3955        let start_func_ty = FuncType::new(store.engine(), [ValType::I32], []);
3956        let (instance, registry) = self.id().get_mut_and_registry(store.0);
3957        let callee = instance
3958            .index_runtime_func_table(registry, start_func_table_idx, start_func_idx as u64)?
3959            .ok_or_else(|| Trap::ThreadNewIndirectUninitialized)?;
3960        if callee.type_index(store.0) != start_func_ty.type_index() {
3961            bail!(Trap::ThreadNewIndirectInvalidType);
3962        }
3963
3964        let token = StoreToken::new(store.as_context_mut());
3965        let start_func = Box::new(
3966            move |store: &mut dyn VMStore, guest_thread: QualifiedThreadId| -> Result<()> {
3967                let old_thread = store.set_thread(guest_thread)?;
3968                log::trace!(
3969                    "thread start: replaced {old_thread:?} with {guest_thread:?} as current thread"
3970                );
3971
3972                let mut store = token.as_context_mut(store);
3973                let mut params = [ValRaw::i32(context)];
3974                // Use call_unchecked rather than call or call_async, as we don't want to run the function
3975                // on a separate fiber if we're running in an async store.
3976                unsafe { callee.call_unchecked(store.as_context_mut(), &mut params)? };
3977
3978                store.0.set_thread(old_thread)?;
3979
3980                let runtime_instance = self.runtime_instance(runtime_instance);
3981
3982                // We're not returning to any caller, so we may need to switch
3983                // or trap if the instance has an outstanding sync-typed call:
3984                store
3985                    .0
3986                    .switch_or_trap_if_may_not_suspend(runtime_instance)?;
3987
3988                store
3989                    .0
3990                    .cleanup_thread(guest_thread, runtime_instance, CleanupTask::Yes)?;
3991
3992                log::trace!("explicit thread {guest_thread:?} completed");
3993                let state = store.0.concurrent_state_mut()?;
3994                if let Some(t) = old_thread.guest() {
3995                    state.get_mut(t.thread)?.state = GuestThreadState::Running;
3996                }
3997                log::trace!("thread start: restored {old_thread:?} as current thread");
3998
3999                Ok(())
4000            },
4001        );
4002
4003        let current_thread = store.0.current_guest_thread()?;
4004        let state = store.0.concurrent_state_mut()?;
4005        let parent_task = current_thread.task;
4006
4007        let new_thread = GuestThread::new_explicit(state, parent_task, start_func)?;
4008        let thread_id = state.push(new_thread)?;
4009        state.get_mut(parent_task)?.threads.insert(thread_id);
4010
4011        log::trace!("new thread with id {thread_id:?} created");
4012
4013        self.add_guest_thread_to_instance_table(thread_id, store.0, runtime_instance)
4014    }
4015
4016    pub(crate) fn resume_thread(
4017        self,
4018        store: &mut StoreOpaque,
4019        runtime_instance: RuntimeComponentInstanceIndex,
4020        thread_idx: u32,
4021        how: ResumeThread,
4022    ) -> Result<bool> {
4023        let thread_id =
4024            GuestThread::from_instance(self.id().get_mut(store), runtime_instance, thread_idx)?;
4025        let state = store.concurrent_state_mut()?;
4026        let guest_thread = QualifiedThreadId::qualify(state, thread_id)?;
4027
4028        if store.current_guest_thread()? == guest_thread {
4029            bail!(Trap::CannotResumeThread);
4030        }
4031
4032        let state = store.concurrent_state_mut()?;
4033        let thread = state.get_mut(guest_thread.thread)?;
4034        let priority = match how {
4035            ResumeThread::Promote | ResumeThread::Resume => Priority::Switch,
4036            ResumeThread::ResumeLater => Priority::Low,
4037        };
4038
4039        match (&how, &thread.state) {
4040            // Promotion is a noop unless the thread is in a ready state.
4041            (ResumeThread::Promote, GuestThreadState::Ready { .. }) => {}
4042            (ResumeThread::Promote, _) => return Ok(false),
4043
4044            // When resuming a thread it must be in a suspended state otherwise
4045            // this operation is a trap.
4046            (
4047                ResumeThread::Resume | ResumeThread::ResumeLater,
4048                GuestThreadState::NotStartedExplicit(_) | GuestThreadState::Suspended(_),
4049            ) => {}
4050            (ResumeThread::Resume | ResumeThread::ResumeLater, _) => {
4051                bail!(Trap::CannotResumeThread)
4052            }
4053        }
4054
4055        match mem::replace(&mut thread.state, GuestThreadState::Running) {
4056            GuestThreadState::NotStartedExplicit(start_func) => {
4057                log::trace!("starting thread {guest_thread:?}");
4058                let guest_call = WorkItem::GuestCall {
4059                    instance: self.runtime_instance(runtime_instance),
4060                    call: GuestCall {
4061                        thread: guest_thread,
4062                        kind: GuestCallKind::StartExplicit(Box::new(move |store| {
4063                            start_func(store, guest_thread)
4064                        })),
4065                    },
4066                };
4067                store
4068                    .concurrent_state_mut()?
4069                    .push_work_item(guest_call, priority)?;
4070            }
4071            GuestThreadState::Suspended(fiber) => {
4072                log::trace!("resuming thread {thread_id:?} that was suspended");
4073                store.concurrent_state_mut()?.push_work_item(
4074                    WorkItem::ResumeFiber {
4075                        instance: self.runtime_instance(runtime_instance),
4076                        thread: guest_thread,
4077                        fiber,
4078                    },
4079                    priority,
4080                )?;
4081            }
4082            GuestThreadState::Ready { fiber } => {
4083                log::trace!("resuming thread {thread_id:?} that was ready");
4084                thread.state = GuestThreadState::Ready { fiber };
4085                store
4086                    .concurrent_state_mut()?
4087                    .promote_thread_work_item(guest_thread)?;
4088            }
4089            other @ (GuestThreadState::NotStartedImplicit
4090            | GuestThreadState::Running
4091            | GuestThreadState::Completed) => {
4092                thread.state = other;
4093            }
4094        }
4095        Ok(true)
4096    }
4097
4098    fn add_guest_thread_to_instance_table(
4099        self,
4100        thread_id: TableId<GuestThread>,
4101        store: &mut StoreOpaque,
4102        runtime_instance: RuntimeComponentInstanceIndex,
4103    ) -> Result<u32> {
4104        let guest_id = store
4105            .instance_state(self.runtime_instance(runtime_instance))
4106            .thread_handle_table()
4107            .guest_thread_insert(thread_id.rep())?;
4108        store
4109            .concurrent_state_mut()?
4110            .get_mut(thread_id)?
4111            .instance_rep = Some(guest_id);
4112        Ok(guest_id)
4113    }
4114
4115    /// Helper function for the `thread.yield`, thread.suspend`,
4116    /// `thread.suspend-then-resume`, `thread.suspend-then-promote`,
4117    /// `thread.yield-then-resume`, and `thread.yield-then-promote` intrinsics.
4118    pub(crate) fn suspension_intrinsic(
4119        self,
4120        store: &mut StoreOpaque,
4121        caller: RuntimeComponentInstanceIndex,
4122        yielding: bool,
4123        to_thread: SuspensionTarget,
4124    ) -> Result<WaitResult> {
4125        let check_suspend = match to_thread {
4126            SuspensionTarget::Promote(thread) => {
4127                !self.resume_thread(store, caller, thread, ResumeThread::Promote)?
4128            }
4129            SuspensionTarget::Resume(thread) => {
4130                if !self.resume_thread(store, caller, thread, ResumeThread::Resume)? {
4131                    bail_bug!(
4132                        "`resume_thread` should only ever return false \
4133                         when `ResumeThread::Promote` is passed to it"
4134                    );
4135                }
4136                false
4137            }
4138            SuspensionTarget::None => true,
4139        };
4140
4141        if check_suspend && !store.switch_if_may_not_suspend(self.runtime_instance(caller))? {
4142            return if yielding {
4143                Ok(WaitResult::Completed)
4144            } else {
4145                Err(Trap::CannotBlockSyncTask.into())
4146            };
4147        }
4148
4149        let guest_thread = store.current_guest_thread()?;
4150
4151        let reason = if yielding {
4152            SuspendReason::Yielding {
4153                thread: guest_thread,
4154            }
4155        } else {
4156            SuspendReason::ExplicitlySuspending {
4157                thread: guest_thread,
4158            }
4159        };
4160
4161        store.suspend(reason)?;
4162
4163        Ok(WaitResult::Completed)
4164    }
4165
4166    /// Helper function for the `waitable-set.wait` and `waitable-set.poll` intrinsics.
4167    fn waitable_check(
4168        self,
4169        store: &mut StoreOpaque,
4170        caller: RuntimeInstance,
4171        check: WaitableCheck,
4172        params: WaitableCheckParams,
4173    ) -> Result<u32> {
4174        let guest_thread = store.current_guest_thread()?;
4175
4176        log::trace!("waitable check for {guest_thread:?}; set {:?}", params.set);
4177
4178        let state = store.concurrent_state_mut()?;
4179        let task = state.get_mut(guest_thread.task)?;
4180
4181        // If we're waiting, and there are no events immediately available,
4182        // suspend the fiber until that changes.
4183        match &check {
4184            WaitableCheck::Wait => {
4185                let set = params.set;
4186
4187                if (task.event.is_none() || matches!(task.event, Some(Event::Cancelled)))
4188                    && state.get_mut(set)?.ready.is_empty()
4189                {
4190                    store.switch_or_trap_if_may_not_suspend(caller)?;
4191
4192                    store.suspend(SuspendReason::Waiting {
4193                        set,
4194                        thread: guest_thread,
4195                    })?;
4196                }
4197            }
4198            WaitableCheck::Poll => {}
4199        }
4200
4201        log::trace!(
4202            "waitable check for {guest_thread:?}; set {:?}, part two",
4203            params.set
4204        );
4205
4206        // Deliver any pending events to the guest and return.
4207        let event = self.get_event(store, guest_thread.task, Some(params.set), false)?;
4208
4209        let (ordinal, handle, result) = match &check {
4210            WaitableCheck::Wait => {
4211                let (event, waitable) = match event {
4212                    Some(p) => p,
4213                    None => bail_bug!("event expected to be present"),
4214                };
4215                let handle = waitable.map(|(_, v)| v).unwrap_or(0);
4216                let (ordinal, result) = event.parts();
4217                (ordinal, handle, result)
4218            }
4219            WaitableCheck::Poll => {
4220                if let Some((event, waitable)) = event {
4221                    let handle = waitable.map(|(_, v)| v).unwrap_or(0);
4222                    let (ordinal, result) = event.parts();
4223                    (ordinal, handle, result)
4224                } else {
4225                    log::trace!(
4226                        "no events ready to deliver via waitable-set.poll to {:?}; set {:?}",
4227                        guest_thread.task,
4228                        params.set
4229                    );
4230                    let (ordinal, result) = Event::None.parts();
4231                    (ordinal, 0, result)
4232                }
4233            }
4234        };
4235        let memory = self.options_memory_mut(store, params.options);
4236        let ptr = crate::component::func::validate_inbounds_dynamic(
4237            &CanonicalAbiInfo::POINTER_PAIR,
4238            memory,
4239            &ValRaw::u32(params.payload),
4240        )?;
4241        memory[ptr + 0..][..4].copy_from_slice(&handle.to_le_bytes());
4242        memory[ptr + 4..][..4].copy_from_slice(&result.to_le_bytes());
4243        Ok(ordinal)
4244    }
4245
4246    /// Implements the `subtask.cancel` intrinsic.
4247    pub(crate) fn subtask_cancel(
4248        self,
4249        store: &mut StoreOpaque,
4250        caller_instance: RuntimeComponentInstanceIndex,
4251        async_: bool,
4252        task_id: u32,
4253    ) -> Result<u32> {
4254        let (rep, is_host) = store
4255            .instance_state(self.runtime_instance(caller_instance))
4256            .handle_table()
4257            .subtask_rep(task_id)?;
4258        let waitable = if is_host {
4259            Waitable::Host(TableId::<HostTask>::new(rep))
4260        } else {
4261            Waitable::Guest(TableId::<GuestTask>::new(rep))
4262        };
4263        let concurrent_state = store.concurrent_state_mut()?;
4264
4265        log::trace!("subtask_cancel {waitable:?} (handle {task_id}; async {async_})");
4266
4267        waitable.trap_if_in_waitable_set(concurrent_state)?;
4268
4269        let needs_block;
4270        if let Waitable::Host(host_task) = waitable {
4271            let state = &mut concurrent_state.get_mut(host_task)?.state;
4272            match mem::replace(state, HostTaskState::CalleeDone { cancelled: true }) {
4273                // If the callee is still running, signal an abort is requested.
4274                //
4275                // After cancelling this falls through to block waiting for the
4276                // host task to actually finish assuming that `async_` is false.
4277                // This blocking behavior resolves the race of `handle.abort()`
4278                // with the task actually getting cancelled or finishing.
4279                HostTaskState::CalleeRunning(handle) => {
4280                    handle.abort();
4281                    needs_block = true;
4282                }
4283
4284                // Cancellation was already requested, so fail as the task can't
4285                // be cancelled twice.
4286                HostTaskState::CalleeDone { cancelled } => {
4287                    if cancelled {
4288                        bail!(Trap::SubtaskCancelAfterTerminal);
4289                    } else {
4290                        // The callee is already done so there's no need to
4291                        // block further for an event.
4292                        needs_block = false;
4293                    }
4294                }
4295
4296                // These states should not be possible for a subtask that's
4297                // visible from the guest, so trap here.
4298                HostTaskState::CalleeStarted | HostTaskState::CalleeFinished(_) => {
4299                    bail_bug!("invalid states for host callee")
4300                }
4301            }
4302        } else {
4303            let guest_task = TableId::<GuestTask>::new(rep);
4304            let task = concurrent_state.get_mut(guest_task)?;
4305            if !task.already_lowered_parameters() {
4306                store.cancel_guest_subtask_without_lowered_parameters(
4307                    self.runtime_instance(caller_instance),
4308                    guest_task,
4309                )?;
4310                return Ok(Status::StartCancelled as u32);
4311            } else if !task.returned_or_cancelled() {
4312                // Started, but not yet returned or cancelled; send the
4313                // `CANCELLED` event
4314                //
4315                // Note that this might overwrite an event that was set earlier
4316                // (e.g. `Event::None` if the task is yielding, or
4317                // `Event::Cancelled` if it was already cancelled), but that's
4318                // okay -- this should supersede the previous state.
4319                task.event = Some(Event::Cancelled);
4320                let runtime_instance = task.instance;
4321                for thread in task.threads.clone() {
4322                    let thread = QualifiedThreadId {
4323                        task: guest_task,
4324                        thread,
4325                    };
4326                    let thread_mut = concurrent_state.get_mut(thread.thread)?;
4327
4328                    let yield_ = |store: &mut StoreOpaque| {
4329                        // While we're yielding, temporarily set
4330                        // `do_not_suspend` to false on the subtask's instance
4331                        // since we'll be getting control back if it does
4332                        // suspend.
4333                        let state = store.instance_state(runtime_instance).concurrent_state();
4334                        let old_do_not_suspend = state.do_not_suspend;
4335                        state.do_not_suspend = false;
4336
4337                        let caller = store.current_guest_thread()?;
4338
4339                        // Temporarily add the waitable to the caller's
4340                        // `sync_call_set` to ensure that (1) it isn't already
4341                        // part of a different set and (2) it can't be added to
4342                        // a different set while we yield to the subtask.
4343                        let state = store.concurrent_state_mut()?;
4344                        let set = state.get_mut(caller.thread)?.sync_call_set;
4345                        waitable.join(state, Some(set))?;
4346
4347                        store.suspend(SuspendReason::YieldingToSubtask { thread: caller })?;
4348
4349                        let state = store.concurrent_state_mut()?;
4350                        waitable.join(state, None)?;
4351
4352                        store
4353                            .instance_state(runtime_instance)
4354                            .concurrent_state()
4355                            .do_not_suspend = old_do_not_suspend;
4356
4357                        Ok::<(), crate::Error>(())
4358                    };
4359
4360                    match thread_mut.wake_on_cancel.take() {
4361                        WakeOnCancel::Waiting(set) => {
4362                            // The thread is in a cancellable wait, so wake it up:
4363                            let item = match concurrent_state.get_mut(set)?.waiting.remove(&thread)
4364                            {
4365                                Some(WaitMode::Callback(instance)) => WorkItem::GuestCall {
4366                                    instance: runtime_instance,
4367                                    call: GuestCall {
4368                                        thread,
4369                                        kind: GuestCallKind::DeliverEvent {
4370                                            instance,
4371                                            set: None,
4372                                        },
4373                                    },
4374                                },
4375                                other => bail_bug!(
4376                                    "expected `Some(WaitMode::Callback(_))`; got `{other:?}`"
4377                                ),
4378                            };
4379                            concurrent_state.set_switch_item(item)?;
4380
4381                            yield_(store)?;
4382
4383                            break;
4384                        }
4385                        WakeOnCancel::Yielding => {
4386                            if concurrent_state.promote_thread_work_item(thread)? {
4387                                yield_(store)?;
4388                                break;
4389                            } else {
4390                                bail_bug!("thread with `WakeOnCancel::Yielding` not promotable");
4391                            }
4392                        }
4393                        WakeOnCancel::None => {}
4394                    }
4395                }
4396
4397                // Guest tasks need to block if they have not yet returned or
4398                // cancelled, even as a result of the event delivery above.
4399                needs_block = !store
4400                    .concurrent_state_mut()?
4401                    .get_mut(guest_task)?
4402                    .returned_or_cancelled()
4403            } else {
4404                needs_block = false;
4405            }
4406        };
4407
4408        // If we need to block waiting on the terminal status of this subtask
4409        // then return immediately in `async` mode, or otherwise wait for the
4410        // event to get signaled through the store.
4411        if needs_block {
4412            if async_ {
4413                return Ok(BLOCKED);
4414            }
4415
4416            // Save and later restore `next_switch_item` during a sync cancel so
4417            // we don't try to switch to it while blocking.
4418            let old_next_switch_item = {
4419                let state = store.concurrent_state_mut()?;
4420                let item = state.next_switch_item.take();
4421                // Note that we store it in the table here rather than directly
4422                // in a local variable to ensure the fiber is disposed of
4423                // properly if we end up trapping or panicking.
4424                state.push(item)?
4425            };
4426
4427            // Wait for this waitable to get signaled with its terminal
4428            // status. Once that's done fall through to the shared code.
4429            store.wait_for_event(self.runtime_instance(caller_instance), waitable)?;
4430
4431            let state = store.concurrent_state_mut()?;
4432            state.next_switch_item = state.delete(old_next_switch_item)?;
4433
4434            // .. fall through to determine what event's in store for us.
4435        }
4436
4437        let event = waitable.take_event(store.concurrent_state_mut()?)?;
4438        if let Some(Event::Subtask {
4439            status: status @ (Status::Returned | Status::ReturnCancelled),
4440        }) = event
4441        {
4442            Ok(status as u32)
4443        } else {
4444            bail!(Trap::SubtaskCancelAfterTerminal);
4445        }
4446    }
4447}
4448
4449/// Trait representing component model ABI async intrinsics and fused adapter
4450/// helper functions.
4451///
4452/// SAFETY (callers): Most of the methods in this trait accept raw pointers,
4453/// which must be valid for at least the duration of the call (and possibly for
4454/// as long as the relevant guest task exists, in the case of `*mut VMFuncRef`
4455/// pointers used for async calls).
4456pub trait VMComponentAsyncStore {
4457    /// A helper function for fused adapter modules involving calls where the
4458    /// one of the caller or callee is async.
4459    ///
4460    /// This helper is not used when the caller and callee both use the sync
4461    /// ABI, only when at least one is async is this used.
4462    unsafe fn prepare_call(
4463        &mut self,
4464        instance: Instance,
4465        memory: *mut VMMemoryDefinition,
4466        start: NonNull<VMFuncRef>,
4467        return_: NonNull<VMFuncRef>,
4468        caller_instance: RuntimeComponentInstanceIndex,
4469        callee_instance: RuntimeComponentInstanceIndex,
4470        task_return_type: TypeTupleIndex,
4471        callee_async: bool,
4472        string_encoding: StringEncoding,
4473        result_count: u32,
4474        storage: *mut ValRaw,
4475        storage_len: usize,
4476    ) -> Result<()>;
4477
4478    /// A helper function for fused adapter modules involving calls where the
4479    /// caller is sync-lowered but the callee is async-lifted.
4480    unsafe fn sync_start(
4481        &mut self,
4482        instance: Instance,
4483        callback: *mut VMFuncRef,
4484        callee: NonNull<VMFuncRef>,
4485        param_count: u32,
4486        storage: *mut MaybeUninit<ValRaw>,
4487        storage_len: usize,
4488    ) -> Result<()>;
4489
4490    /// A helper function for fused adapter modules involving calls where the
4491    /// caller is async-lowered.
4492    unsafe fn async_start(
4493        &mut self,
4494        instance: Instance,
4495        callback: *mut VMFuncRef,
4496        post_return: *mut VMFuncRef,
4497        callee: NonNull<VMFuncRef>,
4498        param_count: u32,
4499        result_count: u32,
4500        flags: u32,
4501    ) -> Result<u32>;
4502
4503    /// The `future.write` intrinsic.
4504    fn future_write(
4505        &mut self,
4506        instance: Instance,
4507        caller: RuntimeComponentInstanceIndex,
4508        ty: TypeFutureTableIndex,
4509        options: OptionsIndex,
4510        future: u32,
4511        address: u32,
4512    ) -> Result<u32>;
4513
4514    /// The `future.read` intrinsic.
4515    fn future_read(
4516        &mut self,
4517        instance: Instance,
4518        caller: RuntimeComponentInstanceIndex,
4519        ty: TypeFutureTableIndex,
4520        options: OptionsIndex,
4521        future: u32,
4522        address: u32,
4523    ) -> Result<u32>;
4524
4525    /// The `future.drop-writable` intrinsic.
4526    fn future_drop_writable(
4527        &mut self,
4528        instance: Instance,
4529        ty: TypeFutureTableIndex,
4530        writer: u32,
4531    ) -> Result<()>;
4532
4533    /// The `stream.write` intrinsic.
4534    fn stream_write(
4535        &mut self,
4536        instance: Instance,
4537        caller: RuntimeComponentInstanceIndex,
4538        ty: TypeStreamTableIndex,
4539        options: OptionsIndex,
4540        stream: u32,
4541        address: u32,
4542        count: u32,
4543    ) -> Result<u32>;
4544
4545    /// The `stream.read` intrinsic.
4546    fn stream_read(
4547        &mut self,
4548        instance: Instance,
4549        caller: RuntimeComponentInstanceIndex,
4550        ty: TypeStreamTableIndex,
4551        options: OptionsIndex,
4552        stream: u32,
4553        address: u32,
4554        count: u32,
4555    ) -> Result<u32>;
4556
4557    /// The "fast-path" implementation of the `stream.write` intrinsic for
4558    /// "flat" (i.e. memcpy-able) payloads.
4559    fn flat_stream_write(
4560        &mut self,
4561        instance: Instance,
4562        caller: RuntimeComponentInstanceIndex,
4563        ty: TypeStreamTableIndex,
4564        options: OptionsIndex,
4565        payload_size: u32,
4566        payload_align: u32,
4567        stream: u32,
4568        address: u32,
4569        count: u32,
4570    ) -> Result<u32>;
4571
4572    /// The "fast-path" implementation of the `stream.read` intrinsic for "flat"
4573    /// (i.e. memcpy-able) payloads.
4574    fn flat_stream_read(
4575        &mut self,
4576        instance: Instance,
4577        caller: RuntimeComponentInstanceIndex,
4578        ty: TypeStreamTableIndex,
4579        options: OptionsIndex,
4580        payload_size: u32,
4581        payload_align: u32,
4582        stream: u32,
4583        address: u32,
4584        count: u32,
4585    ) -> Result<u32>;
4586
4587    /// The `stream.drop-writable` intrinsic.
4588    fn stream_drop_writable(
4589        &mut self,
4590        instance: Instance,
4591        ty: TypeStreamTableIndex,
4592        writer: u32,
4593    ) -> Result<()>;
4594
4595    /// The `error-context.debug-message` intrinsic.
4596    fn error_context_debug_message(
4597        &mut self,
4598        instance: Instance,
4599        ty: TypeComponentLocalErrorContextTableIndex,
4600        options: OptionsIndex,
4601        err_ctx_handle: u32,
4602        debug_msg_address: u32,
4603    ) -> Result<()>;
4604
4605    /// The `thread.new-indirect` intrinsic
4606    fn thread_new_indirect(
4607        &mut self,
4608        instance: Instance,
4609        caller: RuntimeComponentInstanceIndex,
4610        func_ty_idx: TypeFuncIndex,
4611        start_func_table_idx: RuntimeTableIndex,
4612        start_func_idx: u32,
4613        context: i32,
4614    ) -> Result<u32>;
4615}
4616
4617/// SAFETY: See trait docs.
4618impl<T: 'static> VMComponentAsyncStore for StoreInner<T> {
4619    unsafe fn prepare_call(
4620        &mut self,
4621        instance: Instance,
4622        memory: *mut VMMemoryDefinition,
4623        start: NonNull<VMFuncRef>,
4624        return_: NonNull<VMFuncRef>,
4625        caller_instance: RuntimeComponentInstanceIndex,
4626        callee_instance: RuntimeComponentInstanceIndex,
4627        task_return_type: TypeTupleIndex,
4628        callee_async: bool,
4629        string_encoding: StringEncoding,
4630        result_count_or_max_if_async: u32,
4631        storage: *mut ValRaw,
4632        storage_len: usize,
4633    ) -> Result<()> {
4634        // SAFETY: The `wasmtime_cranelift`-generated code that calls
4635        // this method will have ensured that `storage` is a valid
4636        // pointer containing at least `storage_len` items.
4637        let params = unsafe { core::slice::from_raw_parts(storage, storage_len) }.to_vec();
4638
4639        unsafe {
4640            instance.prepare_call(
4641                StoreContextMut(self),
4642                start,
4643                return_,
4644                caller_instance,
4645                callee_instance,
4646                task_return_type,
4647                callee_async,
4648                memory,
4649                string_encoding,
4650                match result_count_or_max_if_async {
4651                    PREPARE_ASYNC_NO_RESULT => CallerInfo::Async {
4652                        params,
4653                        has_result: false,
4654                    },
4655                    PREPARE_ASYNC_WITH_RESULT => CallerInfo::Async {
4656                        params,
4657                        has_result: true,
4658                    },
4659                    result_count => CallerInfo::Sync {
4660                        params,
4661                        result_count,
4662                    },
4663                },
4664            )
4665        }
4666    }
4667
4668    unsafe fn sync_start(
4669        &mut self,
4670        instance: Instance,
4671        callback: *mut VMFuncRef,
4672        callee: NonNull<VMFuncRef>,
4673        param_count: u32,
4674        storage: *mut MaybeUninit<ValRaw>,
4675        storage_len: usize,
4676    ) -> Result<()> {
4677        unsafe {
4678            instance
4679                .start_call(
4680                    StoreContextMut(self),
4681                    callback,
4682                    ptr::null_mut(),
4683                    callee,
4684                    param_count,
4685                    1,
4686                    START_FLAG_ASYNC_CALLEE,
4687                    // SAFETY: The `wasmtime_cranelift`-generated code that calls
4688                    // this method will have ensured that `storage` is a valid
4689                    // pointer containing at least `storage_len` items.
4690                    Some(core::slice::from_raw_parts_mut(storage, storage_len)),
4691                )
4692                .map(drop)
4693        }
4694    }
4695
4696    unsafe fn async_start(
4697        &mut self,
4698        instance: Instance,
4699        callback: *mut VMFuncRef,
4700        post_return: *mut VMFuncRef,
4701        callee: NonNull<VMFuncRef>,
4702        param_count: u32,
4703        result_count: u32,
4704        flags: u32,
4705    ) -> Result<u32> {
4706        unsafe {
4707            instance.start_call(
4708                StoreContextMut(self),
4709                callback,
4710                post_return,
4711                callee,
4712                param_count,
4713                result_count,
4714                flags,
4715                None,
4716            )
4717        }
4718    }
4719
4720    fn future_write(
4721        &mut self,
4722        instance: Instance,
4723        caller: RuntimeComponentInstanceIndex,
4724        ty: TypeFutureTableIndex,
4725        options: OptionsIndex,
4726        future: u32,
4727        address: u32,
4728    ) -> Result<u32> {
4729        instance
4730            .guest_write(
4731                StoreContextMut(self),
4732                caller,
4733                TransmitIndex::Future(ty),
4734                options,
4735                None,
4736                future,
4737                address,
4738                1,
4739            )
4740            .map(|result| result.encode())
4741    }
4742
4743    fn future_read(
4744        &mut self,
4745        instance: Instance,
4746        caller: RuntimeComponentInstanceIndex,
4747        ty: TypeFutureTableIndex,
4748        options: OptionsIndex,
4749        future: u32,
4750        address: u32,
4751    ) -> Result<u32> {
4752        instance
4753            .guest_read(
4754                StoreContextMut(self),
4755                caller,
4756                TransmitIndex::Future(ty),
4757                options,
4758                None,
4759                future,
4760                address,
4761                1,
4762            )
4763            .map(|result| result.encode())
4764    }
4765
4766    fn stream_write(
4767        &mut self,
4768        instance: Instance,
4769        caller: RuntimeComponentInstanceIndex,
4770        ty: TypeStreamTableIndex,
4771        options: OptionsIndex,
4772        stream: u32,
4773        address: u32,
4774        count: u32,
4775    ) -> Result<u32> {
4776        instance
4777            .guest_write(
4778                StoreContextMut(self),
4779                caller,
4780                TransmitIndex::Stream(ty),
4781                options,
4782                None,
4783                stream,
4784                address,
4785                count,
4786            )
4787            .map(|result| result.encode())
4788    }
4789
4790    fn stream_read(
4791        &mut self,
4792        instance: Instance,
4793        caller: RuntimeComponentInstanceIndex,
4794        ty: TypeStreamTableIndex,
4795        options: OptionsIndex,
4796        stream: u32,
4797        address: u32,
4798        count: u32,
4799    ) -> Result<u32> {
4800        instance
4801            .guest_read(
4802                StoreContextMut(self),
4803                caller,
4804                TransmitIndex::Stream(ty),
4805                options,
4806                None,
4807                stream,
4808                address,
4809                count,
4810            )
4811            .map(|result| result.encode())
4812    }
4813
4814    fn future_drop_writable(
4815        &mut self,
4816        instance: Instance,
4817        ty: TypeFutureTableIndex,
4818        writer: u32,
4819    ) -> Result<()> {
4820        instance.guest_drop_writable(self, TransmitIndex::Future(ty), writer)
4821    }
4822
4823    fn flat_stream_write(
4824        &mut self,
4825        instance: Instance,
4826        caller: RuntimeComponentInstanceIndex,
4827        ty: TypeStreamTableIndex,
4828        options: OptionsIndex,
4829        payload_size: u32,
4830        payload_align: u32,
4831        stream: u32,
4832        address: u32,
4833        count: u32,
4834    ) -> Result<u32> {
4835        instance
4836            .guest_write(
4837                StoreContextMut(self),
4838                caller,
4839                TransmitIndex::Stream(ty),
4840                options,
4841                Some(FlatAbi {
4842                    size: payload_size,
4843                    align: payload_align,
4844                }),
4845                stream,
4846                address,
4847                count,
4848            )
4849            .map(|result| result.encode())
4850    }
4851
4852    fn flat_stream_read(
4853        &mut self,
4854        instance: Instance,
4855        caller: RuntimeComponentInstanceIndex,
4856        ty: TypeStreamTableIndex,
4857        options: OptionsIndex,
4858        payload_size: u32,
4859        payload_align: u32,
4860        stream: u32,
4861        address: u32,
4862        count: u32,
4863    ) -> Result<u32> {
4864        instance
4865            .guest_read(
4866                StoreContextMut(self),
4867                caller,
4868                TransmitIndex::Stream(ty),
4869                options,
4870                Some(FlatAbi {
4871                    size: payload_size,
4872                    align: payload_align,
4873                }),
4874                stream,
4875                address,
4876                count,
4877            )
4878            .map(|result| result.encode())
4879    }
4880
4881    fn stream_drop_writable(
4882        &mut self,
4883        instance: Instance,
4884        ty: TypeStreamTableIndex,
4885        writer: u32,
4886    ) -> Result<()> {
4887        instance.guest_drop_writable(self, TransmitIndex::Stream(ty), writer)
4888    }
4889
4890    fn error_context_debug_message(
4891        &mut self,
4892        instance: Instance,
4893        ty: TypeComponentLocalErrorContextTableIndex,
4894        options: OptionsIndex,
4895        err_ctx_handle: u32,
4896        debug_msg_address: u32,
4897    ) -> Result<()> {
4898        instance.error_context_debug_message(
4899            StoreContextMut(self),
4900            ty,
4901            options,
4902            err_ctx_handle,
4903            debug_msg_address,
4904        )
4905    }
4906
4907    fn thread_new_indirect(
4908        &mut self,
4909        instance: Instance,
4910        caller: RuntimeComponentInstanceIndex,
4911        func_ty_idx: TypeFuncIndex,
4912        start_func_table_idx: RuntimeTableIndex,
4913        start_func_idx: u32,
4914        context: i32,
4915    ) -> Result<u32> {
4916        instance.thread_new_indirect(
4917            StoreContextMut(self),
4918            caller,
4919            func_ty_idx,
4920            start_func_table_idx,
4921            start_func_idx,
4922            context,
4923        )
4924    }
4925}
4926
4927type HostTaskFuture = Pin<Box<dyn Future<Output = Result<()>> + Send + 'static>>;
4928
4929/// Runs the given future with the current thread set to `task` each time it is
4930/// polled.
4931async fn run_with_host_task_set<F>(task: TableId<HostTask>, future: F) -> Result<F::Output>
4932where
4933    F: Future,
4934{
4935    let mut future = pin!(future);
4936    future::poll_fn(|cx| {
4937        let old_thread = match tls::get(|store| store.set_thread(task)) {
4938            Ok(thread) => thread,
4939            Err(error) => return Poll::Ready(Err(error)),
4940        };
4941        let result = future.as_mut().poll(cx);
4942        match tls::get(|store| store.set_thread(old_thread)) {
4943            Ok(_) => result.map(Ok),
4944            Err(error) => Poll::Ready(Err(error)),
4945        }
4946    })
4947    .await
4948}
4949
4950/// Represents the state of a pending host task.
4951///
4952/// This is used to represent tasks when the guest calls into the host.
4953pub(crate) struct HostTask {
4954    common: WaitableCommon,
4955
4956    /// State of borrows/etc the host needs to track. Used when the guest passes
4957    /// borrows to the host, for example.
4958    call_context: CallContext,
4959
4960    state: HostTaskState,
4961
4962    group: TaskGroupId,
4963}
4964
4965enum HostTaskState {
4966    /// A host task has been created and it's considered "started".
4967    ///
4968    /// The host task has yet to enter `first_poll` or `poll_and_block` which
4969    /// is where this will get updated further.
4970    CalleeStarted,
4971
4972    /// State used for tasks in `first_poll` meaning that the guest did an async
4973    /// lower of a host async function which is blocked. The specified handle is
4974    /// linked to the future in the main `FuturesUnordered` of a store which is
4975    /// used to cancel it if the guest requests cancellation.
4976    CalleeRunning(JoinHandle),
4977
4978    /// Terminal state used for tasks in `poll_and_block` to store the result of
4979    /// their computation. Note that this state is not used for tasks in
4980    /// `first_poll`.
4981    CalleeFinished(LiftedResult),
4982
4983    /// Terminal state for host tasks meaning that the task was cancelled or the
4984    /// result was taken.
4985    CalleeDone { cancelled: bool },
4986}
4987
4988impl HostTask {
4989    fn new(
4990        concurrent_state: &mut ConcurrentState,
4991        state: HostTaskState,
4992        caller: QualifiedThreadId,
4993    ) -> Result<Self> {
4994        let group = concurrent_state.get_mut(caller.task)?.group;
4995        concurrent_state.increment_group_ref_count(group)?;
4996
4997        Ok(Self {
4998            common: WaitableCommon::default(),
4999            call_context: CallContext::default(),
5000            state,
5001            group,
5002        })
5003    }
5004}
5005
5006impl TableDebug for HostTask {
5007    fn type_name() -> &'static str {
5008        "HostTask"
5009    }
5010}
5011
5012type CallbackFn = Box<dyn Fn(&mut dyn VMStore, Event, u32) -> Result<u32> + Send + Sync + 'static>;
5013
5014/// Represents the caller of a given guest task.
5015enum Caller {
5016    /// The host called the guest task.
5017    Host {
5018        /// If present, may be used to deliver the result.
5019        tx: Option<oneshot::Sender<LiftedResult>>,
5020        /// If true, there's a host future that must be dropped before the task
5021        /// can be deleted.
5022        host_future_present: bool,
5023        /// The host task which called into the guest, or `None` for a call from
5024        /// the top-level host. The task may belong to an entirely unrelated
5025        /// top-level component instance than the one the host called into.
5026        caller: Option<TableId<HostTask>>,
5027    },
5028    /// Another guest thread called the guest task
5029    Guest {
5030        /// The id of the caller
5031        thread: QualifiedThreadId,
5032    },
5033}
5034
5035/// Represents a closure and related canonical ABI parameters required to
5036/// validate a `task.return` call at runtime and lift the result.
5037struct LiftResult {
5038    lift: RawLift,
5039    ty: TypeTupleIndex,
5040    memory: Option<SendSyncPtr<VMMemoryDefinition>>,
5041    string_encoding: StringEncoding,
5042}
5043
5044/// The table ID for a guest thread, qualified by the task to which it belongs.
5045///
5046/// This exists to minimize table lookups and the necessity to pass stores around mutably
5047/// for the common case of identifying the task to which a thread belongs.
5048#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
5049pub(crate) struct QualifiedThreadId {
5050    task: TableId<GuestTask>,
5051    thread: TableId<GuestThread>,
5052}
5053
5054impl QualifiedThreadId {
5055    fn qualify(
5056        state: &mut ConcurrentState,
5057        thread: TableId<GuestThread>,
5058    ) -> Result<QualifiedThreadId> {
5059        Ok(QualifiedThreadId {
5060            task: state.get_mut(thread)?.parent_task,
5061            thread,
5062        })
5063    }
5064}
5065
5066impl fmt::Debug for QualifiedThreadId {
5067    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
5068        f.debug_tuple("QualifiedThreadId")
5069            .field(&self.task.rep())
5070            .field(&self.thread.rep())
5071            .finish()
5072    }
5073}
5074
5075enum GuestThreadState {
5076    NotStartedImplicit,
5077    NotStartedExplicit(
5078        Box<dyn FnOnce(&mut dyn VMStore, QualifiedThreadId) -> Result<()> + Send + Sync>,
5079    ),
5080    Running,
5081    Suspended(StoreFiber<'static>),
5082    Ready {
5083        fiber: StoreFiber<'static>,
5084    },
5085    Completed,
5086}
5087
5088impl fmt::Debug for GuestThreadState {
5089    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
5090        match self {
5091            Self::NotStartedImplicit => f.debug_tuple("NotStartedImplicit").finish(),
5092            Self::NotStartedExplicit(_) => f.debug_tuple("NotStartedExplicit").finish(),
5093            Self::Running => f.debug_tuple("Running").finish(),
5094            Self::Suspended(_) => f.debug_tuple("Suspended").finish(),
5095            Self::Ready { .. } => f.debug_struct("Ready").finish(),
5096            Self::Completed => f.debug_tuple("Completed").finish(),
5097        }
5098    }
5099}
5100
5101#[derive(Copy, Clone, PartialEq, Eq, Debug)]
5102enum WakeOnCancel {
5103    None,
5104    Waiting(TableId<WaitableSet>),
5105    Yielding,
5106}
5107
5108impl WakeOnCancel {
5109    fn is_none(self) -> bool {
5110        matches!(self, WakeOnCancel::None)
5111    }
5112
5113    fn replace(&mut self, other: WakeOnCancel) -> Self {
5114        let old = *self;
5115        *self = other;
5116        old
5117    }
5118
5119    fn take(&mut self) -> Self {
5120        self.replace(WakeOnCancel::None)
5121    }
5122}
5123
5124pub struct GuestThread {
5125    /// Context-local state used to implement the `context.{get,set}`
5126    /// intrinsics.
5127    context: [u32; NUM_COMPONENT_CONTEXT_SLOTS],
5128    /// The owning guest task.
5129    parent_task: TableId<GuestTask>,
5130    /// If non-`None`, indicates that the thread is currently either waiting a
5131    /// waitable set or yielding but may be cancelled and woken immediately.
5132    wake_on_cancel: WakeOnCancel,
5133    /// The execution state of this guest thread
5134    state: GuestThreadState,
5135    /// The index of this thread in the component instance's handle table.
5136    /// This must always be `Some` after initialization.
5137    instance_rep: Option<u32>,
5138    /// Scratch waitable set used to watch subtasks during synchronous calls.
5139    sync_call_set: TableId<WaitableSet>,
5140    /// The old value of `do_not_suspend` prior to the sync-typed task for which
5141    /// this thread was created, if relevant.
5142    old_do_not_suspend: Option<bool>,
5143}
5144
5145impl GuestThread {
5146    /// Retrieve the `GuestThread` corresponding to the specified guest-visible
5147    /// handle.
5148    fn from_instance(
5149        state: Pin<&mut ComponentInstance>,
5150        caller_instance: RuntimeComponentInstanceIndex,
5151        guest_thread: u32,
5152    ) -> Result<TableId<Self>> {
5153        let rep = state.instance_states().0[caller_instance]
5154            .thread_handle_table()
5155            .guest_thread_rep(guest_thread)?;
5156        Ok(TableId::new(rep))
5157    }
5158
5159    fn new_implicit(state: &mut ConcurrentState, parent_task: TableId<GuestTask>) -> Result<Self> {
5160        let sync_call_set = state.push(WaitableSet {
5161            is_sync_call_set: true,
5162            ..WaitableSet::default()
5163        })?;
5164        Ok(Self {
5165            context: [0; NUM_COMPONENT_CONTEXT_SLOTS],
5166            parent_task,
5167            wake_on_cancel: WakeOnCancel::None,
5168            state: GuestThreadState::NotStartedImplicit,
5169            instance_rep: None,
5170            sync_call_set,
5171            old_do_not_suspend: None,
5172        })
5173    }
5174
5175    fn new_explicit(
5176        state: &mut ConcurrentState,
5177        parent_task: TableId<GuestTask>,
5178        start_func: Box<
5179            dyn FnOnce(&mut dyn VMStore, QualifiedThreadId) -> Result<()> + Send + Sync,
5180        >,
5181    ) -> Result<Self> {
5182        let sync_call_set = state.push(WaitableSet {
5183            is_sync_call_set: true,
5184            ..WaitableSet::default()
5185        })?;
5186        Ok(Self {
5187            context: [0; NUM_COMPONENT_CONTEXT_SLOTS],
5188            parent_task,
5189            wake_on_cancel: WakeOnCancel::None,
5190            state: GuestThreadState::NotStartedExplicit(start_func),
5191            instance_rep: None,
5192            sync_call_set,
5193            old_do_not_suspend: None,
5194        })
5195    }
5196}
5197
5198impl TableDebug for GuestThread {
5199    fn type_name() -> &'static str {
5200        "GuestThread"
5201    }
5202}
5203
5204enum SyncResult {
5205    NotProduced,
5206    Produced(Option<ValRaw>),
5207    Taken,
5208}
5209
5210impl SyncResult {
5211    fn take(&mut self) -> Result<Option<Option<ValRaw>>> {
5212        Ok(match mem::replace(self, SyncResult::Taken) {
5213            SyncResult::NotProduced => None,
5214            SyncResult::Produced(val) => Some(val),
5215            SyncResult::Taken => {
5216                bail_bug!("attempted to take a synchronous result that was already taken")
5217            }
5218        })
5219    }
5220}
5221
5222#[derive(Debug)]
5223enum HostFutureState {
5224    NotApplicable,
5225    Live,
5226    Dropped,
5227}
5228
5229/// Represents a pending guest task.
5230pub(crate) struct GuestTask {
5231    /// See `WaitableCommon`
5232    common: WaitableCommon,
5233    /// Closure to lower the parameters passed to this task.
5234    lower_params: Option<RawLower>,
5235    /// See `LiftResult`
5236    lift_result: Option<LiftResult>,
5237    /// A place to stash the type-erased lifted result if it can't be delivered
5238    /// immediately.
5239    result: Option<LiftedResult>,
5240    /// Closure to call the callback function for an async-lifted export, if
5241    /// provided.
5242    callback: Option<CallbackFn>,
5243    /// See `Caller`
5244    caller: Caller,
5245    /// Borrow state for this task.
5246    ///
5247    /// Keeps track of `borrow<T>` received to this task to ensure that
5248    /// everything is dropped by the time it exits.
5249    call_context: CallContext,
5250    /// A place to stash the lowered result for a sync-to-async call until it
5251    /// can be returned to the caller.
5252    sync_result: SyncResult,
5253    /// Whether or not the task has been cancelled (i.e. whether the
5254    /// cancellation request has been delivered to the task, and thus whether
5255    /// the task is permitted to call `task.cancel`).
5256    cancel_request_delivered: bool,
5257    /// Whether or not we've sent a `Status::Starting` event to any current or
5258    /// future waiters for this waitable.
5259    starting_sent: bool,
5260    /// The runtime instance to which the exported function for this guest task
5261    /// belongs.
5262    ///
5263    /// Note that the task may do a sync->sync call via a fused adapter which
5264    /// results in that task executing code in a different instance, and it may
5265    /// call host functions and intrinsics from that other instance.
5266    instance: RuntimeInstance,
5267    /// If present, a pending `Event::None` or `Event::Cancelled` to be
5268    /// delivered to this task.
5269    event: Option<Event>,
5270    /// Whether or not the task has exited.
5271    exited: bool,
5272    /// Threads belonging to this task
5273    threads: HashSet<TableId<GuestThread>>,
5274    /// The state of the host future that represents an async task, which must
5275    /// be dropped before we can delete the task.
5276    host_future_state: HostFutureState,
5277    /// Indicates whether this task was created for a call to an async-typed
5278    /// export.
5279    async_typed: bool,
5280    /// Indicates whether this task was created for a call to an async-lifted
5281    /// export.
5282    async_lifted: bool,
5283
5284    decremented_interesting_task_count: bool,
5285
5286    group: TaskGroupId,
5287}
5288
5289impl GuestTask {
5290    fn already_lowered_parameters(&self) -> bool {
5291        // We reset `lower_params` after we lower the parameters
5292        self.lower_params.is_none()
5293    }
5294
5295    fn returned_or_cancelled(&self) -> bool {
5296        // We reset `lift_result` after we return or exit
5297        self.lift_result.is_none()
5298    }
5299
5300    fn ready_to_delete(&self) -> bool {
5301        let threads_completed = self.threads.is_empty();
5302        let has_sync_result = matches!(self.sync_result, SyncResult::Produced(_));
5303        let pending_completion_event = matches!(
5304            self.common.event,
5305            Some(Event::Subtask {
5306                status: Status::Returned | Status::ReturnCancelled
5307            })
5308        );
5309        let ready = threads_completed
5310            && !has_sync_result
5311            && !pending_completion_event
5312            && !matches!(self.host_future_state, HostFutureState::Live);
5313        log::trace!(
5314            "ready to delete? {ready} (threads_completed: {}, has_sync_result: {}, pending_completion_event: {}, host_future_state: {:?})",
5315            threads_completed,
5316            has_sync_result,
5317            pending_completion_event,
5318            self.host_future_state
5319        );
5320        ready
5321    }
5322
5323    fn new(
5324        state: &mut ConcurrentState,
5325        lower_params: RawLower,
5326        lift_result: LiftResult,
5327        caller: Caller,
5328        callback: Option<CallbackFn>,
5329        instance: RuntimeInstance,
5330        async_typed: bool,
5331        async_lifted: bool,
5332    ) -> Result<QualifiedThreadId> {
5333        let host_future_state = match &caller {
5334            Caller::Guest { .. } => HostFutureState::NotApplicable,
5335            Caller::Host {
5336                host_future_present,
5337                ..
5338            } => {
5339                if *host_future_present {
5340                    HostFutureState::Live
5341                } else {
5342                    HostFutureState::NotApplicable
5343                }
5344            }
5345        };
5346
5347        let group = match caller {
5348            Caller::Guest { thread } => {
5349                let group = state.get_mut(thread.task)?.group;
5350                state.increment_group_ref_count(group)?;
5351                group
5352            }
5353            Caller::Host { .. } => state.make_task_group()?,
5354        };
5355
5356        let task = state.push(Self {
5357            common: WaitableCommon::default(),
5358            lower_params: Some(lower_params),
5359            lift_result: Some(lift_result),
5360            result: None,
5361            callback,
5362            caller,
5363            call_context: CallContext::default(),
5364            sync_result: SyncResult::NotProduced,
5365            cancel_request_delivered: false,
5366            starting_sent: false,
5367            instance,
5368            event: None,
5369            exited: false,
5370            threads: HashSet::new(),
5371            host_future_state,
5372            async_typed,
5373            async_lifted,
5374            decremented_interesting_task_count: false,
5375            group,
5376        })?;
5377        let new_thread = GuestThread::new_implicit(state, task)?;
5378        let thread = state.push(new_thread)?;
5379        state.get_mut(task)?.threads.insert(thread);
5380        state.interesting_tasks += 1;
5381        let thread = QualifiedThreadId { task, thread };
5382        log::trace!("new implicit thread {thread:?} for instance {instance:?}");
5383        Ok(thread)
5384    }
5385}
5386
5387impl TableDebug for GuestTask {
5388    fn type_name() -> &'static str {
5389        "GuestTask"
5390    }
5391}
5392
5393/// Represents state common to all kinds of waitables.
5394#[derive(Default)]
5395struct WaitableCommon {
5396    /// The currently pending event for this waitable, if any.
5397    event: Option<Event>,
5398    /// The set to which this waitable belongs, if any.
5399    set: Option<TableId<WaitableSet>>,
5400    /// The handle with which the guest refers to this waitable, if any.
5401    handle: Option<u32>,
5402}
5403
5404/// Represents a Component Model Async `waitable`.
5405#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
5406enum Waitable {
5407    /// A host task
5408    Host(TableId<HostTask>),
5409    /// A guest task
5410    Guest(TableId<GuestTask>),
5411    /// The read or write end of a stream or future
5412    Transmit(TableId<TransmitHandle>),
5413}
5414
5415impl Waitable {
5416    /// Retrieve the `Waitable` corresponding to the specified guest-visible
5417    /// handle.
5418    fn from_instance(
5419        state: Pin<&mut ComponentInstance>,
5420        caller_instance: RuntimeComponentInstanceIndex,
5421        waitable: u32,
5422    ) -> Result<Self> {
5423        use crate::runtime::vm::component::Waitable;
5424
5425        let (waitable, kind) = state.instance_states().0[caller_instance]
5426            .handle_table()
5427            .waitable_rep(waitable)?;
5428
5429        Ok(match kind {
5430            Waitable::Subtask { is_host: true } => Self::Host(TableId::new(waitable)),
5431            Waitable::Subtask { is_host: false } => Self::Guest(TableId::new(waitable)),
5432            Waitable::Stream | Waitable::Future => Self::Transmit(TableId::new(waitable)),
5433        })
5434    }
5435
5436    /// Retrieve the host-visible identifier for this `Waitable`.
5437    fn rep(&self) -> u32 {
5438        match self {
5439            Self::Host(id) => id.rep(),
5440            Self::Guest(id) => id.rep(),
5441            Self::Transmit(id) => id.rep(),
5442        }
5443    }
5444
5445    /// Move this `Waitable` to the specified set (when `set` is `Some(_)`) or
5446    /// remove it from any set it may currently belong to (when `set` is
5447    /// `None`).
5448    fn join(&self, state: &mut ConcurrentState, set: Option<TableId<WaitableSet>>) -> Result<()> {
5449        log::trace!("waitable {self:?} join set {set:?}");
5450
5451        let old = mem::replace(&mut self.common(state)?.set, set);
5452
5453        if let Some(old) = old {
5454            match *self {
5455                Waitable::Host(id) => state.remove_child(id, old),
5456                Waitable::Guest(id) => state.remove_child(id, old),
5457                Waitable::Transmit(id) => state.remove_child(id, old),
5458            }?;
5459
5460            state.get_mut(old)?.ready.remove(self);
5461        }
5462
5463        if let Some(set) = set {
5464            match *self {
5465                Waitable::Host(id) => state.add_child(id, set),
5466                Waitable::Guest(id) => state.add_child(id, set),
5467                Waitable::Transmit(id) => state.add_child(id, set),
5468            }?;
5469
5470            if self.common(state)?.event.is_some() {
5471                self.mark_ready(state)?;
5472            }
5473        }
5474
5475        Ok(())
5476    }
5477
5478    /// Retrieve mutable access to the `WaitableCommon` for this `Waitable`.
5479    fn common<'a>(&self, state: &'a mut ConcurrentState) -> Result<&'a mut WaitableCommon> {
5480        Ok(match self {
5481            Self::Host(id) => &mut state.get_mut(*id)?.common,
5482            Self::Guest(id) => &mut state.get_mut(*id)?.common,
5483            Self::Transmit(id) => &mut state.get_mut(*id)?.common,
5484        })
5485    }
5486
5487    /// Trap if this waitable is currently a member of a waitable set.
5488    ///
5489    /// A synchronous stream/future/subtask operation may end up blocking on
5490    /// this waitable, so it is not allowed to run while the waitable is also
5491    /// being watched by a waitable set.
5492    fn trap_if_in_waitable_set(&self, state: &mut ConcurrentState) -> Result<()> {
5493        if self.common(state)?.set.is_some() {
5494            bail!(Trap::WaitableSyncAndAsync);
5495        }
5496        Ok(())
5497    }
5498
5499    /// Set or clear the pending event for this waitable and either deliver it
5500    /// to the first waiter, if any, or mark it as ready to be delivered to the
5501    /// next waiter that arrives.
5502    fn set_event(&self, state: &mut ConcurrentState, event: Option<Event>) -> Result<()> {
5503        log::trace!("set event for {self:?}: {event:?}");
5504        self.common(state)?.event = event;
5505        self.mark_ready(state)
5506    }
5507
5508    /// Take the pending event from this waitable, leaving `None` in its place.
5509    fn take_event(&self, state: &mut ConcurrentState) -> Result<Option<Event>> {
5510        let common = self.common(state)?;
5511        let event = common.event.take();
5512        if let Some(set) = self.common(state)?.set {
5513            state.get_mut(set)?.ready.remove(self);
5514        }
5515
5516        Ok(event)
5517    }
5518
5519    /// Deliver the current event for this waitable to the first waiter, if any,
5520    /// or else mark it as ready to be delivered to the next waiter that
5521    /// arrives.
5522    fn mark_ready(&self, state: &mut ConcurrentState) -> Result<()> {
5523        if let Some(set) = self.common(state)?.set {
5524            let set_state = state.get_mut(set)?;
5525            set_state.ready.insert(*self);
5526
5527            if let Some((thread, mode)) = set_state.waiting.pop_first() {
5528                let wake_on_cancel = state.get_mut(thread.thread)?.wake_on_cancel.take();
5529                assert!(wake_on_cancel.is_none() || wake_on_cancel == WakeOnCancel::Waiting(set));
5530
5531                let item = match mode {
5532                    WaitMode::Fiber(fiber) => Some(WorkItem::ResumeFiber {
5533                        instance: state.get_mut(thread.task)?.instance,
5534                        thread,
5535                        fiber,
5536                    }),
5537                    WaitMode::Callback(instance) => Some(WorkItem::GuestCall {
5538                        instance: state.get_mut(thread.task)?.instance,
5539                        call: GuestCall {
5540                            thread,
5541                            kind: GuestCallKind::DeliverEvent {
5542                                instance,
5543                                set: Some(set),
5544                            },
5545                        },
5546                    }),
5547                };
5548
5549                if let Some(item) = item {
5550                    state.push_high_priority(item);
5551                }
5552            }
5553        }
5554        Ok(())
5555    }
5556
5557    /// Remove this waitable from the store's rep table.
5558    fn delete_from(&self, store: &mut StoreOpaque) -> Result<()> {
5559        match self {
5560            Self::Host(task) => {
5561                log::trace!("delete host task {task:?}");
5562                let state = store.concurrent_state_mut()?;
5563                let task = state.delete(*task)?;
5564
5565                state.decrement_group_ref_count(task.group)?;
5566            }
5567            Self::Guest(task) => {
5568                log::trace!("delete guest task {task:?}");
5569                let state = store.concurrent_state_mut()?;
5570                let task = state.delete(*task)?;
5571
5572                state.decrement_group_ref_count(task.group)?;
5573
5574                // When a guest task is created it increments the
5575                // `ConcurrentState::interesting_tasks` counter, and that needs
5576                // to be paired with a decrement. There are a few situations in
5577                // which the decrement needs to happen which don't all funnel
5578                // through here, so in lieu of that at least try to catch issues
5579                // where we forgot to do a decrement.
5580                debug_assert!(task.decremented_interesting_task_count);
5581            }
5582            Self::Transmit(task) => {
5583                store.concurrent_state_mut()?.delete(*task)?;
5584            }
5585        }
5586
5587        Ok(())
5588    }
5589}
5590
5591impl fmt::Debug for Waitable {
5592    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
5593        match self {
5594            Self::Host(id) => write!(f, "{id:?}"),
5595            Self::Guest(id) => write!(f, "{id:?}"),
5596            Self::Transmit(id) => write!(f, "{id:?}"),
5597        }
5598    }
5599}
5600
5601/// Represents a Component Model Async `waitable-set`.
5602#[derive(Default)]
5603struct WaitableSet {
5604    /// Which waitables in this set have pending events, if any.
5605    ready: BTreeSet<Waitable>,
5606    /// Which guest threads are currently waiting on this set, if any.
5607    waiting: BTreeMap<QualifiedThreadId, WaitMode>,
5608    /// Whether this set is a synthetic, internal one meant for handling
5609    /// synchronous calls.
5610    is_sync_call_set: bool,
5611}
5612
5613impl TableDebug for WaitableSet {
5614    fn type_name() -> &'static str {
5615        "WaitableSet"
5616    }
5617}
5618
5619/// Type-erased closure to lower the parameters for a guest task.
5620type RawLower =
5621    Box<dyn FnOnce(&mut dyn VMStore, &mut [MaybeUninit<ValRaw>]) -> Result<()> + Send + Sync>;
5622
5623/// Type-erased closure to lift the result for a guest task.
5624type RawLift = Box<
5625    dyn FnOnce(&mut dyn VMStore, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>> + Send + Sync,
5626>;
5627
5628/// Type erased result of a guest task which may be downcast to the expected
5629/// type by a host caller (or simply ignored in the case of a guest caller; see
5630/// `DummyResult`).
5631type LiftedResult = Box<dyn Any + Send + Sync>;
5632
5633/// Used to return a result from a `LiftFn` when the actual result has already
5634/// been lowered to a guest task's stack and linear memory.
5635struct DummyResult;
5636
5637/// Represents the Component Model Async state of a (sub-)component instance.
5638#[derive(Default)]
5639pub struct ConcurrentInstanceState {
5640    /// Whether backpressure is set for this instance (enabled if >0)
5641    backpressure: u16,
5642    /// Whether this instance can be entered
5643    do_not_enter: bool,
5644    /// Whether this instance may suspend (i.e. whether this instance is
5645    /// currently running a sync-typed function).
5646    do_not_suspend: bool,
5647    /// Pending calls for this instance which require `Self::backpressure` to be
5648    /// zero and/or `Self::do_not_enter` to be false before they can proceed.
5649    pending: BTreeMap<QualifiedThreadId, GuestCallKind>,
5650}
5651
5652impl ConcurrentInstanceState {
5653    pub fn pending_is_empty(&self) -> bool {
5654        self.pending.is_empty()
5655    }
5656}
5657
5658#[derive(Debug, Copy, Clone)]
5659pub(crate) enum CurrentThread {
5660    /// The currently running thread is a guest, identified here with its
5661    /// task/thread id combo.
5662    Guest(QualifiedThreadId),
5663    /// The currently running thread is a host task.
5664    Host(TableId<HostTask>),
5665    /// The currently running thread is a host call whose task has not yet been
5666    /// materialized. The contained ID identifies its guest caller.
5667    DeferredHost(QualifiedThreadId),
5668    /// There is no currently running thread because we are in the main event
5669    /// loop or concurrency is disabled.
5670    None,
5671}
5672
5673impl CurrentThread {
5674    fn guest(&self) -> Option<&QualifiedThreadId> {
5675        match self {
5676            Self::Guest(id) => Some(id),
5677            _ => None,
5678        }
5679    }
5680
5681    fn guest_task(&self) -> Option<TableId<GuestTask>> {
5682        match self {
5683            Self::Guest(id) => Some(id.task),
5684            _ => None,
5685        }
5686    }
5687
5688    fn is_none(&self) -> bool {
5689        matches!(self, Self::None)
5690    }
5691}
5692
5693impl From<QualifiedThreadId> for CurrentThread {
5694    fn from(id: QualifiedThreadId) -> Self {
5695        Self::Guest(id)
5696    }
5697}
5698
5699impl From<TableId<HostTask>> for CurrentThread {
5700    fn from(id: TableId<HostTask>) -> Self {
5701        Self::Host(id)
5702    }
5703}
5704
5705enum Priority {
5706    Switch,
5707    High,
5708    Low,
5709}
5710
5711/// Represents the Component Model Async state of a store.
5712pub struct ConcurrentState {
5713    /// The currently running thread, if any.
5714    ///
5715    /// Note that we lazily materialize threads on-demand and this field is not
5716    /// necessarily up-to-date. The `StoreOpaque::current_thread` method should
5717    /// be preferred over directly accessing this field.
5718    unforced_current_thread: CurrentThread,
5719
5720    /// Borrow state for the deferred host call, if any.
5721    ///
5722    /// This is `Some` if and only if [`Self::unforced_current_thread`] is
5723    /// [`CurrentThread::DeferredHost`]. Materializing the host task moves this
5724    /// context into that task.
5725    deferred_host_call_context: Option<CallContext>,
5726
5727    /// The set of pending host and background tasks, if any.
5728    ///
5729    /// See `ComponentInstance::poll_until` for where we temporarily take this
5730    /// out, poll it, then put it back to avoid any mutable aliasing hazards.
5731    futures: AlwaysMut<Option<FuturesUnordered<HostTaskFuture>>>,
5732    /// The table of waitables, waitable sets, etc.
5733    table: AlwaysMut<ResourceTable>,
5734    /// The next item to switch to if any.
5735    ///
5736    /// This takes precedence over items in the `high_priority` queue below and
5737    /// should be used in cases such as subtask calls and thread resume/promote
5738    /// operations where we must switch to a specific thread at the next turn of
5739    /// the event loop regardless of what happens to be present in the
5740    /// `high_priority` queue.
5741    switch_item: Option<WorkItem>,
5742    /// The item to set `switch_item` to when the current thread suspends.
5743    ///
5744    /// This is used when ever an async-lowered import or `subtask.cancel` is
5745    /// called in order to track the thread to switch back to once the current
5746    /// subtask suspends, if any.
5747    next_switch_item: Option<WorkItem>,
5748    /// The "high priority" work queue for this store's event loop.
5749    high_priority: VecDeque<WorkItem>,
5750    /// The "low priority" work queue for this store's event loop.
5751    low_priority: VecDeque<WorkItem>,
5752    /// A place to stash the reason a fiber is suspending so that the code which
5753    /// resumed it will know under what conditions the fiber should be resumed
5754    /// again.
5755    suspend_reason: Option<SuspendReason>,
5756    /// A cached fiber which is waiting for work to do.
5757    ///
5758    /// This helps us avoid creating a new fiber for each `GuestCall` work item.
5759    worker: Option<StoreFiber<'static>>,
5760    /// A place to stash the work item for which we're resuming a worker fiber.
5761    worker_item: Option<WorkerItem>,
5762
5763    /// Reference counts for all component error contexts
5764    ///
5765    /// NOTE: it is possible the global ref count to be *greater* than the sum of
5766    /// (sub)component ref counts as tracked by `error_context_tables`, for
5767    /// example when the host holds one or more references to error contexts.
5768    ///
5769    /// The key of this primary map is often referred to as the "rep" (i.e. host-side
5770    /// component-wide representation) of the index into concurrent state for a given
5771    /// stored `ErrorContext`.
5772    ///
5773    /// Stated another way, `TypeComponentGlobalErrorContextTableIndex` is essentially the same
5774    /// as a `TableId<ErrorContextState>`.
5775    global_error_context_ref_counts:
5776        BTreeMap<TypeComponentGlobalErrorContextTableIndex, GlobalErrorContextRefCount>,
5777
5778    /// The number of "interesting tasks" currently executing in the store.
5779    ///
5780    /// This tracks the concept of a component instance lifetime as defined in
5781    /// https://github.com/WebAssembly/component-model/pull/643. Specifically
5782    /// all tasks currently increment this counter which then gets decremented
5783    /// when they exit. In the future some tasks might not increment this
5784    /// counter, but for now all do.
5785    ///
5786    /// This is used to implement `Accessor::poll_no_interesting_tasks` to
5787    /// inform the embedder when all tasks have completed. This is then
5788    /// used in wasmtime-wasi-http, for example, to know when an instance is
5789    /// idle.
5790    interesting_tasks: usize,
5791
5792    /// Single waker to notify when `interesting_tasks` reaches 0.
5793    ///
5794    /// Used in the implementation of `Accessor::poll_no_interesting_tasks`.
5795    interesting_tasks_empty_waker: Option<Waker>,
5796
5797    /// Single waker to notify when a component instance goes from
5798    /// not-concurrently-callable to concurrently-callable.
5799    ///
5800    /// Used in the implementation of `Accessor::poll_ready_for_concurrent_call`.
5801    ready_for_concurrent_call_waker: Option<Waker>,
5802
5803    /// Whether the `StoreContextMut::poll_until` event loop is running.
5804    event_loop_running: bool,
5805
5806    /// See [TaskGroupHook].
5807    #[cfg(feature = "task-group-hook")]
5808    task_group_hook: Option<Box<dyn TaskGroupHook>>,
5809}
5810
5811impl Default for ConcurrentState {
5812    fn default() -> Self {
5813        Self {
5814            unforced_current_thread: CurrentThread::None,
5815            deferred_host_call_context: None,
5816            table: AlwaysMut::new(ResourceTable::new()),
5817            futures: AlwaysMut::new(Some(FuturesUnordered::new())),
5818            switch_item: None,
5819            next_switch_item: None,
5820            high_priority: VecDeque::new(),
5821            low_priority: VecDeque::new(),
5822            suspend_reason: None,
5823            worker: None,
5824            worker_item: None,
5825            global_error_context_ref_counts: BTreeMap::new(),
5826            interesting_tasks: 0,
5827            interesting_tasks_empty_waker: None,
5828            ready_for_concurrent_call_waker: None,
5829            event_loop_running: false,
5830            #[cfg(feature = "task-group-hook")]
5831            task_group_hook: None,
5832        }
5833    }
5834}
5835
5836impl ConcurrentState {
5837    /// Take ownership of any fibers and futures owned by this object.
5838    ///
5839    /// This should be used when disposing of the `Store` containing this object
5840    /// in order to gracefully resolve any and all fibers using
5841    /// `StoreFiber::dispose`.  This is necessary to avoid possible
5842    /// use-after-free bugs due to fibers which may still have access to the
5843    /// `Store`.
5844    ///
5845    /// Additionally, the futures collected with this function should be dropped
5846    /// within a `tls::set` call, which will ensure than any futures closing
5847    /// over an `&Accessor` will have access to the store when dropped, allowing
5848    /// e.g. `WithAccessor[AndValue]` instances to be disposed of without
5849    /// panicking.
5850    ///
5851    /// Note that this will leave the object in an inconsistent and unusable
5852    /// state, so it should only be used just prior to dropping it.
5853    pub(crate) fn take_fibers_and_futures(
5854        &mut self,
5855        fibers: &mut Vec<StoreFiber<'static>>,
5856        futures: &mut Vec<FuturesUnordered<HostTaskFuture>>,
5857    ) {
5858        let mut items = Vec::new();
5859        for (_, entry) in self.table.get_mut().iter_mut() {
5860            if let Some(set) = entry.downcast_mut::<WaitableSet>() {
5861                for mode in mem::take(&mut set.waiting).into_values() {
5862                    match mode {
5863                        WaitMode::Fiber(fiber) => {
5864                            fibers.push(fiber);
5865                        }
5866                        WaitMode::Callback(_) => {}
5867                    }
5868                }
5869            } else if let Some(thread) = entry.downcast_mut::<GuestThread>() {
5870                if let GuestThreadState::Suspended(fiber) | GuestThreadState::Ready { fiber, .. } =
5871                    mem::replace(&mut thread.state, GuestThreadState::Completed)
5872                {
5873                    fibers.push(fiber);
5874                }
5875            } else if let Some(item) = entry.downcast_mut::<Option<WorkItem>>() {
5876                if let Some(item) = item.take() {
5877                    items.push(item);
5878                }
5879            }
5880        }
5881
5882        if let Some(fiber) = self.worker.take() {
5883            fibers.push(fiber);
5884        }
5885
5886        let mut handle_item = |item| match item {
5887            WorkItem::ResumeFiber { fiber, .. } => {
5888                fibers.push(fiber);
5889            }
5890            WorkItem::PushFuture(future) => {
5891                self.futures
5892                    .get_mut()
5893                    .as_mut()
5894                    .unwrap()
5895                    .push(future.into_inner());
5896            }
5897            WorkItem::ResumeThread { .. }
5898            | WorkItem::GuestCall { .. }
5899            | WorkItem::WorkerFunction(_) => {}
5900        };
5901
5902        for item in items {
5903            handle_item(item);
5904        }
5905        if let Some(item) = self.switch_item.take() {
5906            handle_item(item);
5907        }
5908        if let Some(item) = self.next_switch_item.take() {
5909            handle_item(item);
5910        }
5911        for item in mem::take(&mut self.high_priority) {
5912            handle_item(item);
5913        }
5914        for item in mem::take(&mut self.low_priority) {
5915            handle_item(item);
5916        }
5917
5918        if let Some(them) = self.futures.get_mut().take() {
5919            futures.push(them);
5920        }
5921    }
5922
5923    #[cfg(feature = "gc")]
5924    pub(crate) fn trace_fiber_roots(
5925        &mut self,
5926        modules: &ModuleRegistry,
5927        unwind: &dyn Unwind,
5928        gc_roots_list: &mut GcRootsList,
5929    ) {
5930        let ConcurrentState {
5931            table,
5932            worker,
5933            switch_item,
5934            next_switch_item,
5935            high_priority,
5936            low_priority,
5937
5938            // TODO(cm-gc): This field contains `ValRaw`s, but they are never GC
5939            // references because the component model doesn't support GC yet. We
5940            // will need to trace these somehow when it does.
5941            futures: _,
5942
5943            // These fields do not contain GC references.
5944            worker_item: _,
5945            unforced_current_thread: _,
5946            deferred_host_call_context: _,
5947            suspend_reason: _,
5948            global_error_context_ref_counts: _,
5949            interesting_tasks: _,
5950            interesting_tasks_empty_waker: _,
5951            ready_for_concurrent_call_waker: _,
5952            event_loop_running: _,
5953            #[cfg(feature = "task-group-hook")]
5954                task_group_hook: _,
5955        } = self;
5956
5957        for (_, entry) in table.get_mut().iter_mut() {
5958            if let Some(set) = entry.downcast_mut::<WaitableSet>() {
5959                for mode in set.waiting.values_mut() {
5960                    match mode {
5961                        WaitMode::Fiber(fiber) => {
5962                            fiber.trace_gc_roots(modules, unwind, gc_roots_list);
5963                        }
5964                        WaitMode::Callback(_) => {}
5965                    }
5966                }
5967            } else if let Some(thread) = entry.downcast_mut::<GuestThread>() {
5968                if let GuestThreadState::Suspended(fiber) | GuestThreadState::Ready { fiber, .. } =
5969                    &mut thread.state
5970                {
5971                    fiber.trace_gc_roots(modules, unwind, gc_roots_list);
5972                }
5973            } else if let Some(Some(WorkItem::ResumeFiber { fiber, .. })) =
5974                entry.downcast_mut::<Option<WorkItem>>()
5975            {
5976                fiber.trace_gc_roots(modules, unwind, gc_roots_list);
5977            }
5978        }
5979
5980        if let Some(fiber) = worker {
5981            fiber.trace_gc_roots(modules, unwind, gc_roots_list);
5982        }
5983
5984        let mut handle_item = |item: &mut WorkItem| match item {
5985            WorkItem::ResumeFiber { fiber, .. } => {
5986                fiber.trace_gc_roots(modules, unwind, gc_roots_list);
5987            }
5988            WorkItem::PushFuture(_future) => {
5989                // TODO(cm-gc): once futures can contain GC roots, we will need
5990                // to trace them.
5991            }
5992            WorkItem::ResumeThread { .. }
5993            | WorkItem::GuestCall { .. }
5994            | WorkItem::WorkerFunction(_) => {}
5995        };
5996
5997        if let Some(item) = switch_item {
5998            handle_item(item);
5999        }
6000        if let Some(item) = next_switch_item {
6001            handle_item(item);
6002        }
6003        for item in high_priority {
6004            handle_item(item);
6005        }
6006        for item in low_priority {
6007            handle_item(item);
6008        }
6009    }
6010
6011    fn push<V: Send + Sync + 'static>(
6012        &mut self,
6013        value: V,
6014    ) -> Result<TableId<V>, ResourceTableError> {
6015        self.table.get_mut().push(value).map(TableId::from)
6016    }
6017
6018    fn get_mut<V: 'static>(&mut self, id: TableId<V>) -> Result<&mut V, ResourceTableError> {
6019        self.table.get_mut().get_mut(&Resource::from(id))
6020    }
6021
6022    pub fn add_child<T: 'static, U: 'static>(
6023        &mut self,
6024        child: TableId<T>,
6025        parent: TableId<U>,
6026    ) -> Result<(), ResourceTableError> {
6027        self.table
6028            .get_mut()
6029            .add_child(Resource::from(child), Resource::from(parent))
6030    }
6031
6032    pub fn remove_child<T: 'static, U: 'static>(
6033        &mut self,
6034        child: TableId<T>,
6035        parent: TableId<U>,
6036    ) -> Result<(), ResourceTableError> {
6037        self.table
6038            .get_mut()
6039            .remove_child(Resource::from(child), Resource::from(parent))
6040    }
6041
6042    fn delete<V: 'static>(&mut self, id: TableId<V>) -> Result<V, ResourceTableError> {
6043        self.table.get_mut().delete(Resource::from(id))
6044    }
6045
6046    fn push_future(&mut self, future: HostTaskFuture) {
6047        // Note that we can't directly push to `ConcurrentState::futures` here
6048        // since this may be called from a future that's being polled inside
6049        // `Self::poll_until`, which temporarily removes the `FuturesUnordered`
6050        // so it has exclusive access while polling it.  Therefore, we push a
6051        // work item to the "high priority" queue, which will actually push to
6052        // `ConcurrentState::futures` later.
6053        self.push_high_priority(WorkItem::PushFuture(AlwaysMut::new(future)));
6054    }
6055
6056    fn set_switch_item(&mut self, item: WorkItem) -> Result<()> {
6057        log::trace!("set switch item: {item:?}");
6058
6059        if self.switch_item.is_some() {
6060            bail_bug!("switch item already set");
6061        }
6062
6063        self.switch_item = Some(item);
6064
6065        Ok(())
6066    }
6067
6068    fn take_next_switch_item(&mut self) -> Result<()> {
6069        if let Some(item) = self.next_switch_item.take() {
6070            self.set_switch_item(item)?;
6071        }
6072        Ok(())
6073    }
6074
6075    fn push_high_priority(&mut self, item: WorkItem) {
6076        log::trace!("push high priority: {item:?}");
6077        self.high_priority.push_front(item);
6078    }
6079
6080    fn push_low_priority(&mut self, item: WorkItem) {
6081        log::trace!("push low priority: {item:?}");
6082        self.low_priority.push_front(item);
6083    }
6084
6085    fn push_work_item(&mut self, item: WorkItem, priority: Priority) -> Result<()> {
6086        match priority {
6087            Priority::Switch => self.set_switch_item(item)?,
6088            Priority::High => self.push_high_priority(item),
6089            Priority::Low => self.push_low_priority(item),
6090        }
6091
6092        Ok(())
6093    }
6094
6095    fn promote_instance_local_thread_work_item(
6096        &mut self,
6097        current_instance: RuntimeInstance,
6098    ) -> Result<bool> {
6099        log::trace!("promote thread work items for {current_instance:?}");
6100
6101        self.promote_work_item_matching(|item: &WorkItem| {
6102            let result = match item {
6103                WorkItem::ResumeThread { instance, .. }
6104                | WorkItem::ResumeFiber { instance, .. }
6105                | WorkItem::GuestCall { instance, .. } => *instance == current_instance,
6106                _ => false,
6107            };
6108
6109            log::trace!("candidate {item:?}: {result}");
6110            result
6111        })
6112    }
6113
6114    fn promote_thread_work_item(&mut self, thread: QualifiedThreadId) -> Result<bool> {
6115        self.promote_work_item_matching(|item: &WorkItem| match item {
6116            WorkItem::ResumeThread {
6117                thread: item_thread,
6118                ..
6119            }
6120            | WorkItem::GuestCall {
6121                call:
6122                    GuestCall {
6123                        thread: item_thread,
6124                        ..
6125                    },
6126                ..
6127            } => *item_thread == thread,
6128            _ => false,
6129        })
6130    }
6131
6132    fn promote_work_item_matching<F>(&mut self, mut predicate: F) -> Result<bool>
6133    where
6134        F: FnMut(&WorkItem) -> bool,
6135    {
6136        // Note the use of `.rev()` below to preserve ordering given that items
6137        // are popped from the back of the `VecDeque`s by `poll_until` and
6138        // pushed to the front by `push_{high,low}_priority`.
6139
6140        for item in mem::take(&mut self.high_priority).into_iter().rev() {
6141            if self.switch_item.is_none() && predicate(&item) {
6142                self.set_switch_item(item)?;
6143            } else {
6144                self.push_high_priority(item);
6145            }
6146        }
6147
6148        if self.switch_item.is_none() {
6149            for item in mem::take(&mut self.low_priority).into_iter().rev() {
6150                if self.switch_item.is_none() && predicate(&item) {
6151                    self.set_switch_item(item)?;
6152                } else {
6153                    self.push_low_priority(item);
6154                }
6155            }
6156        }
6157
6158        Ok(self.switch_item.is_some())
6159    }
6160
6161    /// Used by `ResourceTables` to acquire the current `CallContext` for the
6162    /// specified task.
6163    pub fn call_context(&mut self, task: Scope) -> Result<&mut CallContext> {
6164        match task {
6165            Scope::HostId(task) => {
6166                let task: TableId<HostTask> = TableId::new(task);
6167                Ok(&mut self.get_mut(task)?.call_context)
6168            }
6169            Scope::Id(task) => {
6170                let task: TableId<GuestTask> = TableId::new(task);
6171                Ok(&mut self.get_mut(task)?.call_context)
6172            }
6173        }
6174    }
6175
6176    pub(crate) fn deferred_host_call_context(&mut self) -> Option<&mut CallContext> {
6177        self.deferred_host_call_context.as_mut()
6178    }
6179
6180    fn futures_mut(&mut self) -> Result<&mut FuturesUnordered<HostTaskFuture>> {
6181        match self.futures.get_mut().as_mut() {
6182            Some(f) => Ok(f),
6183            None => bail_bug!("futures field of concurrent state is currently taken"),
6184        }
6185    }
6186
6187    pub(crate) fn table(&mut self) -> &mut ResourceTable {
6188        self.table.get_mut()
6189    }
6190
6191    fn debug_assert_deferred_host_invariant(&self) {
6192        debug_assert_eq!(
6193            self.deferred_host_call_context.is_some(),
6194            matches!(self.unforced_current_thread, CurrentThread::DeferredHost(_)),
6195            "a deferred host thread and call context must exist together",
6196        );
6197    }
6198
6199    fn materialize_host_task(&mut self) -> Result<CurrentThread> {
6200        self.debug_assert_deferred_host_invariant();
6201        let caller = match self.unforced_current_thread {
6202            CurrentThread::DeferredHost(caller) => caller,
6203            thread => return Ok(thread),
6204        };
6205
6206        // Push first so allocation failure leaves the deferred state intact.
6207        let task = HostTask::new(self, HostTaskState::CalleeStarted, caller)?;
6208        let task = self.push(task)?;
6209        let call_context = self
6210            .deferred_host_call_context
6211            .take()
6212            .expect("deferred host call context should be present");
6213        self.get_mut(task)
6214            .expect("newly inserted host task should be present")
6215            .call_context = call_context;
6216        self.unforced_current_thread = CurrentThread::Host(task);
6217        self.debug_assert_deferred_host_invariant();
6218        log::trace!("new host task materialized {task:?}");
6219        Ok(CurrentThread::Host(task))
6220    }
6221
6222    fn materialize_current_host_task_id(&mut self) -> Result<Option<TableId<HostTask>>> {
6223        match self.materialize_host_task()? {
6224            CurrentThread::Host(id) => Ok(Some(id)),
6225            CurrentThread::None => Ok(None),
6226            CurrentThread::Guest(_) => {
6227                bail_bug!("tried to materialize a host task id from a guest thread")
6228            }
6229            CurrentThread::DeferredHost(_) => {
6230                bail_bug!(
6231                    "current thread is a deferred host thread which should have been materialized"
6232                )
6233            }
6234        }
6235    }
6236
6237    pub(crate) fn materialize_current_scope(&mut self) -> Result<Scope> {
6238        match self.materialize_host_task()? {
6239            CurrentThread::Host(id) => Ok(Scope::HostId(id.rep())),
6240            _ => bail_bug!("current scope is not a deferred host scope"),
6241        }
6242    }
6243}
6244
6245/// Provide a type hint to compiler about the shape of a parameter lower
6246/// closure.
6247fn for_any_lower<
6248    F: FnOnce(&mut dyn VMStore, &mut [MaybeUninit<ValRaw>]) -> Result<()> + Send + Sync,
6249>(
6250    fun: F,
6251) -> F {
6252    fun
6253}
6254
6255/// Provide a type hint to compiler about the shape of a result lift closure.
6256fn for_any_lift<
6257    F: FnOnce(&mut dyn VMStore, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>> + Send + Sync,
6258>(
6259    fun: F,
6260) -> F {
6261    fun
6262}
6263
6264fn check_ambient_store(id: StoreId) {
6265    let message = "\
6266        `Future`s which depend on asynchronous component tasks, streams, or \
6267        futures to complete may only be polled from the event loop of the \
6268        store to which they belong.  Please use \
6269        `StoreContextMut::{run_concurrent,spawn}` to poll or await them.\
6270    ";
6271    tls::try_get(|store| {
6272        let matched = match store {
6273            tls::TryGet::Some(store) => store.id() == id,
6274            tls::TryGet::Taken | tls::TryGet::None => false,
6275        };
6276
6277        if !matched {
6278            panic!("{message}")
6279        }
6280    });
6281}
6282
6283fn unpack_callback_code(code: u32) -> (u32, u32) {
6284    (code & 0xF, code >> 4)
6285}
6286
6287/// Helper struct for packaging parameters to be passed to
6288/// `ComponentInstance::waitable_check` for calls to `waitable-set.wait` or
6289/// `waitable-set.poll`.
6290struct WaitableCheckParams {
6291    set: TableId<WaitableSet>,
6292    options: OptionsIndex,
6293    payload: u32,
6294}
6295
6296/// Indicates whether `ComponentInstance::waitable_check` is being called for
6297/// `waitable-set.wait` or `waitable-set.poll`.
6298enum WaitableCheck {
6299    Wait,
6300    Poll,
6301}
6302
6303/// Represents a guest task called from the host, prepared using `prepare_call`.
6304pub(crate) struct PreparedCall<R> {
6305    /// The guest export to be called
6306    handle: Func,
6307    /// The guest thread created by `prepare_call`
6308    thread: QualifiedThreadId,
6309    /// The number of lowered core Wasm parameters to pass to the call.
6310    param_count: usize,
6311    /// The `oneshot::Receiver` to which the result of the call will be
6312    /// delivered when it is available.
6313    rx: oneshot::Receiver<LiftedResult>,
6314    /// The instance that this call is prepared for.
6315    runtime_instance: RuntimeInstance,
6316    _phantom: PhantomData<R>,
6317}
6318
6319impl<R> PreparedCall<R> {
6320    /// Get a copy of the `TaskId` for this `PreparedCall`.
6321    pub(crate) fn task_id(&self) -> TaskId {
6322        TaskId {
6323            task: self.thread.task,
6324            runtime_instance: self.runtime_instance,
6325        }
6326    }
6327}
6328
6329/// Represents a task created by `prepare_call`.
6330pub(crate) struct TaskId {
6331    task: TableId<GuestTask>,
6332    runtime_instance: RuntimeInstance,
6333}
6334
6335impl TaskId {
6336    /// The host future for an async task was dropped. If the parameters have not been lowered yet,
6337    /// it is no longer valid to do so, as the lowering closure would see a dangling pointer. In this case,
6338    /// we delete the task eagerly. Otherwise, there may be running threads, or ones that are suspended
6339    /// and can be resumed by other tasks for this component, so we mark the future as dropped
6340    /// and delete the task when all threads are done.
6341    pub(crate) fn host_future_dropped(&self, store: &mut StoreOpaque) -> Result<()> {
6342        let task = store.concurrent_state_mut()?.get_mut(self.task)?;
6343        let delete = if !task.already_lowered_parameters() {
6344            store.cancel_guest_subtask_without_lowered_parameters(
6345                self.runtime_instance,
6346                self.task,
6347            )?;
6348            true
6349        } else {
6350            task.host_future_state = HostFutureState::Dropped;
6351            task.ready_to_delete()
6352        };
6353        if delete {
6354            Waitable::Guest(self.task).delete_from(store)?
6355        }
6356        Ok(())
6357    }
6358}
6359
6360/// Prepare a call to the specified exported Wasm function, providing functions
6361/// for lowering the parameters and lifting the result.
6362///
6363/// To enqueue the returned `PreparedCall` in the `ComponentInstance`'s event
6364/// loop, use `stage_call`.
6365pub(crate) fn prepare_call<T, R>(
6366    mut store: StoreContextMut<T>,
6367    handle: Func,
6368    param_count: usize,
6369    host_future_present: bool,
6370    lower_params: impl FnOnce(StoreContextMut<T>, &mut [MaybeUninit<ValRaw>]) -> Result<()>
6371    + Send
6372    + Sync
6373    + 'static,
6374    lift_result: impl FnOnce(&mut StoreOpaque, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>>
6375    + Send
6376    + Sync
6377    + 'static,
6378) -> Result<PreparedCall<R>> {
6379    if !store.0.may_enter() {
6380        bail!(Trap::CannotEnterComponent);
6381    }
6382
6383    let (options, _flags, ty, raw_options) = handle.abi_info(store.0);
6384
6385    let instance = handle.instance().id().get(store.0);
6386    let options = &instance.component().env_component().options[options];
6387    let ty = &instance.component().types()[ty];
6388    let async_typed = ty.async_;
6389    let async_lifted = raw_options.async_;
6390    let task_return_type = ty.results;
6391    let component_instance = raw_options.instance;
6392    let callback = options.callback.map(|i| instance.runtime_callback(i));
6393    let memory = options
6394        .memory()
6395        .map(|i| instance.runtime_memory(i))
6396        .map(SendSyncPtr::new);
6397    let string_encoding = options.string_encoding;
6398    let token = StoreToken::new(store.as_context_mut());
6399    let caller = store.0.materialize_host_task_id()?;
6400    let state = store.0.concurrent_state_mut()?;
6401
6402    let (tx, rx) = oneshot::channel();
6403
6404    let instance = handle.instance().runtime_instance(component_instance);
6405    let thread = GuestTask::new(
6406        state,
6407        Box::new(for_any_lower(move |store, params| {
6408            lower_params(token.as_context_mut(store), params)
6409        })),
6410        LiftResult {
6411            lift: Box::new(for_any_lift(move |store, result| {
6412                lift_result(store, result)
6413            })),
6414            ty: task_return_type,
6415            memory,
6416            string_encoding,
6417        },
6418        Caller::Host {
6419            tx: Some(tx),
6420            host_future_present,
6421            caller,
6422        },
6423        callback.map(|callback| {
6424            let callback = SendSyncPtr::new(callback);
6425            let instance = handle.instance();
6426            Box::new(move |store: &mut dyn VMStore, event, handle| {
6427                let store = token.as_context_mut(store);
6428                // SAFETY: Per the contract of `prepare_call`, the callback
6429                // will remain valid at least as long is this task exists.
6430                unsafe { instance.call_callback(store, callback, event, handle) }
6431            }) as CallbackFn
6432        }),
6433        instance,
6434        async_typed,
6435        async_lifted,
6436    )?;
6437
6438    Ok(PreparedCall {
6439        handle,
6440        thread,
6441        param_count,
6442        runtime_instance: instance,
6443        rx,
6444        _phantom: PhantomData,
6445    })
6446}
6447
6448pub(crate) struct StagedCall<R> {
6449    store: StoreId,
6450    rx: oneshot::Receiver<LiftedResult>,
6451    _marker: PhantomData<fn() -> R>,
6452    group: TaskGroupId,
6453}
6454
6455impl<R> StagedCall<R> {
6456    /// Queue a call previously prepared using `prepare_call` to be run as part of
6457    /// the associated `ComponentInstance`'s event loop.
6458    ///
6459    /// The returned future will resolve to the result once it is available, but
6460    /// must only be polled via the instance's event loop. See
6461    /// `StoreContextMut::run_concurrent` for details.
6462    pub(crate) fn new<T: 'static>(
6463        mut store: StoreContextMut<T>,
6464        prepared: PreparedCall<R>,
6465    ) -> Result<StagedCall<R>> {
6466        let PreparedCall {
6467            handle,
6468            thread,
6469            param_count,
6470            rx,
6471            ..
6472        } = prepared;
6473
6474        stage_call0(store.as_context_mut(), handle, thread, param_count)?;
6475
6476        Ok(StagedCall {
6477            store: store.0.id(),
6478            rx,
6479            _marker: PhantomData,
6480            group: store.0.concurrent_state_mut()?.get_mut(thread.task)?.group,
6481        })
6482    }
6483}
6484
6485impl<R> Future for StagedCall<R>
6486where
6487    R: 'static,
6488{
6489    type Output = Result<R>;
6490
6491    fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
6492        check_ambient_store(self.store);
6493        Pin::new(&mut self.rx).poll(cx).map(|result| match result {
6494            Ok(r) => match r.downcast() {
6495                Ok(r) => Ok(*r),
6496                Err(_) => bail_bug!("wrong type of value produced"),
6497            },
6498            Err(oneshot::Canceled) => bail_bug!("channel erroneously dropped"),
6499        })
6500    }
6501}
6502
6503/// Queue a call previously prepared using `prepare_call` to be run as part of
6504/// the associated `ComponentInstance`'s event loop.
6505fn stage_call0<T: 'static>(
6506    store: StoreContextMut<T>,
6507    handle: Func,
6508    guest_thread: QualifiedThreadId,
6509    param_count: usize,
6510) -> Result<()> {
6511    let (_options, _, _ty, raw_options) = handle.abi_info(store.0);
6512    let is_concurrent = raw_options.async_;
6513    let callback = raw_options.callback;
6514    let instance = handle.instance();
6515    let callee = handle.lifted_core_func(store.0);
6516    let post_return = raw_options
6517        .post_return
6518        .map(|i| instance.id().get(store.0).runtime_post_return(i));
6519    let callback = callback.map(|i| {
6520        let instance = instance.id().get(store.0);
6521        SendSyncPtr::new(instance.runtime_callback(i))
6522    });
6523
6524    log::trace!("queueing call {guest_thread:?}");
6525
6526    // SAFETY: `callee`, `callback`, and `post_return` are valid pointers
6527    // (with signatures appropriate for this call) and will remain valid as
6528    // long as this instance is valid.
6529    unsafe {
6530        instance.stage_call(
6531            store,
6532            guest_thread,
6533            SendSyncPtr::new(callee),
6534            param_count,
6535            1,
6536            is_concurrent,
6537            callback,
6538            post_return.map(SendSyncPtr::new),
6539            true,
6540        )
6541    }
6542}