1use 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
128const BLOCKED: u32 = 0xffff_ffff;
131
132#[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 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#[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 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
208mod callback_code {
210 pub const EXIT: u32 = 0;
211 pub const YIELD: u32 = 1;
212 pub const WAIT: u32 = 2;
213}
214
215const START_FLAG_ASYNC_CALLEE: u32 = wasmtime_environ::component::START_FLAG_ASYNC_CALLEE as u32;
219
220pub 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 pub fn new(store: StoreContextMut<'a, T>, get_data: fn(&mut T) -> D::Data<'_>) -> Self {
237 Self { store, get_data }
238 }
239
240 pub fn data_mut(&mut self) -> &mut T {
242 self.store.data_mut()
243 }
244
245 pub fn get(&mut self) -> D::Data<'_> {
247 (self.get_data)(self.data_mut())
248 }
249
250 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 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
295pub 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
362pub trait AsAccessor {
379 type Data: 'static;
381
382 type AccessorData: HasData + ?Sized;
385
386 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
408const _: () = {
431 const fn assert<T: Send + Sync>() {}
432 assert::<Accessor<UnsafeCell<u32>>>();
433};
434
435impl<T> Accessor<T> {
436 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 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 pub fn getter(&self) -> fn(&mut T) -> D::Data<'_> {
485 self.get_data
486 }
487
488 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 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 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 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
626pub trait AccessorTask<'fut, T, D = HasSelf<T>>:
648 AsyncFnOnce(&Accessor<T, D>) -> Result<()> + Send + 'static
649where
650 D: HasData + ?Sized,
651{
652 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
669enum CallerInfo {
672 Async {
674 params: Vec<ValRaw>,
675 has_result: bool,
676 },
677 Sync {
679 params: Vec<ValRaw>,
680 result_count: u32,
681 },
682}
683
684enum WaitMode {
686 Fiber(StoreFiber<'static>),
688 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#[derive(Debug)]
704enum SuspendReason {
705 Waiting {
708 set: TableId<WaitableSet>,
709 thread: QualifiedThreadId,
710 },
711 YieldingToSubtask { thread: QualifiedThreadId },
714 NeedWork,
717 Yielding { thread: QualifiedThreadId },
720 ExplicitlySuspending { thread: QualifiedThreadId },
723}
724
725enum GuestCallKind {
727 DeliverEvent {
730 instance: Instance,
732 set: Option<TableId<WaitableSet>>,
737 },
738 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#[derive(Copy, Clone, Debug)]
763pub enum SuspensionTarget {
764 Resume(u32),
765 Promote(u32),
766 None,
767}
768
769#[derive(Copy, Clone, Debug)]
771pub enum ResumeThread {
772 Promote,
773 Resume,
774 ResumeLater,
775}
776
777#[derive(Debug)]
779struct GuestCall {
780 thread: QualifiedThreadId,
781 kind: GuestCallKind,
782}
783
784impl GuestCall {
785 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
816enum WorkerItem {
818 GuestCall(GuestCall),
819 Function(AlwaysMut<Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send>>),
820}
821
822enum WorkItem {
825 PushFuture(AlwaysMut<HostTaskFuture>),
827 ResumeFiber {
829 instance: RuntimeInstance,
830 thread: QualifiedThreadId,
831 fiber: StoreFiber<'static>,
832 },
833 ResumeThread {
835 instance: RuntimeInstance,
836 thread: QualifiedThreadId,
837 },
838 GuestCall {
840 instance: RuntimeInstance,
841 call: GuestCall,
842 },
843 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#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
875pub(crate) enum WaitResult {
876 Cancelled,
877 Completed,
878}
879
880pub(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 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 Poll::Ready(result) => return result,
913
914 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 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 Waitable::Host(task).join(store.concurrent_state_mut()?, None)?;
962 task
963 }
964 };
965
966 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
977fn 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 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 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 #[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 #[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 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 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 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 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 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 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 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 const MAX_TURNS_WITHOUT_YIELD: usize = 128;
1354 let mut turns_without_yield = 0;
1355
1356 loop {
1357 let futures = self.0.concurrent_state_mut()?.futures.get_mut().take();
1361 let mut reset = Reset {
1362 store: self.as_context_mut(),
1363 futures,
1364 };
1365 let mut next = match reset.futures.as_mut() {
1366 Some(f) => pin!(f.next()),
1367 None => bail_bug!("concurrent state missing futures field"),
1368 };
1369
1370 enum PollResult<R> {
1371 Complete(R),
1372 ProcessWork {
1373 ready: Option<WorkItem>,
1374 low_priority: bool,
1375 },
1376 }
1377
1378 let result = future::poll_fn(|cx| {
1379 if let Poll::Ready(value) = tls::set(reset.store.0, || future.as_mut().poll(cx)) {
1382 return Poll::Ready(Ok(PollResult::Complete(value)));
1383 }
1384
1385 let next = match tls::set(reset.store.0, || next.as_mut().poll(cx)) {
1389 Poll::Ready(Some(output)) => {
1390 match output {
1391 Err(e) => return Poll::Ready(Err(e)),
1392 Ok(()) => {}
1393 }
1394 Poll::Ready(true)
1395 }
1396 Poll::Ready(None) => Poll::Ready(false),
1397 Poll::Pending => Poll::Pending,
1398 };
1399
1400 let state = reset.store.0.concurrent_state_mut()?;
1415 let mut ready = state.switch_item.take();
1416 let mut low_priority = false;
1417 if ready.is_none() {
1418 ready = state.high_priority.pop_back();
1419 if ready.is_none() {
1420 ready = state.low_priority.pop_back();
1421 low_priority = true;
1422 }
1423 }
1424 if ready.is_some() {
1425 return Poll::Ready(Ok(PollResult::ProcessWork {
1426 ready,
1427 low_priority,
1428 }));
1429 }
1430
1431 return match next {
1435 Poll::Ready(true) => {
1436 Poll::Ready(Ok(PollResult::ProcessWork {
1442 ready: None,
1443 low_priority: false,
1444 }))
1445 }
1446 Poll::Ready(false) => {
1447 if let Poll::Ready(value) =
1451 tls::set(reset.store.0, || future.as_mut().poll(cx))
1452 {
1453 Poll::Ready(Ok(PollResult::Complete(value)))
1454 } else {
1455 if trap_on_idle {
1461 Poll::Ready(Err(if reset.store.0.any_may_not_suspend()? {
1468 Trap::CannotBlockSyncTask.into()
1469 } else {
1470 Trap::AsyncDeadlock.into()
1472 }))
1473 } else {
1474 Poll::Pending
1478 }
1479 }
1480 }
1481 Poll::Pending => Poll::Pending,
1486 };
1487 })
1488 .await;
1489
1490 drop(reset);
1494
1495 match result? {
1496 PollResult::Complete(value) => break Ok(value),
1499 PollResult::ProcessWork {
1502 ready,
1503 low_priority,
1504 } => {
1505 struct Dispose<'a, T: 'static> {
1506 store: StoreContextMut<'a, T>,
1507 ready: Option<WorkItem>,
1508 }
1509
1510 impl<'a, T> Drop for Dispose<'a, T> {
1511 fn drop(&mut self) {
1512 if let Some(item) = self.ready.take() {
1513 match item {
1514 WorkItem::ResumeFiber { mut fiber, .. } => {
1515 fiber.dispose(self.store.0)
1516 }
1517 WorkItem::PushFuture(future) => {
1518 tls::set(self.store.0, move || drop(future))
1519 }
1520 _ => {}
1521 }
1522 }
1523 }
1524 }
1525
1526 let mut dispose = Dispose {
1527 store: self.as_context_mut(),
1528 ready,
1529 };
1530
1531 if low_priority {
1553 dispose.store.0.yield_now().await;
1554 turns_without_yield = 0;
1555 }
1556
1557 if let Some(item) = dispose.ready.take() {
1558 dispose
1559 .store
1560 .as_context_mut()
1561 .handle_work_item(item)
1562 .await?;
1563 }
1564
1565 turns_without_yield += 1;
1566 if turns_without_yield == MAX_TURNS_WITHOUT_YIELD {
1567 turns_without_yield = 0;
1568 dispose.store.0.yield_now().await;
1569 }
1570 }
1571 }
1572 }
1573 }
1574
1575 async fn handle_work_item(self, item: WorkItem) -> Result<()> {
1577 log::trace!("handle work item {item:?}");
1578 match item {
1579 WorkItem::PushFuture(future) => {
1580 self.0
1581 .concurrent_state_mut()?
1582 .futures_mut()?
1583 .push(future.into_inner());
1584 }
1585 WorkItem::ResumeFiber { fiber, .. } => {
1586 self.0.resume_fiber(fiber).await?;
1587 }
1588 WorkItem::ResumeThread { thread, .. } => {
1589 if let GuestThreadState::Ready { fiber, .. } = mem::replace(
1590 &mut self.0.concurrent_state_mut()?.get_mut(thread.thread)?.state,
1591 GuestThreadState::Running,
1592 ) {
1593 self.0.resume_fiber(fiber).await?;
1594 } else {
1595 bail_bug!("cannot resume non-pending thread {thread:?}");
1596 }
1597 }
1598 WorkItem::GuestCall { call, .. } => {
1599 if call.is_ready(self.0)? {
1600 self.0
1601 .concurrent_state_mut()?
1602 .get_mut(call.thread.thread)?
1603 .wake_on_cancel = WakeOnCancel::None;
1604 self.run_on_worker(WorkerItem::GuestCall(call)).await?;
1605 } else {
1606 let state = self.0.concurrent_state_mut()?;
1607 let task = state.get_mut(call.thread.task)?;
1608 if !task.starting_sent {
1609 task.starting_sent = true;
1610 if let GuestCallKind::StartImplicit(_) = &call.kind {
1611 Waitable::Guest(call.thread.task).set_event(
1612 state,
1613 Some(Event::Subtask {
1614 status: Status::Starting,
1615 }),
1616 )?;
1617 }
1618 }
1619
1620 let instance = state.get_mut(call.thread.task)?.instance;
1621 self.0
1622 .instance_state(instance)
1623 .concurrent_state()
1624 .pending
1625 .insert(call.thread, call.kind);
1626
1627 self.0.concurrent_state_mut()?.take_next_switch_item()?;
1631 }
1632 }
1633 WorkItem::WorkerFunction(fun) => {
1634 self.run_on_worker(WorkerItem::Function(fun)).await?;
1635 }
1636 }
1637
1638 Ok(())
1639 }
1640
1641 async fn run_on_worker(self, item: WorkerItem) -> Result<()> {
1643 let worker = if let Some(fiber) = self.0.concurrent_state_mut()?.worker.take() {
1644 fiber
1645 } else {
1646 unsafe {
1665 fiber::make_fiber_unchecked(self.0, move |store| {
1666 loop {
1667 let Some(item) = store.concurrent_state_mut()?.worker_item.take() else {
1668 bail_bug!("worker_item not present when resuming fiber")
1669 };
1670 match item {
1671 WorkerItem::GuestCall(call) => handle_guest_call(store, call)?,
1672 WorkerItem::Function(fun) => fun.into_inner()(store)?,
1673 }
1674
1675 store.suspend(SuspendReason::NeedWork)?;
1676 }
1677 })?
1678 }
1679 };
1680
1681 let worker_item = &mut self.0.concurrent_state_mut()?.worker_item;
1682 assert!(worker_item.is_none());
1683 *worker_item = Some(item);
1684
1685 self.0.resume_fiber(worker).await
1686 }
1687
1688 pub(crate) fn wrap_call<F, R>(self, closure: F) -> impl Future<Output = Result<R>> + 'static
1693 where
1694 T: 'static,
1695 F: FnOnce(&Accessor<T>) -> Pin<Box<dyn Future<Output = Result<R>> + Send + '_>>
1696 + Send
1697 + Sync
1698 + 'static,
1699 R: Send + Sync + 'static,
1700 {
1701 let token = StoreToken::new(self);
1702 async move {
1703 let mut accessor = Accessor::new(token);
1704 closure(&mut accessor).await
1705 }
1706 }
1707
1708 pub(crate) async fn start_instance(
1709 &mut self,
1710 instance: ModuleInstance,
1711 ) -> Result<ModuleInstance> {
1712 let (tx, rx) = oneshot::channel();
1713 let token = StoreToken::new(self.as_context_mut());
1714 self.0.queue_task(move |store| {
1715 _ = tx.send(
1716 instance
1717 .start_raw(&mut token.as_context_mut(store))
1718 .map(|()| instance),
1719 );
1720 Ok(())
1721 })?;
1722 self.as_context_mut()
1723 .run_concurrent_trap_on_idle(async |_| {
1724 rx.await
1725 .map_err(|_| format_err!("oneshot channel canceled"))
1726 })
1727 .await??
1728 }
1729}
1730
1731pub type EnteredHostTask = Option<QualifiedThreadId>;
1738
1739impl StoreOpaque {
1740 #[inline]
1744 pub(crate) fn current_thread(&mut self) -> Result<CurrentThread> {
1745 if !self.concurrency_support() {
1747 return Ok(CurrentThread::None);
1748 }
1749
1750 if !self
1753 .vm_store_context_mut()
1754 .current_thread_mut()
1755 .is_deferred()
1756 {
1757 return Ok(self
1758 .concurrent_state_mut_already_forced_current_thread()
1759 .unforced_current_thread);
1760 }
1761
1762 self.force_deferred_current_thread()
1763 }
1764
1765 #[cold]
1768 fn force_deferred_current_thread(&mut self) -> Result<CurrentThread> {
1769 let state = self.concurrent_state_mut_without_forcing_current_thread();
1778 let id = match state.unforced_current_thread.guest_task() {
1779 Some(task) => state.get_mut(task)?.instance.instance,
1780 None => bail_bug!("deferred component-model thread with non-guest base"),
1781 };
1782
1783 let mut frames = Vec::new();
1786 let mut cur = *self.vm_store_context_mut().current_thread_mut();
1787 while let Some(ptr) = cur.as_deferred() {
1788 let deferred = unsafe { ptr.as_non_null().as_ref() };
1793 frames.push((
1794 deferred.callee_async != 0,
1795 deferred.callee_instance,
1796 deferred.saved_context,
1797 ));
1798 cur = deferred.parent;
1799 }
1800
1801 *self.vm_store_context_mut().current_thread_mut() = VMLazyThread::forced();
1805
1806 let current_context = *self.vm_store_context_mut().component_context_mut();
1809
1810 for (callee_async, callee_instance, saved_context) in frames.into_iter().rev() {
1814 *self.vm_store_context_mut().component_context_mut() = saved_context;
1818 let callee = RuntimeInstance {
1819 instance: id,
1820 index: RuntimeComponentInstanceIndex::from_u32(callee_instance),
1821 };
1822 self.enter_guest_sync_call(callee_async, callee)?;
1823 }
1824
1825 *self.vm_store_context_mut().component_context_mut() = current_context;
1827
1828 Ok(self
1829 .concurrent_state_mut_without_forcing_current_thread()
1830 .unforced_current_thread)
1831 }
1832
1833 fn current_guest_thread(&mut self) -> Result<QualifiedThreadId> {
1834 match self.current_thread()?.guest() {
1835 Some(id) => Ok(*id),
1836 None => bail_bug!("current thread is not a guest thread"),
1837 }
1838 }
1839
1840 pub(crate) fn current_materialized_host_task(&mut self) -> Result<Option<TableId<HostTask>>> {
1844 match self.current_thread()? {
1845 CurrentThread::Host(id) => Ok(Some(id)),
1846 CurrentThread::DeferredHost(_) | CurrentThread::None => Ok(None),
1847 _ => bail_bug!("current thread is not a host thread"),
1848 }
1849 }
1850
1851 fn materialize_host_task_id(&mut self) -> Result<Option<TableId<HostTask>>> {
1854 Ok(self
1855 .concurrent_state_mut()?
1856 .materialize_current_host_task_id()?)
1857 }
1858
1859 fn enter_sync_call(&mut self, callee: RuntimeInstance) -> Result<()> {
1860 log::trace!("enter sync-typed call {callee:?}");
1861 let state = self.instance_state(callee).concurrent_state();
1862 let old_do_not_suspend = state.do_not_suspend;
1863 state.do_not_suspend = true;
1864
1865 let thread = self.current_guest_thread()?;
1866 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
1867 if thread.old_do_not_suspend.is_some() {
1868 bail_bug!("current thread already has `old_do_not_suspend` value");
1869 }
1870
1871 thread.old_do_not_suspend = Some(old_do_not_suspend);
1872
1873 Ok(())
1874 }
1875
1876 fn exit_sync_call(&mut self, callee: RuntimeInstance) -> Result<()> {
1877 log::trace!("exit sync-typed call {callee:?}");
1878 let thread = self.current_guest_thread()?;
1879 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
1880 let Some(old_do_not_suspend) = thread.old_do_not_suspend.take() else {
1881 bail_bug!("current thread missing `old_do_not_suspend` value");
1882 };
1883 let state = self.instance_state(callee).concurrent_state();
1884 state.do_not_suspend = old_do_not_suspend;
1885 Ok(())
1886 }
1887
1888 pub(crate) fn enter_guest_sync_call(
1900 &mut self,
1901 callee_async_typed: bool,
1902 callee: RuntimeInstance,
1903 ) -> Result<()> {
1904 log::trace!("enter sync-lifted call {callee:?}");
1905 if !self.concurrency_support() {
1906 return self.enter_call_not_concurrent();
1907 }
1908
1909 let thread = self.current_thread()?;
1910 let caller = if let Some(thread) = thread.guest() {
1911 Caller::Guest { thread: *thread }
1912 } else {
1913 Caller::Host {
1914 tx: None,
1915 host_future_present: false,
1916 caller: self.materialize_host_task_id()?,
1917 }
1918 };
1919
1920 let state = self.concurrent_state_mut()?;
1921 let guest_thread = GuestTask::new(
1922 state,
1923 Box::new(move |_, _| bail_bug!("cannot lower params in sync call")),
1924 LiftResult {
1925 lift: Box::new(move |_, _| bail_bug!("cannot lift result in sync call")),
1926 ty: TypeTupleIndex::reserved_value(),
1927 memory: None,
1928 string_encoding: StringEncoding::Utf8,
1929 },
1930 caller,
1931 None,
1932 callee,
1933 callee_async_typed,
1934 true,
1935 )?;
1936
1937 Instance::from_wasmtime(self, callee.instance).add_guest_thread_to_instance_table(
1938 guest_thread.thread,
1939 self,
1940 callee.index,
1941 )?;
1942 self.set_thread(guest_thread)?;
1943
1944 if !callee_async_typed {
1945 self.enter_sync_call(callee)?;
1946 }
1947
1948 Ok(())
1949 }
1950
1951 pub(crate) fn exit_guest_sync_call(&mut self) -> Result<()> {
1959 if !self.concurrency_support() {
1960 return Ok(self.exit_call_not_concurrent());
1961 }
1962
1963 let thread = match self.current_thread()?.guest() {
1964 Some(t) => *t,
1965 None => bail_bug!("expected task when exiting"),
1966 };
1967 let task = self.concurrent_state_mut()?.get_mut(thread.task)?;
1968 let instance = task.instance;
1969
1970 let caller = match &task.caller {
1971 &Caller::Guest { thread } => thread.into(),
1972 &Caller::Host { caller, .. } => caller
1973 .map(CurrentThread::Host)
1974 .unwrap_or(CurrentThread::None),
1975 };
1976 task.lift_result = None;
1977 task.exited = true;
1978 let async_typed = task.async_typed;
1979
1980 if !async_typed {
1981 self.exit_sync_call(instance)?;
1982 }
1983
1984 self.set_thread(caller)?;
1985
1986 log::trace!("exit sync-lifted call {instance:?}");
1987
1988 if async_typed {
1989 self.switch_or_trap_if_may_not_suspend(instance)?;
1994 }
1995
1996 self.cleanup_thread(thread, instance, CleanupTask::Yes)?;
1997
1998 Ok(())
1999 }
2000
2001 pub(crate) fn host_task_create(&mut self) -> Result<EnteredHostTask> {
2008 if !self.concurrency_support() {
2009 self.enter_call_not_concurrent()?;
2010 return Ok(None);
2011 }
2012 let caller = self.current_guest_thread()?;
2013 log::trace!("new deferred host task with caller {caller:?}");
2014
2015 self.set_thread(CurrentThread::DeferredHost(caller))?;
2016 let state = self.concurrent_state_mut()?;
2017 debug_assert!(state.deferred_host_call_context.is_none());
2018 state.deferred_host_call_context = Some(CallContext::default());
2019 state.debug_assert_deferred_host_invariant();
2020 Ok(Some(caller))
2021 }
2022
2023 pub(crate) fn host_task_delete(
2030 &mut self,
2031 original_task: EnteredHostTask,
2032 materialized_task: Option<TableId<HostTask>>,
2033 ) -> Result<()> {
2034 match original_task {
2035 Some(caller) => {
2036 self.set_thread(caller)?;
2037 if materialized_task.is_none() {
2038 let state = self.concurrent_state_mut()?;
2039 let context = state
2040 .deferred_host_call_context
2041 .take()
2042 .expect("deferred host call context should be present");
2043 debug_assert!(context.is_empty());
2044 state.debug_assert_deferred_host_invariant();
2045 }
2046 log::trace!(
2047 "delete host task with caller {original_task:?} and materialized as {materialized_task:?}"
2048 );
2049 if let Some(task) = materialized_task {
2050 Waitable::Host(task).delete_from(self)?;
2051 }
2052 }
2053 None => {
2054 debug_assert!(materialized_task.is_none());
2055 self.exit_call_not_concurrent();
2056 }
2057 }
2058 Ok(())
2059 }
2060
2061 fn instance_state(&mut self, instance: RuntimeInstance) -> &mut InstanceState {
2064 self.component_instance_mut(instance.instance)
2065 .instance_state(instance.index)
2066 }
2067
2068 pub(crate) fn set_thread(&mut self, thread: impl Into<CurrentThread>) -> Result<CurrentThread> {
2074 let thread = thread.into();
2075 let state = self.concurrent_state_mut()?;
2076 state.debug_assert_deferred_host_invariant();
2077 let old_thread = mem::replace(&mut state.unforced_current_thread, thread);
2078
2079 state.handle_thread_switch(old_thread, thread)?;
2080
2081 if let Some(old_thread) = old_thread.guest() {
2089 let old_context = *self.vm_store_context_mut().component_context_mut();
2090 self.concurrent_state_mut()?
2091 .get_mut(old_thread.thread)?
2092 .context = old_context;
2093 }
2094 if cfg!(debug_assertions) {
2095 *self.vm_store_context_mut().component_context_mut() =
2096 [u32::MAX; NUM_COMPONENT_CONTEXT_SLOTS];
2097 }
2098 if let Some(thread) = thread.guest() {
2099 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
2100 let context = thread.context;
2101 if cfg!(debug_assertions) {
2102 thread.context = [u32::MAX; NUM_COMPONENT_CONTEXT_SLOTS];
2103 }
2104 *self.vm_store_context_mut().component_context_mut() = context;
2105 }
2106
2107 *self.vm_store_context_mut().current_thread_mut() = if thread.is_none() {
2109 VMLazyThread::none()
2110 } else {
2111 VMLazyThread::forced()
2112 };
2113
2114 Ok(old_thread)
2115 }
2116
2117 fn switch_or_trap_if_may_not_suspend(&mut self, instance: RuntimeInstance) -> Result<()> {
2119 if self.switch_if_may_not_suspend(instance)? {
2120 Ok(())
2121 } else {
2122 Err(Trap::CannotBlockSyncTask.into())
2123 }
2124 }
2125
2126 fn switch_if_may_not_suspend(&mut self, instance: RuntimeInstance) -> Result<bool> {
2130 self.concurrent_state_mut()?;
2134
2135 Ok(!self.concurrency_support()
2136 || !self
2137 .instance_state(instance)
2138 .concurrent_state()
2139 .do_not_suspend
2140 || self
2141 .concurrent_state_mut()?
2142 .promote_instance_local_thread_work_item(instance)?)
2143 }
2144
2145 fn enter_instance(&mut self, instance: RuntimeInstance) {
2149 log::trace!("enter {instance:?}");
2150 self.instance_state(instance)
2151 .concurrent_state()
2152 .do_not_enter = true;
2153 }
2154
2155 fn exit_instance(&mut self, instance: RuntimeInstance) -> Result<()> {
2159 log::trace!("exit {instance:?}");
2160 self.instance_state(instance)
2161 .concurrent_state()
2162 .do_not_enter = false;
2163 self.partition_pending(instance)
2164 }
2165
2166 fn partition_pending(&mut self, instance: RuntimeInstance) -> Result<()> {
2174 for (thread, kind) in
2175 mem::take(&mut self.instance_state(instance).concurrent_state().pending).into_iter()
2176 {
2177 let call = GuestCall { thread, kind };
2178 if call.is_ready(self)? {
2179 self.concurrent_state_mut()?
2180 .push_high_priority(WorkItem::GuestCall { instance, call });
2181 } else {
2182 self.instance_state(instance)
2183 .concurrent_state()
2184 .pending
2185 .insert(call.thread, call.kind);
2186 }
2187 }
2188
2189 if let Some(waker) = self
2190 .concurrent_state_mut()?
2191 .ready_for_concurrent_call_waker
2192 .take()
2193 {
2194 waker.wake();
2195 }
2196
2197 Ok(())
2198 }
2199
2200 pub(crate) fn backpressure_modify(
2202 &mut self,
2203 caller_instance: RuntimeInstance,
2204 modify: impl FnOnce(u16) -> Option<u16>,
2205 ) -> Result<()> {
2206 let state = self.instance_state(caller_instance).concurrent_state();
2207 let old = state.backpressure;
2208 let new = modify(old).ok_or_else(|| Trap::BackpressureOverflow)?;
2209 state.backpressure = new;
2210
2211 if old > 0 && new == 0 {
2212 self.partition_pending(caller_instance)?;
2215 }
2216
2217 Ok(())
2218 }
2219
2220 async fn resume_fiber(&mut self, fiber: StoreFiber<'static>) -> Result<()> {
2223 let old_thread = self.current_thread()?;
2224 log::trace!("resume_fiber: save current thread {old_thread:?}");
2225
2226 let fiber = fiber::resolve_or_release(self, fiber).await?;
2227
2228 self.set_thread(old_thread)?;
2229
2230 let state = self.concurrent_state_mut()?;
2231
2232 if let Some(ot) = old_thread.guest() {
2233 state.get_mut(ot.thread)?.state = GuestThreadState::Running;
2234 }
2235 log::trace!("resume_fiber: restore current thread {old_thread:?}");
2236
2237 if let Some(mut fiber) = fiber {
2238 log::trace!("resume_fiber: suspend reason {:?}", &state.suspend_reason);
2239 let reason = match state.suspend_reason.take() {
2241 Some(r) => r,
2242 None => bail_bug!("suspend reason missing when resuming fiber"),
2243 };
2244 match reason {
2245 SuspendReason::NeedWork => {
2246 if state.worker.is_none() {
2247 state.worker = Some(fiber);
2248 } else {
2249 fiber.dispose(self);
2250 }
2251 }
2252 SuspendReason::Yielding { thread } => {
2253 state.get_mut(thread.thread)?.state = GuestThreadState::Ready { fiber };
2254 let instance = state.get_mut(thread.task)?.instance;
2255 state.push_low_priority(WorkItem::ResumeThread { instance, thread });
2256 }
2257 SuspendReason::ExplicitlySuspending { thread } => {
2258 state.get_mut(thread.thread)?.state = GuestThreadState::Suspended(fiber);
2259 }
2260 SuspendReason::Waiting { set, thread } => {
2261 let old = state
2262 .get_mut(set)?
2263 .waiting
2264 .insert(thread, WaitMode::Fiber(fiber));
2265 assert!(old.is_none());
2266 }
2267 SuspendReason::YieldingToSubtask { thread } => {
2268 let item = WorkItem::ResumeFiber {
2277 instance: state.get_mut(thread.task)?.instance,
2278 thread,
2279 fiber,
2280 };
2281
2282 if state.next_switch_item.replace(item).is_some() {
2283 bail_bug!(
2286 "`ConcurrentState::next_switch_item` was already `Some(_)` when \
2287 a thread wanted to wait on a subtask"
2288 );
2289 }
2290 }
2291 };
2292 } else {
2293 log::trace!("resume_fiber: fiber has exited");
2294 }
2295
2296 Ok(())
2297 }
2298
2299 fn suspend(&mut self, reason: SuspendReason) -> Result<()> {
2305 log::trace!("suspend fiber: {reason:?}");
2306
2307 let state = self.concurrent_state_mut()?;
2308
2309 let (save_and_restore_thread, save_and_restore_next_switch_item) = match &reason {
2316 SuspendReason::Yielding { .. }
2317 | SuspendReason::Waiting { .. }
2318 | SuspendReason::ExplicitlySuspending { .. } => {
2319 if state.switch_item.is_none() {
2322 state.take_next_switch_item()?;
2323 }
2324
2325 (true, false)
2326 }
2327 SuspendReason::YieldingToSubtask { .. } => (true, true),
2328 SuspendReason::NeedWork => (false, false),
2329 };
2330
2331 let old_next_switch_item = if save_and_restore_next_switch_item {
2332 let item = state.next_switch_item.take();
2333 Some(state.push(item)?)
2337 } else {
2338 None
2339 };
2340
2341 let old_guest_thread = if save_and_restore_thread {
2342 self.current_thread()?
2343 } else {
2344 CurrentThread::None
2345 };
2346
2347 let suspend_reason = &mut self.concurrent_state_mut()?.suspend_reason;
2348 assert!(suspend_reason.is_none());
2349 *suspend_reason = Some(reason);
2350
2351 if !self.fiber_async_state_mut().can_block() {
2354 return Err(format_err!("future dropped"));
2355 }
2356
2357 self.with_blocking(|_, cx| cx.suspend(StoreFiberYield::ReleaseStore))?;
2358
2359 if save_and_restore_thread {
2360 self.set_thread(old_guest_thread)?;
2361 }
2362
2363 if let Some(item) = old_next_switch_item {
2364 let state = self.concurrent_state_mut()?;
2365 state.next_switch_item = state.delete(item)?;
2366 }
2367
2368 Ok(())
2369 }
2370
2371 fn wait_for_event(
2372 &mut self,
2373 caller_instance: RuntimeInstance,
2374 waitable: Waitable,
2375 ) -> Result<()> {
2376 let caller = self.current_guest_thread()?;
2377 let state = self.concurrent_state_mut()?;
2378
2379 waitable.trap_if_in_waitable_set(state)?;
2380
2381 let set = state.get_mut(caller.thread)?.sync_call_set;
2382 waitable.join(state, Some(set))?;
2383
2384 self.switch_or_trap_if_may_not_suspend(caller_instance)?;
2385
2386 self.suspend(SuspendReason::Waiting {
2387 set,
2388 thread: caller,
2389 })?;
2390 let state = self.concurrent_state_mut()?;
2391
2392 waitable.join(state, None)
2393 }
2394
2395 fn cleanup_thread(
2417 &mut self,
2418 guest_thread: QualifiedThreadId,
2419 runtime_instance: RuntimeInstance,
2420 cleanup_task: CleanupTask,
2421 ) -> Result<()> {
2422 let state = self.concurrent_state_mut()?;
2423 state.take_next_switch_item()?;
2426 let thread_data = state.get_mut(guest_thread.thread)?;
2427 let sync_call_set = thread_data.sync_call_set;
2428 if let Some(guest_id) = thread_data.instance_rep {
2429 self.instance_state(runtime_instance)
2430 .thread_handle_table()
2431 .guest_thread_remove(guest_id)?;
2432 }
2433 let state = self.concurrent_state_mut()?;
2434
2435 for waitable in mem::take(&mut state.get_mut(sync_call_set)?.ready) {
2437 if let Some(Event::Subtask {
2438 status: Status::Returned | Status::ReturnCancelled,
2439 }) = waitable.common(self.concurrent_state_mut()?)?.event
2440 {
2441 waitable.delete_from(self)?;
2442 }
2443 }
2444
2445 let state = self.concurrent_state_mut()?;
2446 state.delete(guest_thread.thread)?;
2447 state.delete(sync_call_set)?;
2448 let task = state.get_mut(guest_thread.task)?;
2449 task.threads.remove(&guest_thread.thread);
2450
2451 if task.threads.is_empty() && !task.returned_or_cancelled() {
2452 bail!(Trap::NoAsyncResult);
2453 }
2454 let ready_to_delete = task.ready_to_delete();
2455
2456 if !task.decremented_interesting_task_count && task.exited && task.returned_or_cancelled() {
2457 task.decremented_interesting_task_count = true;
2458
2459 debug_assert!(state.interesting_tasks > 0);
2460 state.interesting_tasks -= 1;
2461 if state.interesting_tasks == 0
2462 && let Some(waker) = state.interesting_tasks_empty_waker.take()
2463 {
2464 waker.wake();
2465 }
2466 }
2467
2468 match cleanup_task {
2469 CleanupTask::Yes => {
2470 if ready_to_delete {
2471 Waitable::Guest(guest_thread.task).delete_from(self)?;
2472 }
2473 }
2474 CleanupTask::No => {}
2475 }
2476
2477 Ok(())
2478 }
2479
2480 fn cancel_guest_subtask_without_lowered_parameters(
2493 &mut self,
2494 caller_instance: RuntimeInstance,
2495 guest_task: TableId<GuestTask>,
2496 ) -> Result<()> {
2497 let concurrent_state = self.concurrent_state_mut()?;
2498 let task = concurrent_state.get_mut(guest_task)?;
2499 assert!(!task.already_lowered_parameters());
2500 task.lower_params = None;
2504 task.lift_result = None;
2505 task.exited = true;
2506 let instance = task.instance;
2507
2508 assert_eq!(1, task.threads.len());
2511 let thread = *task.threads.iter().next().unwrap();
2512 self.cleanup_thread(
2513 QualifiedThreadId {
2514 task: guest_task,
2515 thread,
2516 },
2517 caller_instance,
2518 CleanupTask::No,
2519 )?;
2520
2521 let pending = &mut self.instance_state(instance).concurrent_state().pending;
2523 let pending_count = pending.len();
2524 pending.retain(|thread, _| thread.task != guest_task);
2525 if pending.len() == pending_count {
2527 bail!(Trap::SubtaskCancelAfterTerminal);
2528 }
2529 Ok(())
2530 }
2531
2532 pub(crate) fn current_scope(&mut self) -> Result<Option<CurrentScope>> {
2535 if !self.concurrency_support() {
2536 return Ok(self
2537 .current_scope_id_not_concurrent()?
2538 .map(|id| CurrentScope::Id(Scope::Id(id))));
2539 }
2540
2541 Ok(match self.current_thread()? {
2542 CurrentThread::Guest(id) => Some(CurrentScope::Id(Scope::Id(id.task.rep()))),
2543 CurrentThread::Host(id) => Some(CurrentScope::Id(Scope::HostId(id.rep()))),
2544 CurrentThread::DeferredHost(_) => Some(CurrentScope::DeferredHost),
2545 CurrentThread::None => return Ok(None),
2546 })
2547 }
2548
2549 pub(crate) fn queue_task(
2550 &mut self,
2551 task: impl FnOnce(&mut dyn VMStore) -> Result<()> + Send + 'static,
2552 ) -> Result<()> {
2553 self.concurrent_state_mut()?
2554 .push_high_priority(WorkItem::WorkerFunction(AlwaysMut::new(Box::new(task))));
2555 Ok(())
2556 }
2557
2558 fn any_may_not_suspend(&mut self) -> Result<bool> {
2567 Ok(self
2575 .concurrent_state_mut()?
2576 .table
2577 .get_mut()
2578 .iter_mut()
2579 .filter_map(|(_, entry)| {
2580 if let Some(task) = entry.downcast_ref::<GuestTask>() {
2581 Some(task.instance)
2582 } else {
2583 None
2584 }
2585 })
2586 .collect::<Vec<_>>()
2587 .into_iter()
2588 .any(|instance| {
2589 self.instance_state(instance)
2590 .concurrent_state()
2591 .do_not_suspend
2592 }))
2593 }
2594}
2595
2596enum CleanupTask {
2597 Yes,
2598 No,
2599}
2600
2601impl Instance {
2602 fn get_event(
2605 self,
2606 store: &mut StoreOpaque,
2607 guest_task: TableId<GuestTask>,
2608 set: Option<TableId<WaitableSet>>,
2609 cancellable: bool,
2610 ) -> Result<Option<(Event, Option<(Waitable, u32)>)>> {
2611 let state = store.concurrent_state_mut()?;
2612
2613 let task = state.get_mut(guest_task)?;
2614 let event = &mut task.event;
2615 if let Some(ev) = event
2616 && (cancellable || !matches!(ev, Event::Cancelled))
2617 {
2618 log::trace!("deliver event {ev:?} to {guest_task:?}");
2619
2620 if matches!(ev, Event::Cancelled) {
2621 task.cancel_request_delivered = true;
2622 }
2623
2624 let ev = *ev;
2625 *event = None;
2626 return Ok(Some((ev, None)));
2627 }
2628
2629 let set = match set {
2630 Some(set) => set,
2631 None => return Ok(None),
2632 };
2633 let waitable = match state.get_mut(set)?.ready.pop_first() {
2634 Some(v) => v,
2635 None => return Ok(None),
2636 };
2637
2638 let common = waitable.common(state)?;
2639 let handle = match common.handle {
2640 Some(h) => h,
2641 None => bail_bug!("handle not set when delivering event"),
2642 };
2643 let event = match common.event.take() {
2644 Some(e) => e,
2645 None => bail_bug!("event not set when delivering event"),
2646 };
2647
2648 log::trace!(
2649 "deliver event {event:?} to {guest_task:?} for {waitable:?} (handle {handle}); set {set:?}"
2650 );
2651
2652 waitable.on_delivery(store, self, event)?;
2653
2654 Ok(Some((event, Some((waitable, handle)))))
2655 }
2656
2657 fn handle_callback_code(
2663 self,
2664 store: &mut StoreOpaque,
2665 guest_thread: QualifiedThreadId,
2666 runtime_instance: RuntimeComponentInstanceIndex,
2667 code: u32,
2668 ) -> Result<()> {
2669 let (code, set) = unpack_callback_code(code);
2670
2671 log::trace!("received callback code from {guest_thread:?}: {code} (set: {set})");
2672
2673 let state = store.concurrent_state_mut()?;
2674
2675 state.take_next_switch_item()?;
2676
2677 let get_set = |store: &mut StoreOpaque, handle| -> Result<_> {
2678 let set = store
2679 .instance_state(self.runtime_instance(runtime_instance))
2680 .handle_table()
2681 .waitable_set_rep(handle)?;
2682
2683 Ok(TableId::<WaitableSet>::new(set))
2684 };
2685
2686 match code {
2687 callback_code::EXIT => {
2688 log::trace!("implicit thread {guest_thread:?} completed");
2689 let task = store.concurrent_state_mut()?.get_mut(guest_thread.task)?;
2690 task.exited = true;
2691 task.callback = None;
2692
2693 let runtime_instance = self.runtime_instance(runtime_instance);
2694
2695 store.switch_or_trap_if_may_not_suspend(runtime_instance)?;
2700
2701 store.cleanup_thread(guest_thread, runtime_instance, CleanupTask::Yes)?;
2702 }
2703 callback_code::YIELD => {
2704 let old = state
2707 .get_mut(guest_thread.thread)?
2708 .wake_on_cancel
2709 .replace(WakeOnCancel::Yielding);
2710 if !old.is_none() {
2711 bail_bug!("thread unexpectedly had wake_on_cancel set");
2712 }
2713
2714 let task = state.get_mut(guest_thread.task)?;
2715 if let Some(event) = task.event {
2720 assert!(matches!(event, Event::None | Event::Cancelled));
2721 } else {
2722 task.event = Some(Event::None);
2723 }
2724 let call = GuestCall {
2725 thread: guest_thread,
2726 kind: GuestCallKind::DeliverEvent {
2727 instance: self,
2728 set: None,
2729 },
2730 };
2731 state.push_low_priority(WorkItem::GuestCall {
2734 instance: self.runtime_instance(runtime_instance),
2735 call,
2736 });
2737 }
2738 callback_code::WAIT => {
2739 let set = get_set(store, set)?;
2740 let state = store.concurrent_state_mut()?;
2741
2742 if state.get_mut(guest_thread.task)?.event.is_some()
2743 || !state.get_mut(set)?.ready.is_empty()
2744 {
2745 state.push_high_priority(WorkItem::GuestCall {
2747 instance: self.runtime_instance(runtime_instance),
2748 call: GuestCall {
2749 thread: guest_thread,
2750 kind: GuestCallKind::DeliverEvent {
2751 instance: self,
2752 set: Some(set),
2753 },
2754 },
2755 });
2756 } else {
2757 let old = state
2765 .get_mut(guest_thread.thread)?
2766 .wake_on_cancel
2767 .replace(WakeOnCancel::Waiting(set));
2768 if !old.is_none() {
2769 bail_bug!("thread unexpectedly had wake_on_cancel set");
2770 }
2771 let old = state
2772 .get_mut(set)?
2773 .waiting
2774 .insert(guest_thread, WaitMode::Callback(self));
2775 if !old.is_none() {
2776 bail_bug!("set's waiting set already had this thread registered");
2777 }
2778 }
2779 }
2780 _ => bail!(Trap::UnsupportedCallbackCode),
2781 }
2782
2783 Ok(())
2784 }
2785
2786 unsafe fn stage_call<T: 'static>(
2793 self,
2794 mut store: StoreContextMut<T>,
2795 guest_thread: QualifiedThreadId,
2796 callee: SendSyncPtr<VMFuncRef>,
2797 param_count: usize,
2798 result_count: usize,
2799 async_: bool,
2800 callback: Option<SendSyncPtr<VMFuncRef>>,
2801 post_return: Option<SendSyncPtr<VMFuncRef>>,
2802 host_caller: bool,
2803 ) -> Result<()> {
2804 unsafe fn make_call<T: 'static>(
2819 store: StoreContextMut<T>,
2820 guest_thread: QualifiedThreadId,
2821 callee: SendSyncPtr<VMFuncRef>,
2822 param_count: usize,
2823 result_count: usize,
2824 ) -> impl FnOnce(&mut dyn VMStore) -> Result<[MaybeUninit<ValRaw>; MAX_FLAT_PARAMS]>
2825 + Send
2826 + Sync
2827 + 'static
2828 + use<T> {
2829 let token = StoreToken::new(store);
2830 move |store: &mut dyn VMStore| {
2831 let mut storage = [MaybeUninit::uninit(); MAX_FLAT_PARAMS];
2832
2833 store
2834 .concurrent_state_mut()?
2835 .get_mut(guest_thread.thread)?
2836 .state = GuestThreadState::Running;
2837 let task = store.concurrent_state_mut()?.get_mut(guest_thread.task)?;
2838 let lower = match task.lower_params.take() {
2839 Some(l) => l,
2840 None => bail_bug!("lower_params missing"),
2841 };
2842
2843 lower(store, &mut storage[..param_count])?;
2844
2845 let mut store = token.as_context_mut(store);
2846
2847 unsafe {
2850 crate::Func::call_unchecked_raw(
2851 &mut store,
2852 callee.as_non_null(),
2853 NonNull::new(
2854 &mut storage[..param_count.max(result_count)]
2855 as *mut [MaybeUninit<ValRaw>] as _,
2856 )
2857 .unwrap(),
2858 )?;
2859 }
2860
2861 Ok(storage)
2862 }
2863 }
2864
2865 let call = unsafe {
2869 make_call(
2870 store.as_context_mut(),
2871 guest_thread,
2872 callee,
2873 param_count,
2874 result_count,
2875 )
2876 };
2877
2878 let callee_instance = store
2879 .0
2880 .concurrent_state_mut()?
2881 .get_mut(guest_thread.task)?
2882 .instance;
2883
2884 let fun = if callback.is_some() {
2885 assert!(async_);
2886
2887 Box::new(move |store: &mut dyn VMStore| {
2888 self.add_guest_thread_to_instance_table(
2889 guest_thread.thread,
2890 store,
2891 callee_instance.index,
2892 )?;
2893 let old_thread = store.set_thread(guest_thread)?;
2894 log::trace!(
2895 "stackless call: replaced {old_thread:?} with {guest_thread:?} as current thread"
2896 );
2897
2898 store.enter_instance(callee_instance);
2899
2900 let storage = call(store)?;
2907
2908 store.exit_instance(callee_instance)?;
2909
2910 store.set_thread(old_thread)?;
2911 let state = store.concurrent_state_mut()?;
2912 if let Some(t) = old_thread.guest() {
2913 state.get_mut(t.thread)?.state = GuestThreadState::Running;
2914 }
2915 log::trace!("stackless call: restored {old_thread:?} as current thread");
2916
2917 let code = unsafe { storage[0].assume_init() }.get_i32() as u32;
2920
2921 self.handle_callback_code(store, guest_thread, callee_instance.index, code)
2922 }) as Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send + Sync>
2923 } else {
2924 let token = StoreToken::new(store.as_context_mut());
2925 Box::new(move |store: &mut dyn VMStore| {
2926 self.add_guest_thread_to_instance_table(
2927 guest_thread.thread,
2928 store,
2929 callee_instance.index,
2930 )?;
2931 let old_thread = store.set_thread(guest_thread)?;
2932 log::trace!(
2933 "sync/async-stackful call: replaced {old_thread:?} with {guest_thread:?} as current thread",
2934 );
2935 let flags = self.id().get(store).instance_flags(callee_instance.index);
2936
2937 let callee_async_typed = store
2938 .concurrent_state_mut()?
2939 .get_mut(guest_thread.task)?
2940 .async_typed;
2941
2942 if !async_ && callee_async_typed {
2946 store.enter_instance(callee_instance);
2947 }
2948
2949 if !callee_async_typed {
2950 store.enter_sync_call(callee_instance)?;
2951 }
2952
2953 let storage = call(store)?;
2960
2961 if !callee_async_typed {
2962 store.exit_sync_call(callee_instance)?;
2963 }
2964
2965 if !async_ {
2966 if callee_async_typed {
2972 store.exit_instance(callee_instance)?;
2973 }
2974
2975 let lift = {
2976 let state = store.concurrent_state_mut()?;
2977 if !state.get_mut(guest_thread.task)?.result.is_none() {
2978 bail_bug!("task has already produced a result");
2979 }
2980
2981 match state.get_mut(guest_thread.task)?.lift_result.take() {
2982 Some(lift) => lift,
2983 None => bail_bug!("lift_result field is missing"),
2984 }
2985 };
2986
2987 let result = (lift.lift)(store, unsafe {
2990 mem::transmute::<&[MaybeUninit<ValRaw>], &[ValRaw]>(
2991 &storage[..result_count],
2992 )
2993 })?;
2994
2995 let post_return_arg = match result_count {
2996 0 => ValRaw::i32(0),
2997 1 => unsafe { storage[0].assume_init() },
3000 _ => unreachable!(),
3001 };
3002
3003 unsafe {
3004 call_post_return(
3005 token.as_context_mut(store),
3006 post_return.map(|v| v.as_non_null()),
3007 post_return_arg,
3008 flags,
3009 )?;
3010 }
3011
3012 self.task_complete(store, guest_thread.task, result, Status::Returned)?;
3013 }
3014
3015 store.set_thread(old_thread)?;
3016
3017 store
3018 .concurrent_state_mut()?
3019 .get_mut(guest_thread.task)?
3020 .exited = true;
3021
3022 log::trace!(
3023 "clean up thread; async lifted? {async_} async typed? {callee_async_typed}"
3024 );
3025
3026 if callee_async_typed {
3027 store.switch_or_trap_if_may_not_suspend(callee_instance)?;
3032 }
3033
3034 store.cleanup_thread(guest_thread, callee_instance, CleanupTask::Yes)?;
3036 Ok(())
3037 })
3038 };
3039
3040 store.0.concurrent_state_mut()?.push_work_item(
3041 WorkItem::GuestCall {
3042 instance: callee_instance,
3043 call: GuestCall {
3044 thread: guest_thread,
3045 kind: GuestCallKind::StartImplicit(fun),
3046 },
3047 },
3048 if host_caller {
3049 Priority::High
3050 } else {
3051 Priority::Switch
3052 },
3053 )?;
3054
3055 Ok(())
3056 }
3057
3058 unsafe fn prepare_call<T: 'static>(
3071 self,
3072 mut store: StoreContextMut<T>,
3073 start: NonNull<VMFuncRef>,
3074 return_: NonNull<VMFuncRef>,
3075 caller_instance: RuntimeComponentInstanceIndex,
3076 callee_instance: RuntimeComponentInstanceIndex,
3077 task_return_type: TypeTupleIndex,
3078 callee_async_typed: bool,
3079 memory: *mut VMMemoryDefinition,
3080 string_encoding: StringEncoding,
3081 caller_info: CallerInfo,
3082 ) -> Result<()> {
3083 enum ResultInfo {
3084 Heap { results: u32 },
3085 Stack { result_count: u32 },
3086 }
3087
3088 let result_info = match &caller_info {
3089 CallerInfo::Async {
3090 has_result: true,
3091 params,
3092 } => ResultInfo::Heap {
3093 results: match params.last() {
3094 Some(r) => r.get_u32(),
3095 None => bail_bug!("retptr missing"),
3096 },
3097 },
3098 CallerInfo::Async {
3099 has_result: false, ..
3100 } => ResultInfo::Stack { result_count: 0 },
3101 CallerInfo::Sync {
3102 result_count,
3103 params,
3104 } if *result_count > u32::try_from(MAX_FLAT_RESULTS)? => ResultInfo::Heap {
3105 results: match params.last() {
3106 Some(r) => r.get_u32(),
3107 None => bail_bug!("arg ptr missing"),
3108 },
3109 },
3110 CallerInfo::Sync { result_count, .. } => ResultInfo::Stack {
3111 result_count: *result_count,
3112 },
3113 };
3114
3115 let sync_caller = matches!(caller_info, CallerInfo::Sync { .. });
3116
3117 let start = SendSyncPtr::new(start);
3121 let return_ = SendSyncPtr::new(return_);
3122 let token = StoreToken::new(store.as_context_mut());
3123 let old_thread = store.0.current_guest_thread()?;
3124
3125 let state = store.0.concurrent_state_mut()?;
3126
3127 debug_assert_eq!(
3128 state.get_mut(old_thread.task)?.instance,
3129 self.runtime_instance(caller_instance)
3130 );
3131
3132 let guest_thread = GuestTask::new(
3133 state,
3134 Box::new(move |store, dst| {
3135 let mut store = token.as_context_mut(store);
3136 assert!(dst.len() <= MAX_FLAT_PARAMS);
3137 let mut src = [MaybeUninit::uninit(); MAX_FLAT_PARAMS + 1];
3139 let count = match caller_info {
3140 CallerInfo::Async { params, has_result } => {
3144 let params = ¶ms[..params.len() - usize::from(has_result)];
3145 for (param, src) in params.iter().zip(&mut src) {
3146 src.write(*param);
3147 }
3148 params.len()
3149 }
3150
3151 CallerInfo::Sync { params, .. } => {
3153 for (param, src) in params.iter().zip(&mut src) {
3154 src.write(*param);
3155 }
3156 params.len()
3157 }
3158 };
3159 unsafe {
3166 crate::Func::call_unchecked_raw(
3167 &mut store,
3168 start.as_non_null(),
3169 NonNull::new(
3170 &mut src[..count.max(dst.len())] as *mut [MaybeUninit<ValRaw>] as _,
3171 )
3172 .unwrap(),
3173 )?;
3174 }
3175 dst.copy_from_slice(&src[..dst.len()]);
3176 let task = store.0.current_guest_thread()?.task;
3177 let state = store.0.concurrent_state_mut()?;
3178 Waitable::Guest(task).set_event(
3179 state,
3180 Some(Event::Subtask {
3181 status: Status::Started,
3182 }),
3183 )?;
3184 Ok(())
3185 }),
3186 LiftResult {
3187 lift: Box::new(move |store, src| {
3188 let mut store = token.as_context_mut(store);
3191 let mut my_src = src.to_owned(); if let ResultInfo::Heap { results } = &result_info {
3193 my_src.push(ValRaw::u32(*results));
3194 }
3195
3196 unsafe {
3203 crate::Func::call_unchecked_raw(
3204 &mut store,
3205 return_.as_non_null(),
3206 my_src.as_mut_slice().into(),
3207 )?;
3208 }
3209
3210 let thread = store.0.current_guest_thread()?;
3211 let state = store.0.concurrent_state_mut()?;
3212 if sync_caller {
3213 state.get_mut(thread.task)?.sync_result = SyncResult::Produced(
3214 if let ResultInfo::Stack { result_count } = &result_info {
3215 match result_count {
3216 0 => None,
3217 1 => Some(my_src[0]),
3218 _ => unreachable!(),
3219 }
3220 } else {
3221 None
3222 },
3223 );
3224 }
3225 Ok(Box::new(DummyResult) as Box<dyn Any + Send + Sync>)
3226 }),
3227 ty: task_return_type,
3228 memory: NonNull::new(memory).map(SendSyncPtr::new),
3229 string_encoding,
3230 },
3231 Caller::Guest { thread: old_thread },
3232 None,
3233 self.runtime_instance(callee_instance),
3234 callee_async_typed,
3235 false,
3238 )?;
3239
3240 store.0.set_thread(guest_thread)?;
3243 log::trace!("pushed {guest_thread:?} as current thread; old thread was {old_thread:?}");
3244
3245 Ok(())
3246 }
3247
3248 unsafe fn call_callback<T>(
3253 self,
3254 mut store: StoreContextMut<T>,
3255 function: SendSyncPtr<VMFuncRef>,
3256 event: Event,
3257 handle: u32,
3258 ) -> Result<u32> {
3259 let (ordinal, result) = event.parts();
3260 let params = &mut [
3261 ValRaw::u32(ordinal),
3262 ValRaw::u32(handle),
3263 ValRaw::u32(result),
3264 ];
3265 unsafe {
3270 crate::Func::call_unchecked_raw(
3271 &mut store,
3272 function.as_non_null(),
3273 params.as_mut_slice().into(),
3274 )?;
3275 }
3276 Ok(params[0].get_u32())
3277 }
3278
3279 unsafe fn start_call<T: 'static>(
3292 self,
3293 mut store: StoreContextMut<T>,
3294 callback: *mut VMFuncRef,
3295 post_return: *mut VMFuncRef,
3296 callee: NonNull<VMFuncRef>,
3297 param_count: u32,
3298 result_count: u32,
3299 flags: u32,
3300 storage: Option<&mut [MaybeUninit<ValRaw>]>,
3301 ) -> Result<u32> {
3302 let token = StoreToken::new(store.as_context_mut());
3303 let async_caller = storage.is_none();
3304 let guest_thread = store.0.current_guest_thread()?;
3305 let state = store.0.concurrent_state_mut()?;
3306
3307 if !state.event_loop_running {
3308 bail_bug!("Instance::start_call called without a running event loop");
3309 }
3310
3311 let callee = SendSyncPtr::new(callee);
3312 let param_count = usize::try_from(param_count)?;
3313 assert!(param_count <= MAX_FLAT_PARAMS);
3314 let result_count = usize::try_from(result_count)?;
3315 assert!(result_count <= MAX_FLAT_RESULTS);
3316
3317 let task = state.get_mut(guest_thread.task)?;
3318 let callee_async_typed = task.async_typed;
3319 let callee_instance = task.instance;
3320
3321 task.async_lifted = (flags & START_FLAG_ASYNC_CALLEE) != 0;
3322
3323 if let Some(callback) = NonNull::new(callback) {
3324 let callback = SendSyncPtr::new(callback);
3328 task.callback = Some(Box::new(move |store, event, handle| {
3329 let store = token.as_context_mut(store);
3330 unsafe { self.call_callback::<T>(store, callback, event, handle) }
3331 }));
3332 }
3333
3334 let Caller::Guest { thread: caller } = &task.caller else {
3335 bail_bug!("start_call unexpectedly invoked for host->guest call");
3338 };
3339 let caller = *caller;
3340 let caller_instance = state.get_mut(caller.task)?.instance;
3341
3342 unsafe {
3344 self.stage_call(
3345 store.as_context_mut(),
3346 guest_thread,
3347 callee,
3348 param_count,
3349 result_count,
3350 (flags & START_FLAG_ASYNC_CALLEE) != 0,
3351 NonNull::new(callback).map(SendSyncPtr::new),
3352 NonNull::new(post_return).map(SendSyncPtr::new),
3353 false,
3354 )?;
3355 }
3356
3357 let old_do_not_suspend = if callee_async_typed {
3358 let state = store.0.instance_state(callee_instance).concurrent_state();
3365 let old_do_not_suspend = state.do_not_suspend;
3366 state.do_not_suspend = false;
3367 Some(old_do_not_suspend)
3368 } else {
3369 None
3370 };
3371
3372 let state = store.0.concurrent_state_mut()?;
3373
3374 let guest_waitable = Waitable::Guest(guest_thread.task);
3377 let old_set = guest_waitable.common(state)?.set;
3378 let set = state.get_mut(caller.thread)?.sync_call_set;
3379 guest_waitable.join(state, Some(set))?;
3380
3381 store.0.set_thread(CurrentThread::None)?;
3382
3383 let mut yielded = false;
3399 let (status, waitable) = loop {
3400 store.0.suspend(if yielded {
3401 SuspendReason::Waiting {
3402 set,
3403 thread: caller,
3404 }
3405 } else {
3406 yielded = true;
3407 SuspendReason::YieldingToSubtask { thread: caller }
3408 })?;
3409
3410 if let Some(old_do_not_suspend) = old_do_not_suspend {
3411 store
3412 .0
3413 .instance_state(callee_instance)
3414 .concurrent_state()
3415 .do_not_suspend = old_do_not_suspend;
3416 }
3417
3418 let state = store.0.concurrent_state_mut()?;
3419
3420 log::trace!("taking event for {:?}", guest_thread.task);
3421 let event = guest_waitable.take_event(state)?;
3422 let Some(Event::Subtask { status }) = event else {
3423 bail_bug!("subtasks should only get subtask events, got {event:?}")
3424 };
3425
3426 log::trace!("status {status:?} for {:?}", guest_thread.task);
3427
3428 if status == Status::Returned {
3429 break (status, None);
3431 } else if async_caller {
3432 let handle = store
3436 .0
3437 .instance_state(caller_instance)
3438 .handle_table()
3439 .subtask_insert_guest(guest_thread.task.rep())?;
3440 store
3441 .0
3442 .concurrent_state_mut()?
3443 .get_mut(guest_thread.task)?
3444 .common
3445 .handle = Some(handle);
3446 break (status, Some(handle));
3447 } else {
3448 store.0.switch_or_trap_if_may_not_suspend(caller_instance)?;
3452 }
3453 };
3454
3455 guest_waitable.join(store.0.concurrent_state_mut()?, old_set)?;
3456
3457 store.0.set_thread(caller)?;
3459 store
3460 .0
3461 .concurrent_state_mut()?
3462 .get_mut(caller.thread)?
3463 .state = GuestThreadState::Running;
3464 log::trace!("popped current thread {guest_thread:?}; new thread is {caller:?}");
3465
3466 if let Some(storage) = storage {
3467 let state = store.0.concurrent_state_mut()?;
3471 let task = state.get_mut(guest_thread.task)?;
3472 if let Some(result) = task.sync_result.take()? {
3473 if let Some(result) = result {
3474 storage[0] = MaybeUninit::new(result);
3475 }
3476
3477 if task.exited && task.ready_to_delete() {
3478 Waitable::Guest(guest_thread.task).delete_from(store.0)?;
3479 }
3480 }
3481 }
3482
3483 Ok(status.pack(waitable))
3484 }
3485
3486 pub(crate) fn first_poll<T: 'static, R: Send + 'static>(
3502 self,
3503 mut store: StoreContextMut<'_, T>,
3504 host_task: EnteredHostTask,
3505 future: impl Future<Output = Result<R>> + Send + 'static,
3506 lower: impl FnOnce(StoreContextMut<T>, Option<R>, bool, Option<TableId<HostTask>>) -> Result<()>
3507 + Send
3508 + 'static,
3509 ) -> Result<u32> {
3510 let token = StoreToken::new(store.as_context_mut());
3511
3512 let (join_handle, future) = JoinHandle::run(future);
3515 let mut future = Box::pin(future);
3516
3517 let poll = tls::set(store.0, || {
3522 future
3523 .as_mut()
3524 .poll(&mut Context::from_waker(&Waker::noop()))
3525 });
3526
3527 match poll {
3528 Poll::Ready(result) => {
3530 let result = result.transpose()?;
3531 let task = store.0.current_materialized_host_task()?;
3534 lower(store.as_context_mut(), result, true, task)?;
3535 return Ok(Status::Returned.pack(None));
3536 }
3537
3538 Poll::Pending => {}
3540 }
3541
3542 let Some(task) = store.0.materialize_host_task_id()? else {
3546 bail_bug!("current thread is not a host thread")
3547 };
3548 {
3549 let state = &mut store.0.concurrent_state_mut()?.get_mut(task)?.state;
3550 assert!(matches!(state, HostTaskState::CalleeStarted));
3551 *state = HostTaskState::CalleeRunning(join_handle);
3552 }
3553
3554 let future = Box::pin(async move {
3562 let result = match run_with_host_task_set(task, future).await? {
3563 Some(result) => Some(result?),
3564 None => None,
3565 };
3566 let on_complete = move |store: &mut dyn VMStore| {
3567 let mut store = token.as_context_mut(store);
3571 let old = store.0.set_thread(task)?;
3572
3573 let status = if result.is_some() {
3574 Status::Returned
3575 } else {
3576 Status::ReturnCancelled
3577 };
3578
3579 lower(store.as_context_mut(), result, false, Some(task))?;
3580 let state = store.0.concurrent_state_mut()?;
3581 match &mut state.get_mut(task)?.state {
3582 pending @ HostTaskState::CalleeCancelling => {
3585 *pending = HostTaskState::CalleeDone { cancelled: true };
3586 }
3587
3588 other => *other = HostTaskState::CalleeDone { cancelled: false },
3590 }
3591 Waitable::Host(task).set_event(state, Some(Event::Subtask { status }))?;
3592
3593 store.0.set_thread(old)?;
3594 Ok(())
3595 };
3596
3597 tls::get(move |store| {
3602 store
3603 .concurrent_state_mut()?
3604 .push_high_priority(WorkItem::WorkerFunction(AlwaysMut::new(Box::new(
3605 on_complete,
3606 ))));
3607 Ok(())
3608 })
3609 });
3610
3611 let caller = match host_task {
3614 Some(caller) => caller,
3615 None => bail_bug!("host task wasn't created but should have been"),
3616 };
3617 let state = store.0.concurrent_state_mut()?;
3618 state.push_future(future);
3619 let instance = state.get_mut(caller.task)?.instance;
3620 let handle = store
3621 .0
3622 .instance_state(instance)
3623 .handle_table()
3624 .subtask_insert_host(task.rep())?;
3625 store.0.concurrent_state_mut()?.get_mut(task)?.common.handle = Some(handle);
3626 log::trace!("assign {task:?} handle {handle} for {caller:?} instance {instance:?}");
3627
3628 store.0.set_thread(caller)?;
3632 Ok(Status::Started.pack(Some(handle)))
3633 }
3634
3635 pub(crate) fn task_return(
3638 self,
3639 store: &mut dyn VMStore,
3640 ty: TypeTupleIndex,
3641 options: OptionsIndex,
3642 storage: &[ValRaw],
3643 ) -> Result<()> {
3644 let guest_thread = store.current_guest_thread()?;
3645 let state = store.concurrent_state_mut()?;
3646 let lift = state
3647 .get_mut(guest_thread.task)?
3648 .lift_result
3649 .take()
3650 .ok_or_else(|| Trap::TaskCancelOrReturnTwice)?;
3651 if !state.get_mut(guest_thread.task)?.result.is_none() {
3652 bail_bug!("task result unexpectedly already set");
3653 }
3654
3655 let CanonicalOptions {
3656 string_encoding,
3657 data_model,
3658 ..
3659 } = &self.id().get(store).component().env_component().options[options];
3660
3661 let invalid = ty != lift.ty
3662 || string_encoding != &lift.string_encoding
3663 || match data_model {
3664 CanonicalOptionsDataModel::LinearMemory(opts) => match opts.memory {
3665 Some(memory) => {
3666 let expected = lift.memory.map(|v| v.as_ptr()).unwrap_or(ptr::null_mut());
3667 let actual = self.id().get(store).runtime_memory(memory);
3668 expected != actual.as_ptr()
3669 }
3670 None => false,
3673 },
3674 CanonicalOptionsDataModel::Gc { .. } => true,
3676 };
3677
3678 if invalid {
3679 bail!(Trap::TaskReturnInvalid);
3680 }
3681
3682 log::trace!("task.return for {guest_thread:?}");
3683
3684 let result = (lift.lift)(store, storage)?;
3685 self.task_complete(store, guest_thread.task, result, Status::Returned)
3686 }
3687
3688 pub(crate) fn task_cancel(self, store: &mut StoreOpaque) -> Result<()> {
3690 let guest_thread = store.current_guest_thread()?;
3691 let state = store.concurrent_state_mut()?;
3692 let task = state.get_mut(guest_thread.task)?;
3693 if !task.cancel_request_delivered {
3694 bail!(Trap::TaskCancelNotCancelled);
3695 }
3696 _ = task
3697 .lift_result
3698 .take()
3699 .ok_or_else(|| Trap::TaskCancelOrReturnTwice)?;
3700
3701 if !task.result.is_none() {
3702 bail_bug!("task result should not bet set yet");
3703 }
3704
3705 log::trace!("task.cancel for {guest_thread:?}");
3706
3707 self.task_complete(
3708 store,
3709 guest_thread.task,
3710 Box::new(DummyResult),
3711 Status::ReturnCancelled,
3712 )
3713 }
3714
3715 fn task_complete(
3721 self,
3722 store: &mut StoreOpaque,
3723 guest_task: TableId<GuestTask>,
3724 result: Box<dyn Any + Send + Sync>,
3725 status: Status,
3726 ) -> Result<()> {
3727 store
3728 .component_resource_tables(Some(self))?
3729 .validate_scope_exit()?;
3730
3731 let state = store.concurrent_state_mut()?;
3732 let task = state.get_mut(guest_task)?;
3733
3734 if let Caller::Host { tx, .. } = &mut task.caller {
3735 if let Some(tx) = tx.take() {
3736 _ = tx.send(result);
3737 }
3738 } else {
3739 task.result = Some(result);
3740 Waitable::Guest(guest_task).set_event(state, Some(Event::Subtask { status }))?;
3741 }
3742
3743 Ok(())
3744 }
3745
3746 pub(crate) fn waitable_set_new(
3748 self,
3749 store: &mut StoreOpaque,
3750 caller_instance: RuntimeComponentInstanceIndex,
3751 ) -> Result<u32> {
3752 let set = store.concurrent_state_mut()?.push(WaitableSet::default())?;
3753 let handle = store
3754 .instance_state(self.runtime_instance(caller_instance))
3755 .handle_table()
3756 .waitable_set_insert(set.rep())?;
3757 log::trace!("new waitable set {set:?} (handle {handle})");
3758 Ok(handle)
3759 }
3760
3761 pub(crate) fn waitable_set_drop(
3763 self,
3764 store: &mut StoreOpaque,
3765 caller_instance: RuntimeComponentInstanceIndex,
3766 set: u32,
3767 ) -> Result<()> {
3768 let rep = store
3769 .instance_state(self.runtime_instance(caller_instance))
3770 .handle_table()
3771 .waitable_set_remove(set)?;
3772
3773 log::trace!("drop waitable set {rep} (handle {set})");
3774
3775 if !store
3779 .concurrent_state_mut()?
3780 .get_mut(TableId::<WaitableSet>::new(rep))?
3781 .waiting
3782 .is_empty()
3783 {
3784 bail!(Trap::WaitableSetDropHasWaiters);
3785 }
3786
3787 store
3788 .concurrent_state_mut()?
3789 .delete(TableId::<WaitableSet>::new(rep))?;
3790
3791 Ok(())
3792 }
3793
3794 pub(crate) fn waitable_join(
3796 self,
3797 store: &mut StoreOpaque,
3798 caller_instance: RuntimeComponentInstanceIndex,
3799 waitable_handle: u32,
3800 set_handle: u32,
3801 ) -> Result<()> {
3802 let mut instance = self.id().get_mut(store);
3803 let waitable =
3804 Waitable::from_instance(instance.as_mut(), caller_instance, waitable_handle)?;
3805
3806 let set = if set_handle == 0 {
3807 None
3808 } else {
3809 let set = instance.instance_states().0[caller_instance]
3810 .handle_table()
3811 .waitable_set_rep(set_handle)?;
3812
3813 let state = store.concurrent_state_mut()?;
3814 if let Some(old) = waitable.common(state)?.set
3815 && state.get_mut(old)?.is_sync_call_set
3816 {
3817 bail!(Trap::WaitableSyncAndAsync);
3818 }
3819
3820 Some(TableId::<WaitableSet>::new(set))
3821 };
3822
3823 log::trace!(
3824 "waitable {waitable:?} (handle {waitable_handle}) join set {set:?} (handle {set_handle})",
3825 );
3826
3827 waitable.join(store.concurrent_state_mut()?, set)
3828 }
3829
3830 pub(crate) fn subtask_drop(
3832 self,
3833 store: &mut StoreOpaque,
3834 caller_instance: RuntimeComponentInstanceIndex,
3835 task_id: u32,
3836 ) -> Result<()> {
3837 self.waitable_join(store, caller_instance, task_id, 0)?;
3838
3839 let (rep, is_host) = store
3840 .instance_state(self.runtime_instance(caller_instance))
3841 .handle_table()
3842 .subtask_remove(task_id)?;
3843
3844 let concurrent_state = store.concurrent_state_mut()?;
3845 let (waitable, delete) = if is_host {
3846 let id = TableId::<HostTask>::new(rep);
3847 let task = concurrent_state.get_mut(id)?;
3848 match &task.state {
3849 HostTaskState::CalleeRunning(_) | HostTaskState::CalleeCancelling => {
3850 bail!(Trap::SubtaskDropNotResolved)
3851 }
3852 HostTaskState::CalleeDone { .. } => {}
3853 HostTaskState::CalleeStarted | HostTaskState::CalleeFinished(_) => {
3854 bail_bug!("invalid state for callee in `subtask.drop`")
3855 }
3856 }
3857
3858 (Waitable::Host(id), true)
3859 } else {
3860 let id = TableId::<GuestTask>::new(rep);
3861 let task = concurrent_state.get_mut(id)?;
3862 if task.lift_result.is_some() {
3863 bail!(Trap::SubtaskDropNotResolved);
3864 }
3865 (
3866 Waitable::Guest(id),
3867 concurrent_state.get_mut(id)?.ready_to_delete(),
3868 )
3869 };
3870
3871 waitable.common(concurrent_state)?.handle = None;
3872
3873 if waitable.take_event(concurrent_state)?.is_some() {
3876 bail!(Trap::SubtaskDropNotResolved);
3877 }
3878
3879 if delete {
3880 waitable.delete_from(store)?;
3881 }
3882
3883 log::trace!("subtask_drop {waitable:?} (handle {task_id})");
3884 Ok(())
3885 }
3886
3887 pub(crate) fn waitable_set_wait(
3889 self,
3890 store: &mut StoreOpaque,
3891 options: OptionsIndex,
3892 set: u32,
3893 payload: u32,
3894 ) -> Result<u32> {
3895 let &CanonicalOptions {
3896 instance: caller_instance,
3897 ..
3898 } = &self.id().get(store).component().env_component().options[options];
3899 let caller = self.runtime_instance(caller_instance);
3900 let rep = store
3901 .instance_state(self.runtime_instance(caller_instance))
3902 .handle_table()
3903 .waitable_set_rep(set)?;
3904
3905 self.waitable_check(
3906 store,
3907 caller,
3908 WaitableCheck::Wait,
3909 WaitableCheckParams {
3910 set: TableId::new(rep),
3911 options,
3912 payload,
3913 },
3914 )
3915 }
3916
3917 pub(crate) fn waitable_set_poll(
3919 self,
3920 store: &mut StoreOpaque,
3921 options: OptionsIndex,
3922 set: u32,
3923 payload: u32,
3924 ) -> Result<u32> {
3925 let &CanonicalOptions {
3926 instance: caller_instance,
3927 ..
3928 } = &self.id().get(store).component().env_component().options[options];
3929 let caller = self.runtime_instance(caller_instance);
3930 let rep = store
3931 .instance_state(caller)
3932 .handle_table()
3933 .waitable_set_rep(set)?;
3934
3935 self.waitable_check(
3936 store,
3937 caller,
3938 WaitableCheck::Poll,
3939 WaitableCheckParams {
3940 set: TableId::new(rep),
3941 options,
3942 payload,
3943 },
3944 )
3945 }
3946
3947 pub(crate) fn thread_index(&self, store: &mut dyn VMStore) -> Result<u32> {
3949 let thread_id = store.current_guest_thread()?.thread;
3950 match store
3951 .concurrent_state_mut()?
3952 .get_mut(thread_id)?
3953 .instance_rep
3954 {
3955 Some(r) => Ok(r),
3956 None => bail_bug!("thread should have instance_rep by now"),
3957 }
3958 }
3959
3960 pub(crate) fn thread_new_indirect<T: 'static>(
3962 self,
3963 mut store: StoreContextMut<T>,
3964 runtime_instance: RuntimeComponentInstanceIndex,
3965 _func_ty_idx: TypeFuncIndex, start_func_table_idx: RuntimeTableIndex,
3967 start_func_idx: u32,
3968 context: i32,
3969 ) -> Result<u32> {
3970 log::trace!("creating new thread");
3971
3972 let start_func_ty = FuncType::new(store.engine(), [ValType::I32], []);
3973 let (instance, registry) = self.id().get_mut_and_registry(store.0);
3974 let callee = instance
3975 .index_runtime_func_table(registry, start_func_table_idx, start_func_idx as u64)?
3976 .ok_or_else(|| Trap::ThreadNewIndirectUninitialized)?;
3977 if callee.type_index(store.0) != start_func_ty.type_index() {
3978 bail!(Trap::ThreadNewIndirectInvalidType);
3979 }
3980
3981 let token = StoreToken::new(store.as_context_mut());
3982 let start_func = Box::new(
3983 move |store: &mut dyn VMStore, guest_thread: QualifiedThreadId| -> Result<()> {
3984 let old_thread = store.set_thread(guest_thread)?;
3985 log::trace!(
3986 "thread start: replaced {old_thread:?} with {guest_thread:?} as current thread"
3987 );
3988
3989 let mut store = token.as_context_mut(store);
3990 let mut params = [ValRaw::i32(context)];
3991 unsafe { callee.call_unchecked(store.as_context_mut(), &mut params)? };
3994
3995 store.0.set_thread(old_thread)?;
3996
3997 let runtime_instance = self.runtime_instance(runtime_instance);
3998
3999 store
4002 .0
4003 .switch_or_trap_if_may_not_suspend(runtime_instance)?;
4004
4005 store
4006 .0
4007 .cleanup_thread(guest_thread, runtime_instance, CleanupTask::Yes)?;
4008
4009 log::trace!("explicit thread {guest_thread:?} completed");
4010 let state = store.0.concurrent_state_mut()?;
4011 if let Some(t) = old_thread.guest() {
4012 state.get_mut(t.thread)?.state = GuestThreadState::Running;
4013 }
4014 log::trace!("thread start: restored {old_thread:?} as current thread");
4015
4016 Ok(())
4017 },
4018 );
4019
4020 let current_thread = store.0.current_guest_thread()?;
4021 let state = store.0.concurrent_state_mut()?;
4022 let parent_task = current_thread.task;
4023
4024 let new_thread = GuestThread::new_explicit(state, parent_task, start_func)?;
4025 let thread_id = state.push(new_thread)?;
4026 state.get_mut(parent_task)?.threads.insert(thread_id);
4027
4028 log::trace!("new thread with id {thread_id:?} created");
4029
4030 self.add_guest_thread_to_instance_table(thread_id, store.0, runtime_instance)
4031 }
4032
4033 pub(crate) fn resume_thread(
4034 self,
4035 store: &mut StoreOpaque,
4036 runtime_instance: RuntimeComponentInstanceIndex,
4037 thread_idx: u32,
4038 how: ResumeThread,
4039 ) -> Result<bool> {
4040 let thread_id =
4041 GuestThread::from_instance(self.id().get_mut(store), runtime_instance, thread_idx)?;
4042 let state = store.concurrent_state_mut()?;
4043 let guest_thread = QualifiedThreadId::qualify(state, thread_id)?;
4044
4045 if store.current_guest_thread()? == guest_thread {
4046 bail!(Trap::CannotResumeThread);
4047 }
4048
4049 let state = store.concurrent_state_mut()?;
4050 let thread = state.get_mut(guest_thread.thread)?;
4051 let priority = match how {
4052 ResumeThread::Promote | ResumeThread::Resume => Priority::Switch,
4053 ResumeThread::ResumeLater => Priority::Low,
4054 };
4055
4056 match (&how, &thread.state) {
4057 (ResumeThread::Promote, GuestThreadState::Ready { .. }) => {}
4059 (ResumeThread::Promote, _) => return Ok(false),
4060
4061 (
4064 ResumeThread::Resume | ResumeThread::ResumeLater,
4065 GuestThreadState::NotStartedExplicit(_) | GuestThreadState::Suspended(_),
4066 ) => {}
4067 (ResumeThread::Resume | ResumeThread::ResumeLater, _) => {
4068 bail!(Trap::CannotResumeThread)
4069 }
4070 }
4071
4072 match mem::replace(&mut thread.state, GuestThreadState::Running) {
4073 GuestThreadState::NotStartedExplicit(start_func) => {
4074 log::trace!("starting thread {guest_thread:?}");
4075 let guest_call = WorkItem::GuestCall {
4076 instance: self.runtime_instance(runtime_instance),
4077 call: GuestCall {
4078 thread: guest_thread,
4079 kind: GuestCallKind::StartExplicit(Box::new(move |store| {
4080 start_func(store, guest_thread)
4081 })),
4082 },
4083 };
4084 store
4085 .concurrent_state_mut()?
4086 .push_work_item(guest_call, priority)?;
4087 }
4088 GuestThreadState::Suspended(fiber) => {
4089 log::trace!("resuming thread {thread_id:?} that was suspended");
4090 store.concurrent_state_mut()?.push_work_item(
4091 WorkItem::ResumeFiber {
4092 instance: self.runtime_instance(runtime_instance),
4093 thread: guest_thread,
4094 fiber,
4095 },
4096 priority,
4097 )?;
4098 }
4099 GuestThreadState::Ready { fiber } => {
4100 log::trace!("resuming thread {thread_id:?} that was ready");
4101 thread.state = GuestThreadState::Ready { fiber };
4102 store
4103 .concurrent_state_mut()?
4104 .promote_thread_work_item(guest_thread)?;
4105 }
4106 other @ (GuestThreadState::NotStartedImplicit
4107 | GuestThreadState::Running
4108 | GuestThreadState::Completed) => {
4109 thread.state = other;
4110 }
4111 }
4112 Ok(true)
4113 }
4114
4115 fn add_guest_thread_to_instance_table(
4116 self,
4117 thread_id: TableId<GuestThread>,
4118 store: &mut StoreOpaque,
4119 runtime_instance: RuntimeComponentInstanceIndex,
4120 ) -> Result<u32> {
4121 let guest_id = store
4122 .instance_state(self.runtime_instance(runtime_instance))
4123 .thread_handle_table()
4124 .guest_thread_insert(thread_id.rep())?;
4125 store
4126 .concurrent_state_mut()?
4127 .get_mut(thread_id)?
4128 .instance_rep = Some(guest_id);
4129 Ok(guest_id)
4130 }
4131
4132 pub(crate) fn suspension_intrinsic(
4136 self,
4137 store: &mut StoreOpaque,
4138 caller: RuntimeComponentInstanceIndex,
4139 yielding: bool,
4140 to_thread: SuspensionTarget,
4141 ) -> Result<WaitResult> {
4142 let check_suspend = match to_thread {
4143 SuspensionTarget::Promote(thread) => {
4144 !self.resume_thread(store, caller, thread, ResumeThread::Promote)?
4145 }
4146 SuspensionTarget::Resume(thread) => {
4147 if !self.resume_thread(store, caller, thread, ResumeThread::Resume)? {
4148 bail_bug!(
4149 "`resume_thread` should only ever return false \
4150 when `ResumeThread::Promote` is passed to it"
4151 );
4152 }
4153 false
4154 }
4155 SuspensionTarget::None => true,
4156 };
4157
4158 if check_suspend && !store.switch_if_may_not_suspend(self.runtime_instance(caller))? {
4159 return if yielding {
4160 Ok(WaitResult::Completed)
4161 } else {
4162 Err(Trap::CannotBlockSyncTask.into())
4163 };
4164 }
4165
4166 let guest_thread = store.current_guest_thread()?;
4167
4168 let reason = if yielding {
4169 SuspendReason::Yielding {
4170 thread: guest_thread,
4171 }
4172 } else {
4173 SuspendReason::ExplicitlySuspending {
4174 thread: guest_thread,
4175 }
4176 };
4177
4178 store.suspend(reason)?;
4179
4180 Ok(WaitResult::Completed)
4181 }
4182
4183 fn waitable_check(
4185 self,
4186 store: &mut StoreOpaque,
4187 caller: RuntimeInstance,
4188 check: WaitableCheck,
4189 params: WaitableCheckParams,
4190 ) -> Result<u32> {
4191 let guest_thread = store.current_guest_thread()?;
4192
4193 log::trace!("waitable check for {guest_thread:?}; set {:?}", params.set);
4194
4195 let state = store.concurrent_state_mut()?;
4196 let task = state.get_mut(guest_thread.task)?;
4197
4198 match &check {
4201 WaitableCheck::Wait => {
4202 let set = params.set;
4203
4204 if (task.event.is_none() || matches!(task.event, Some(Event::Cancelled)))
4205 && state.get_mut(set)?.ready.is_empty()
4206 {
4207 store.switch_or_trap_if_may_not_suspend(caller)?;
4208
4209 store.suspend(SuspendReason::Waiting {
4210 set,
4211 thread: guest_thread,
4212 })?;
4213 }
4214 }
4215 WaitableCheck::Poll => {}
4216 }
4217
4218 log::trace!(
4219 "waitable check for {guest_thread:?}; set {:?}, part two",
4220 params.set
4221 );
4222
4223 let event = self.get_event(store, guest_thread.task, Some(params.set), false)?;
4225
4226 let (ordinal, handle, result) = match &check {
4227 WaitableCheck::Wait => {
4228 let (event, waitable) = match event {
4229 Some(p) => p,
4230 None => bail_bug!("event expected to be present"),
4231 };
4232 let handle = waitable.map(|(_, v)| v).unwrap_or(0);
4233 let (ordinal, result) = event.parts();
4234 (ordinal, handle, result)
4235 }
4236 WaitableCheck::Poll => {
4237 if let Some((event, waitable)) = event {
4238 let handle = waitable.map(|(_, v)| v).unwrap_or(0);
4239 let (ordinal, result) = event.parts();
4240 (ordinal, handle, result)
4241 } else {
4242 log::trace!(
4243 "no events ready to deliver via waitable-set.poll to {:?}; set {:?}",
4244 guest_thread.task,
4245 params.set
4246 );
4247 let (ordinal, result) = Event::None.parts();
4248 (ordinal, 0, result)
4249 }
4250 }
4251 };
4252 let memory = self.options_memory_mut(store, params.options);
4253 let ptr = crate::component::func::validate_inbounds_dynamic(
4254 &CanonicalAbiInfo::POINTER_PAIR,
4255 memory,
4256 &ValRaw::u32(params.payload),
4257 )?;
4258 memory[ptr + 0..][..4].copy_from_slice(&handle.to_le_bytes());
4259 memory[ptr + 4..][..4].copy_from_slice(&result.to_le_bytes());
4260 Ok(ordinal)
4261 }
4262
4263 pub(crate) fn subtask_cancel(
4265 self,
4266 store: &mut StoreOpaque,
4267 caller_instance: RuntimeComponentInstanceIndex,
4268 async_: bool,
4269 task_id: u32,
4270 ) -> Result<u32> {
4271 let (rep, is_host) = store
4272 .instance_state(self.runtime_instance(caller_instance))
4273 .handle_table()
4274 .subtask_rep(task_id)?;
4275 let waitable = if is_host {
4276 Waitable::Host(TableId::<HostTask>::new(rep))
4277 } else {
4278 Waitable::Guest(TableId::<GuestTask>::new(rep))
4279 };
4280 let concurrent_state = store.concurrent_state_mut()?;
4281
4282 log::trace!("subtask_cancel {waitable:?} (handle {task_id}; async {async_})");
4283
4284 waitable.trap_if_in_waitable_set(concurrent_state)?;
4285
4286 let needs_block;
4287 if let Waitable::Host(host_task) = waitable {
4288 let state = &mut concurrent_state.get_mut(host_task)?.state;
4289 match state {
4290 HostTaskState::CalleeRunning(handle) => {
4297 handle.abort();
4298 *state = HostTaskState::CalleeCancelling;
4299 needs_block = true;
4300 }
4301
4302 HostTaskState::CalleeCancelling | HostTaskState::CalleeDone { cancelled: true } => {
4305 bail!(Trap::SubtaskCancelAfterTerminal);
4306 }
4307 HostTaskState::CalleeDone { cancelled: false } => {
4308 *state = HostTaskState::CalleeDone { cancelled: true };
4311 needs_block = false;
4312 }
4313
4314 HostTaskState::CalleeStarted | HostTaskState::CalleeFinished(_) => {
4317 bail_bug!("invalid states for host callee")
4318 }
4319 }
4320 } else {
4321 let guest_task = TableId::<GuestTask>::new(rep);
4322 let task = concurrent_state.get_mut(guest_task)?;
4323 if !task.already_lowered_parameters() {
4324 store.cancel_guest_subtask_without_lowered_parameters(
4325 self.runtime_instance(caller_instance),
4326 guest_task,
4327 )?;
4328 return Ok(Status::StartCancelled as u32);
4329 } else if !task.returned_or_cancelled() {
4330 task.event = Some(Event::Cancelled);
4338 let runtime_instance = task.instance;
4339 for thread in task.threads.clone() {
4340 let thread = QualifiedThreadId {
4341 task: guest_task,
4342 thread,
4343 };
4344 let thread_mut = concurrent_state.get_mut(thread.thread)?;
4345
4346 let yield_ = |store: &mut StoreOpaque| {
4347 let state = store.instance_state(runtime_instance).concurrent_state();
4352 let old_do_not_suspend = state.do_not_suspend;
4353 state.do_not_suspend = false;
4354
4355 let caller = store.current_guest_thread()?;
4356
4357 let state = store.concurrent_state_mut()?;
4362 let set = state.get_mut(caller.thread)?.sync_call_set;
4363 waitable.join(state, Some(set))?;
4364
4365 store.suspend(SuspendReason::YieldingToSubtask { thread: caller })?;
4366
4367 let state = store.concurrent_state_mut()?;
4368 waitable.join(state, None)?;
4369
4370 store
4371 .instance_state(runtime_instance)
4372 .concurrent_state()
4373 .do_not_suspend = old_do_not_suspend;
4374
4375 Ok::<(), crate::Error>(())
4376 };
4377
4378 match thread_mut.wake_on_cancel.take() {
4379 WakeOnCancel::Waiting(set) => {
4380 let item = match concurrent_state.get_mut(set)?.waiting.remove(&thread)
4382 {
4383 Some(WaitMode::Callback(instance)) => WorkItem::GuestCall {
4384 instance: runtime_instance,
4385 call: GuestCall {
4386 thread,
4387 kind: GuestCallKind::DeliverEvent {
4388 instance,
4389 set: None,
4390 },
4391 },
4392 },
4393 other => bail_bug!(
4394 "expected `Some(WaitMode::Callback(_))`; got `{other:?}`"
4395 ),
4396 };
4397 concurrent_state.set_switch_item(item)?;
4398
4399 yield_(store)?;
4400
4401 break;
4402 }
4403 WakeOnCancel::Yielding => {
4404 if concurrent_state.promote_thread_work_item(thread)? {
4405 yield_(store)?;
4406 break;
4407 } else {
4408 bail_bug!("thread with `WakeOnCancel::Yielding` not promotable");
4409 }
4410 }
4411 WakeOnCancel::None => {}
4412 }
4413 }
4414
4415 needs_block = !store
4418 .concurrent_state_mut()?
4419 .get_mut(guest_task)?
4420 .returned_or_cancelled()
4421 } else {
4422 needs_block = false;
4423 }
4424 };
4425
4426 if needs_block {
4430 if async_ {
4431 return Ok(BLOCKED);
4432 }
4433
4434 let old_next_switch_item = {
4437 let state = store.concurrent_state_mut()?;
4438 let item = state.next_switch_item.take();
4439 state.push(item)?
4443 };
4444
4445 store.wait_for_event(self.runtime_instance(caller_instance), waitable)?;
4448
4449 let state = store.concurrent_state_mut()?;
4450 state.next_switch_item = state.delete(old_next_switch_item)?;
4451
4452 }
4454
4455 let event = waitable.take_event(store.concurrent_state_mut()?)?;
4456 if let Some(Event::Subtask {
4457 status: status @ (Status::Returned | Status::ReturnCancelled),
4458 }) = event
4459 {
4460 Ok(status as u32)
4461 } else {
4462 bail!(Trap::SubtaskCancelAfterTerminal);
4463 }
4464 }
4465}
4466
4467pub trait VMComponentAsyncStore {
4475 unsafe fn prepare_call(
4481 &mut self,
4482 instance: Instance,
4483 memory: *mut VMMemoryDefinition,
4484 start: NonNull<VMFuncRef>,
4485 return_: NonNull<VMFuncRef>,
4486 caller_instance: RuntimeComponentInstanceIndex,
4487 callee_instance: RuntimeComponentInstanceIndex,
4488 task_return_type: TypeTupleIndex,
4489 callee_async: bool,
4490 string_encoding: StringEncoding,
4491 result_count: u32,
4492 storage: *mut ValRaw,
4493 storage_len: usize,
4494 ) -> Result<()>;
4495
4496 unsafe fn sync_start(
4499 &mut self,
4500 instance: Instance,
4501 callback: *mut VMFuncRef,
4502 callee: NonNull<VMFuncRef>,
4503 param_count: u32,
4504 storage: *mut MaybeUninit<ValRaw>,
4505 storage_len: usize,
4506 ) -> Result<()>;
4507
4508 unsafe fn async_start(
4511 &mut self,
4512 instance: Instance,
4513 callback: *mut VMFuncRef,
4514 post_return: *mut VMFuncRef,
4515 callee: NonNull<VMFuncRef>,
4516 param_count: u32,
4517 result_count: u32,
4518 flags: u32,
4519 ) -> Result<u32>;
4520
4521 fn future_write(
4523 &mut self,
4524 instance: Instance,
4525 caller: RuntimeComponentInstanceIndex,
4526 ty: TypeFutureTableIndex,
4527 options: OptionsIndex,
4528 future: u32,
4529 address: u32,
4530 ) -> Result<u32>;
4531
4532 fn future_read(
4534 &mut self,
4535 instance: Instance,
4536 caller: RuntimeComponentInstanceIndex,
4537 ty: TypeFutureTableIndex,
4538 options: OptionsIndex,
4539 future: u32,
4540 address: u32,
4541 ) -> Result<u32>;
4542
4543 fn future_drop_writable(
4545 &mut self,
4546 instance: Instance,
4547 ty: TypeFutureTableIndex,
4548 writer: u32,
4549 ) -> Result<()>;
4550
4551 fn stream_write(
4553 &mut self,
4554 instance: Instance,
4555 caller: RuntimeComponentInstanceIndex,
4556 ty: TypeStreamTableIndex,
4557 options: OptionsIndex,
4558 stream: u32,
4559 address: u32,
4560 count: u32,
4561 ) -> Result<u32>;
4562
4563 fn stream_read(
4565 &mut self,
4566 instance: Instance,
4567 caller: RuntimeComponentInstanceIndex,
4568 ty: TypeStreamTableIndex,
4569 options: OptionsIndex,
4570 stream: u32,
4571 address: u32,
4572 count: u32,
4573 ) -> Result<u32>;
4574
4575 fn flat_stream_write(
4578 &mut self,
4579 instance: Instance,
4580 caller: RuntimeComponentInstanceIndex,
4581 ty: TypeStreamTableIndex,
4582 options: OptionsIndex,
4583 payload_size: u32,
4584 payload_align: u32,
4585 stream: u32,
4586 address: u32,
4587 count: u32,
4588 ) -> Result<u32>;
4589
4590 fn flat_stream_read(
4593 &mut self,
4594 instance: Instance,
4595 caller: RuntimeComponentInstanceIndex,
4596 ty: TypeStreamTableIndex,
4597 options: OptionsIndex,
4598 payload_size: u32,
4599 payload_align: u32,
4600 stream: u32,
4601 address: u32,
4602 count: u32,
4603 ) -> Result<u32>;
4604
4605 fn stream_drop_writable(
4607 &mut self,
4608 instance: Instance,
4609 ty: TypeStreamTableIndex,
4610 writer: u32,
4611 ) -> Result<()>;
4612
4613 fn error_context_debug_message(
4615 &mut self,
4616 instance: Instance,
4617 ty: TypeComponentLocalErrorContextTableIndex,
4618 options: OptionsIndex,
4619 err_ctx_handle: u32,
4620 debug_msg_address: u32,
4621 ) -> Result<()>;
4622
4623 fn thread_new_indirect(
4625 &mut self,
4626 instance: Instance,
4627 caller: RuntimeComponentInstanceIndex,
4628 func_ty_idx: TypeFuncIndex,
4629 start_func_table_idx: RuntimeTableIndex,
4630 start_func_idx: u32,
4631 context: i32,
4632 ) -> Result<u32>;
4633}
4634
4635impl<T: 'static> VMComponentAsyncStore for StoreInner<T> {
4637 unsafe fn prepare_call(
4638 &mut self,
4639 instance: Instance,
4640 memory: *mut VMMemoryDefinition,
4641 start: NonNull<VMFuncRef>,
4642 return_: NonNull<VMFuncRef>,
4643 caller_instance: RuntimeComponentInstanceIndex,
4644 callee_instance: RuntimeComponentInstanceIndex,
4645 task_return_type: TypeTupleIndex,
4646 callee_async: bool,
4647 string_encoding: StringEncoding,
4648 result_count_or_max_if_async: u32,
4649 storage: *mut ValRaw,
4650 storage_len: usize,
4651 ) -> Result<()> {
4652 let params = unsafe { core::slice::from_raw_parts(storage, storage_len) }.to_vec();
4656
4657 unsafe {
4658 instance.prepare_call(
4659 StoreContextMut(self),
4660 start,
4661 return_,
4662 caller_instance,
4663 callee_instance,
4664 task_return_type,
4665 callee_async,
4666 memory,
4667 string_encoding,
4668 match result_count_or_max_if_async {
4669 PREPARE_ASYNC_NO_RESULT => CallerInfo::Async {
4670 params,
4671 has_result: false,
4672 },
4673 PREPARE_ASYNC_WITH_RESULT => CallerInfo::Async {
4674 params,
4675 has_result: true,
4676 },
4677 result_count => CallerInfo::Sync {
4678 params,
4679 result_count,
4680 },
4681 },
4682 )
4683 }
4684 }
4685
4686 unsafe fn sync_start(
4687 &mut self,
4688 instance: Instance,
4689 callback: *mut VMFuncRef,
4690 callee: NonNull<VMFuncRef>,
4691 param_count: u32,
4692 storage: *mut MaybeUninit<ValRaw>,
4693 storage_len: usize,
4694 ) -> Result<()> {
4695 unsafe {
4696 instance
4697 .start_call(
4698 StoreContextMut(self),
4699 callback,
4700 ptr::null_mut(),
4701 callee,
4702 param_count,
4703 1,
4704 START_FLAG_ASYNC_CALLEE,
4705 Some(core::slice::from_raw_parts_mut(storage, storage_len)),
4709 )
4710 .map(drop)
4711 }
4712 }
4713
4714 unsafe fn async_start(
4715 &mut self,
4716 instance: Instance,
4717 callback: *mut VMFuncRef,
4718 post_return: *mut VMFuncRef,
4719 callee: NonNull<VMFuncRef>,
4720 param_count: u32,
4721 result_count: u32,
4722 flags: u32,
4723 ) -> Result<u32> {
4724 unsafe {
4725 instance.start_call(
4726 StoreContextMut(self),
4727 callback,
4728 post_return,
4729 callee,
4730 param_count,
4731 result_count,
4732 flags,
4733 None,
4734 )
4735 }
4736 }
4737
4738 fn future_write(
4739 &mut self,
4740 instance: Instance,
4741 caller: RuntimeComponentInstanceIndex,
4742 ty: TypeFutureTableIndex,
4743 options: OptionsIndex,
4744 future: u32,
4745 address: u32,
4746 ) -> Result<u32> {
4747 instance
4748 .guest_write(
4749 StoreContextMut(self),
4750 caller,
4751 TransmitIndex::Future(ty),
4752 options,
4753 None,
4754 future,
4755 address,
4756 1,
4757 )
4758 .map(|result| result.encode())
4759 }
4760
4761 fn future_read(
4762 &mut self,
4763 instance: Instance,
4764 caller: RuntimeComponentInstanceIndex,
4765 ty: TypeFutureTableIndex,
4766 options: OptionsIndex,
4767 future: u32,
4768 address: u32,
4769 ) -> Result<u32> {
4770 instance
4771 .guest_read(
4772 StoreContextMut(self),
4773 caller,
4774 TransmitIndex::Future(ty),
4775 options,
4776 None,
4777 future,
4778 address,
4779 1,
4780 )
4781 .map(|result| result.encode())
4782 }
4783
4784 fn stream_write(
4785 &mut self,
4786 instance: Instance,
4787 caller: RuntimeComponentInstanceIndex,
4788 ty: TypeStreamTableIndex,
4789 options: OptionsIndex,
4790 stream: u32,
4791 address: u32,
4792 count: u32,
4793 ) -> Result<u32> {
4794 instance
4795 .guest_write(
4796 StoreContextMut(self),
4797 caller,
4798 TransmitIndex::Stream(ty),
4799 options,
4800 None,
4801 stream,
4802 address,
4803 count,
4804 )
4805 .map(|result| result.encode())
4806 }
4807
4808 fn stream_read(
4809 &mut self,
4810 instance: Instance,
4811 caller: RuntimeComponentInstanceIndex,
4812 ty: TypeStreamTableIndex,
4813 options: OptionsIndex,
4814 stream: u32,
4815 address: u32,
4816 count: u32,
4817 ) -> Result<u32> {
4818 instance
4819 .guest_read(
4820 StoreContextMut(self),
4821 caller,
4822 TransmitIndex::Stream(ty),
4823 options,
4824 None,
4825 stream,
4826 address,
4827 count,
4828 )
4829 .map(|result| result.encode())
4830 }
4831
4832 fn future_drop_writable(
4833 &mut self,
4834 instance: Instance,
4835 ty: TypeFutureTableIndex,
4836 writer: u32,
4837 ) -> Result<()> {
4838 instance.guest_drop_writable(self, TransmitIndex::Future(ty), writer)
4839 }
4840
4841 fn flat_stream_write(
4842 &mut self,
4843 instance: Instance,
4844 caller: RuntimeComponentInstanceIndex,
4845 ty: TypeStreamTableIndex,
4846 options: OptionsIndex,
4847 payload_size: u32,
4848 payload_align: u32,
4849 stream: u32,
4850 address: u32,
4851 count: u32,
4852 ) -> Result<u32> {
4853 instance
4854 .guest_write(
4855 StoreContextMut(self),
4856 caller,
4857 TransmitIndex::Stream(ty),
4858 options,
4859 Some(FlatAbi {
4860 size: payload_size,
4861 align: payload_align,
4862 }),
4863 stream,
4864 address,
4865 count,
4866 )
4867 .map(|result| result.encode())
4868 }
4869
4870 fn flat_stream_read(
4871 &mut self,
4872 instance: Instance,
4873 caller: RuntimeComponentInstanceIndex,
4874 ty: TypeStreamTableIndex,
4875 options: OptionsIndex,
4876 payload_size: u32,
4877 payload_align: u32,
4878 stream: u32,
4879 address: u32,
4880 count: u32,
4881 ) -> Result<u32> {
4882 instance
4883 .guest_read(
4884 StoreContextMut(self),
4885 caller,
4886 TransmitIndex::Stream(ty),
4887 options,
4888 Some(FlatAbi {
4889 size: payload_size,
4890 align: payload_align,
4891 }),
4892 stream,
4893 address,
4894 count,
4895 )
4896 .map(|result| result.encode())
4897 }
4898
4899 fn stream_drop_writable(
4900 &mut self,
4901 instance: Instance,
4902 ty: TypeStreamTableIndex,
4903 writer: u32,
4904 ) -> Result<()> {
4905 instance.guest_drop_writable(self, TransmitIndex::Stream(ty), writer)
4906 }
4907
4908 fn error_context_debug_message(
4909 &mut self,
4910 instance: Instance,
4911 ty: TypeComponentLocalErrorContextTableIndex,
4912 options: OptionsIndex,
4913 err_ctx_handle: u32,
4914 debug_msg_address: u32,
4915 ) -> Result<()> {
4916 instance.error_context_debug_message(
4917 StoreContextMut(self),
4918 ty,
4919 options,
4920 err_ctx_handle,
4921 debug_msg_address,
4922 )
4923 }
4924
4925 fn thread_new_indirect(
4926 &mut self,
4927 instance: Instance,
4928 caller: RuntimeComponentInstanceIndex,
4929 func_ty_idx: TypeFuncIndex,
4930 start_func_table_idx: RuntimeTableIndex,
4931 start_func_idx: u32,
4932 context: i32,
4933 ) -> Result<u32> {
4934 instance.thread_new_indirect(
4935 StoreContextMut(self),
4936 caller,
4937 func_ty_idx,
4938 start_func_table_idx,
4939 start_func_idx,
4940 context,
4941 )
4942 }
4943}
4944
4945type HostTaskFuture = Pin<Box<dyn Future<Output = Result<()>> + Send + 'static>>;
4946
4947async fn run_with_host_task_set<F>(task: TableId<HostTask>, future: F) -> Result<F::Output>
4950where
4951 F: Future,
4952{
4953 let mut future = pin!(future);
4954 future::poll_fn(|cx| {
4955 let old_thread = match tls::get(|store| store.set_thread(task)) {
4956 Ok(thread) => thread,
4957 Err(error) => return Poll::Ready(Err(error)),
4958 };
4959 let result = future.as_mut().poll(cx);
4960 match tls::get(|store| store.set_thread(old_thread)) {
4961 Ok(_) => result.map(Ok),
4962 Err(error) => Poll::Ready(Err(error)),
4963 }
4964 })
4965 .await
4966}
4967
4968pub(crate) struct HostTask {
4972 common: WaitableCommon,
4973
4974 call_context: CallContext,
4977
4978 state: HostTaskState,
4979
4980 group: TaskGroupId,
4981}
4982
4983enum HostTaskState {
4984 CalleeStarted,
4989
4990 CalleeRunning(JoinHandle),
4995
4996 CalleeCancelling,
5000
5001 CalleeFinished(LiftedResult),
5005
5006 CalleeDone { cancelled: bool },
5009}
5010
5011impl HostTask {
5012 fn new(
5013 concurrent_state: &mut ConcurrentState,
5014 state: HostTaskState,
5015 caller: QualifiedThreadId,
5016 ) -> Result<Self> {
5017 let group = concurrent_state.get_mut(caller.task)?.group;
5018 concurrent_state.increment_group_ref_count(group)?;
5019
5020 Ok(Self {
5021 common: WaitableCommon::default(),
5022 call_context: CallContext::default(),
5023 state,
5024 group,
5025 })
5026 }
5027}
5028
5029impl TableDebug for HostTask {
5030 fn type_name() -> &'static str {
5031 "HostTask"
5032 }
5033}
5034
5035type CallbackFn = Box<dyn Fn(&mut dyn VMStore, Event, u32) -> Result<u32> + Send + Sync + 'static>;
5036
5037enum Caller {
5039 Host {
5041 tx: Option<oneshot::Sender<LiftedResult>>,
5043 host_future_present: bool,
5046 caller: Option<TableId<HostTask>>,
5050 },
5051 Guest {
5053 thread: QualifiedThreadId,
5055 },
5056}
5057
5058struct LiftResult {
5061 lift: RawLift,
5062 ty: TypeTupleIndex,
5063 memory: Option<SendSyncPtr<VMMemoryDefinition>>,
5064 string_encoding: StringEncoding,
5065}
5066
5067#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
5072pub(crate) struct QualifiedThreadId {
5073 task: TableId<GuestTask>,
5074 thread: TableId<GuestThread>,
5075}
5076
5077impl QualifiedThreadId {
5078 fn qualify(
5079 state: &mut ConcurrentState,
5080 thread: TableId<GuestThread>,
5081 ) -> Result<QualifiedThreadId> {
5082 Ok(QualifiedThreadId {
5083 task: state.get_mut(thread)?.parent_task,
5084 thread,
5085 })
5086 }
5087}
5088
5089impl fmt::Debug for QualifiedThreadId {
5090 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
5091 f.debug_tuple("QualifiedThreadId")
5092 .field(&self.task.rep())
5093 .field(&self.thread.rep())
5094 .finish()
5095 }
5096}
5097
5098enum GuestThreadState {
5099 NotStartedImplicit,
5100 NotStartedExplicit(
5101 Box<dyn FnOnce(&mut dyn VMStore, QualifiedThreadId) -> Result<()> + Send + Sync>,
5102 ),
5103 Running,
5104 Suspended(StoreFiber<'static>),
5105 Ready {
5106 fiber: StoreFiber<'static>,
5107 },
5108 Completed,
5109}
5110
5111impl fmt::Debug for GuestThreadState {
5112 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
5113 match self {
5114 Self::NotStartedImplicit => f.debug_tuple("NotStartedImplicit").finish(),
5115 Self::NotStartedExplicit(_) => f.debug_tuple("NotStartedExplicit").finish(),
5116 Self::Running => f.debug_tuple("Running").finish(),
5117 Self::Suspended(_) => f.debug_tuple("Suspended").finish(),
5118 Self::Ready { .. } => f.debug_struct("Ready").finish(),
5119 Self::Completed => f.debug_tuple("Completed").finish(),
5120 }
5121 }
5122}
5123
5124#[derive(Copy, Clone, PartialEq, Eq, Debug)]
5125enum WakeOnCancel {
5126 None,
5127 Waiting(TableId<WaitableSet>),
5128 Yielding,
5129}
5130
5131impl WakeOnCancel {
5132 fn is_none(self) -> bool {
5133 matches!(self, WakeOnCancel::None)
5134 }
5135
5136 fn replace(&mut self, other: WakeOnCancel) -> Self {
5137 let old = *self;
5138 *self = other;
5139 old
5140 }
5141
5142 fn take(&mut self) -> Self {
5143 self.replace(WakeOnCancel::None)
5144 }
5145}
5146
5147pub struct GuestThread {
5148 context: [u32; NUM_COMPONENT_CONTEXT_SLOTS],
5151 parent_task: TableId<GuestTask>,
5153 wake_on_cancel: WakeOnCancel,
5156 state: GuestThreadState,
5158 instance_rep: Option<u32>,
5161 sync_call_set: TableId<WaitableSet>,
5163 old_do_not_suspend: Option<bool>,
5166}
5167
5168impl GuestThread {
5169 fn from_instance(
5172 state: Pin<&mut ComponentInstance>,
5173 caller_instance: RuntimeComponentInstanceIndex,
5174 guest_thread: u32,
5175 ) -> Result<TableId<Self>> {
5176 let rep = state.instance_states().0[caller_instance]
5177 .thread_handle_table()
5178 .guest_thread_rep(guest_thread)?;
5179 Ok(TableId::new(rep))
5180 }
5181
5182 fn new_implicit(state: &mut ConcurrentState, parent_task: TableId<GuestTask>) -> Result<Self> {
5183 let sync_call_set = state.push(WaitableSet {
5184 is_sync_call_set: true,
5185 ..WaitableSet::default()
5186 })?;
5187 Ok(Self {
5188 context: [0; NUM_COMPONENT_CONTEXT_SLOTS],
5189 parent_task,
5190 wake_on_cancel: WakeOnCancel::None,
5191 state: GuestThreadState::NotStartedImplicit,
5192 instance_rep: None,
5193 sync_call_set,
5194 old_do_not_suspend: None,
5195 })
5196 }
5197
5198 fn new_explicit(
5199 state: &mut ConcurrentState,
5200 parent_task: TableId<GuestTask>,
5201 start_func: Box<
5202 dyn FnOnce(&mut dyn VMStore, QualifiedThreadId) -> Result<()> + Send + Sync,
5203 >,
5204 ) -> Result<Self> {
5205 let sync_call_set = state.push(WaitableSet {
5206 is_sync_call_set: true,
5207 ..WaitableSet::default()
5208 })?;
5209 Ok(Self {
5210 context: [0; NUM_COMPONENT_CONTEXT_SLOTS],
5211 parent_task,
5212 wake_on_cancel: WakeOnCancel::None,
5213 state: GuestThreadState::NotStartedExplicit(start_func),
5214 instance_rep: None,
5215 sync_call_set,
5216 old_do_not_suspend: None,
5217 })
5218 }
5219}
5220
5221impl TableDebug for GuestThread {
5222 fn type_name() -> &'static str {
5223 "GuestThread"
5224 }
5225}
5226
5227enum SyncResult {
5228 NotProduced,
5229 Produced(Option<ValRaw>),
5230 Taken,
5231}
5232
5233impl SyncResult {
5234 fn take(&mut self) -> Result<Option<Option<ValRaw>>> {
5235 Ok(match mem::replace(self, SyncResult::Taken) {
5236 SyncResult::NotProduced => None,
5237 SyncResult::Produced(val) => Some(val),
5238 SyncResult::Taken => {
5239 bail_bug!("attempted to take a synchronous result that was already taken")
5240 }
5241 })
5242 }
5243}
5244
5245#[derive(Debug)]
5246enum HostFutureState {
5247 NotApplicable,
5248 Live,
5249 Dropped,
5250}
5251
5252pub(crate) struct GuestTask {
5254 common: WaitableCommon,
5256 lower_params: Option<RawLower>,
5258 lift_result: Option<LiftResult>,
5260 result: Option<LiftedResult>,
5263 callback: Option<CallbackFn>,
5266 caller: Caller,
5268 call_context: CallContext,
5273 sync_result: SyncResult,
5276 cancel_request_delivered: bool,
5280 starting_sent: bool,
5283 instance: RuntimeInstance,
5290 event: Option<Event>,
5293 exited: bool,
5295 threads: HashSet<TableId<GuestThread>>,
5297 host_future_state: HostFutureState,
5300 async_typed: bool,
5303 async_lifted: bool,
5306
5307 decremented_interesting_task_count: bool,
5308
5309 group: TaskGroupId,
5310}
5311
5312impl GuestTask {
5313 fn already_lowered_parameters(&self) -> bool {
5314 self.lower_params.is_none()
5316 }
5317
5318 fn returned_or_cancelled(&self) -> bool {
5319 self.lift_result.is_none()
5321 }
5322
5323 fn ready_to_delete(&self) -> bool {
5324 let threads_completed = self.threads.is_empty();
5325 let has_sync_result = matches!(self.sync_result, SyncResult::Produced(_));
5326 let pending_completion_event = matches!(
5327 self.common.event,
5328 Some(Event::Subtask {
5329 status: Status::Returned | Status::ReturnCancelled
5330 })
5331 );
5332 let ready = threads_completed
5333 && !has_sync_result
5334 && !pending_completion_event
5335 && !matches!(self.host_future_state, HostFutureState::Live);
5336 log::trace!(
5337 "ready to delete? {ready} (threads_completed: {}, has_sync_result: {}, pending_completion_event: {}, host_future_state: {:?})",
5338 threads_completed,
5339 has_sync_result,
5340 pending_completion_event,
5341 self.host_future_state
5342 );
5343 ready
5344 }
5345
5346 fn new(
5347 state: &mut ConcurrentState,
5348 lower_params: RawLower,
5349 lift_result: LiftResult,
5350 caller: Caller,
5351 callback: Option<CallbackFn>,
5352 instance: RuntimeInstance,
5353 async_typed: bool,
5354 async_lifted: bool,
5355 ) -> Result<QualifiedThreadId> {
5356 let host_future_state = match &caller {
5357 Caller::Guest { .. } => HostFutureState::NotApplicable,
5358 Caller::Host {
5359 host_future_present,
5360 ..
5361 } => {
5362 if *host_future_present {
5363 HostFutureState::Live
5364 } else {
5365 HostFutureState::NotApplicable
5366 }
5367 }
5368 };
5369
5370 let group = match caller {
5371 Caller::Guest { thread } => {
5372 let group = state.get_mut(thread.task)?.group;
5373 state.increment_group_ref_count(group)?;
5374 group
5375 }
5376 Caller::Host { .. } => state.make_task_group()?,
5377 };
5378
5379 let task = state.push(Self {
5380 common: WaitableCommon::default(),
5381 lower_params: Some(lower_params),
5382 lift_result: Some(lift_result),
5383 result: None,
5384 callback,
5385 caller,
5386 call_context: CallContext::default(),
5387 sync_result: SyncResult::NotProduced,
5388 cancel_request_delivered: false,
5389 starting_sent: false,
5390 instance,
5391 event: None,
5392 exited: false,
5393 threads: HashSet::new(),
5394 host_future_state,
5395 async_typed,
5396 async_lifted,
5397 decremented_interesting_task_count: false,
5398 group,
5399 })?;
5400 let new_thread = GuestThread::new_implicit(state, task)?;
5401 let thread = state.push(new_thread)?;
5402 state.get_mut(task)?.threads.insert(thread);
5403 state.interesting_tasks += 1;
5404 let thread = QualifiedThreadId { task, thread };
5405 log::trace!("new implicit thread {thread:?} for instance {instance:?}");
5406 Ok(thread)
5407 }
5408}
5409
5410impl TableDebug for GuestTask {
5411 fn type_name() -> &'static str {
5412 "GuestTask"
5413 }
5414}
5415
5416#[derive(Default)]
5418struct WaitableCommon {
5419 event: Option<Event>,
5421 set: Option<TableId<WaitableSet>>,
5423 handle: Option<u32>,
5425}
5426
5427#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
5429enum Waitable {
5430 Host(TableId<HostTask>),
5432 Guest(TableId<GuestTask>),
5434 Transmit(TableId<TransmitHandle>),
5436}
5437
5438impl Waitable {
5439 fn from_instance(
5442 state: Pin<&mut ComponentInstance>,
5443 caller_instance: RuntimeComponentInstanceIndex,
5444 waitable: u32,
5445 ) -> Result<Self> {
5446 use crate::runtime::vm::component::Waitable;
5447
5448 let (waitable, kind) = state.instance_states().0[caller_instance]
5449 .handle_table()
5450 .waitable_rep(waitable)?;
5451
5452 Ok(match kind {
5453 Waitable::Subtask { is_host: true } => Self::Host(TableId::new(waitable)),
5454 Waitable::Subtask { is_host: false } => Self::Guest(TableId::new(waitable)),
5455 Waitable::Stream | Waitable::Future => Self::Transmit(TableId::new(waitable)),
5456 })
5457 }
5458
5459 fn rep(&self) -> u32 {
5461 match self {
5462 Self::Host(id) => id.rep(),
5463 Self::Guest(id) => id.rep(),
5464 Self::Transmit(id) => id.rep(),
5465 }
5466 }
5467
5468 fn join(&self, state: &mut ConcurrentState, set: Option<TableId<WaitableSet>>) -> Result<()> {
5472 log::trace!("waitable {self:?} join set {set:?}");
5473
5474 let old = mem::replace(&mut self.common(state)?.set, set);
5475
5476 if let Some(old) = old {
5477 match *self {
5478 Waitable::Host(id) => state.remove_child(id, old),
5479 Waitable::Guest(id) => state.remove_child(id, old),
5480 Waitable::Transmit(id) => state.remove_child(id, old),
5481 }?;
5482
5483 state.get_mut(old)?.ready.remove(self);
5484 }
5485
5486 if let Some(set) = set {
5487 match *self {
5488 Waitable::Host(id) => state.add_child(id, set),
5489 Waitable::Guest(id) => state.add_child(id, set),
5490 Waitable::Transmit(id) => state.add_child(id, set),
5491 }?;
5492
5493 if self.common(state)?.event.is_some() {
5494 self.mark_ready(state)?;
5495 }
5496 }
5497
5498 Ok(())
5499 }
5500
5501 fn common<'a>(&self, state: &'a mut ConcurrentState) -> Result<&'a mut WaitableCommon> {
5503 Ok(match self {
5504 Self::Host(id) => &mut state.get_mut(*id)?.common,
5505 Self::Guest(id) => &mut state.get_mut(*id)?.common,
5506 Self::Transmit(id) => &mut state.get_mut(*id)?.common,
5507 })
5508 }
5509
5510 fn trap_if_in_waitable_set(&self, state: &mut ConcurrentState) -> Result<()> {
5516 if self.common(state)?.set.is_some() {
5517 bail!(Trap::WaitableSyncAndAsync);
5518 }
5519 Ok(())
5520 }
5521
5522 fn set_event(&self, state: &mut ConcurrentState, event: Option<Event>) -> Result<()> {
5526 log::trace!("set event for {self:?}: {event:?}");
5527 self.common(state)?.event = event;
5528 self.mark_ready(state)
5529 }
5530
5531 fn take_event(&self, state: &mut ConcurrentState) -> Result<Option<Event>> {
5533 let common = self.common(state)?;
5534 let event = common.event.take();
5535 if let Some(set) = self.common(state)?.set {
5536 state.get_mut(set)?.ready.remove(self);
5537 }
5538
5539 Ok(event)
5540 }
5541
5542 fn mark_ready(&self, state: &mut ConcurrentState) -> Result<()> {
5546 if let Some(set) = self.common(state)?.set {
5547 let set_state = state.get_mut(set)?;
5548 set_state.ready.insert(*self);
5549
5550 if let Some((thread, mode)) = set_state.waiting.pop_first() {
5551 let wake_on_cancel = state.get_mut(thread.thread)?.wake_on_cancel.take();
5552 assert!(wake_on_cancel.is_none() || wake_on_cancel == WakeOnCancel::Waiting(set));
5553
5554 let item = match mode {
5555 WaitMode::Fiber(fiber) => Some(WorkItem::ResumeFiber {
5556 instance: state.get_mut(thread.task)?.instance,
5557 thread,
5558 fiber,
5559 }),
5560 WaitMode::Callback(instance) => Some(WorkItem::GuestCall {
5561 instance: state.get_mut(thread.task)?.instance,
5562 call: GuestCall {
5563 thread,
5564 kind: GuestCallKind::DeliverEvent {
5565 instance,
5566 set: Some(set),
5567 },
5568 },
5569 }),
5570 };
5571
5572 if let Some(item) = item {
5573 state.push_high_priority(item);
5574 }
5575 }
5576 }
5577 Ok(())
5578 }
5579
5580 fn delete_from(&self, store: &mut StoreOpaque) -> Result<()> {
5582 match self {
5583 Self::Host(task) => {
5584 log::trace!("delete host task {task:?}");
5585 let state = store.concurrent_state_mut()?;
5586 let task = state.delete(*task)?;
5587
5588 state.decrement_group_ref_count(task.group)?;
5589 }
5590 Self::Guest(task) => {
5591 log::trace!("delete guest task {task:?}");
5592 let state = store.concurrent_state_mut()?;
5593 let task = state.delete(*task)?;
5594
5595 state.decrement_group_ref_count(task.group)?;
5596
5597 debug_assert!(task.decremented_interesting_task_count);
5604 }
5605 Self::Transmit(task) => {
5606 store.concurrent_state_mut()?.delete(*task)?;
5607 }
5608 }
5609
5610 Ok(())
5611 }
5612}
5613
5614impl fmt::Debug for Waitable {
5615 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
5616 match self {
5617 Self::Host(id) => write!(f, "{id:?}"),
5618 Self::Guest(id) => write!(f, "{id:?}"),
5619 Self::Transmit(id) => write!(f, "{id:?}"),
5620 }
5621 }
5622}
5623
5624#[derive(Default)]
5626struct WaitableSet {
5627 ready: BTreeSet<Waitable>,
5629 waiting: BTreeMap<QualifiedThreadId, WaitMode>,
5631 is_sync_call_set: bool,
5634}
5635
5636impl TableDebug for WaitableSet {
5637 fn type_name() -> &'static str {
5638 "WaitableSet"
5639 }
5640}
5641
5642type RawLower =
5644 Box<dyn FnOnce(&mut dyn VMStore, &mut [MaybeUninit<ValRaw>]) -> Result<()> + Send + Sync>;
5645
5646type RawLift = Box<
5648 dyn FnOnce(&mut dyn VMStore, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>> + Send + Sync,
5649>;
5650
5651type LiftedResult = Box<dyn Any + Send + Sync>;
5655
5656struct DummyResult;
5659
5660#[derive(Default)]
5662pub struct ConcurrentInstanceState {
5663 backpressure: u16,
5665 do_not_enter: bool,
5667 do_not_suspend: bool,
5670 pending: BTreeMap<QualifiedThreadId, GuestCallKind>,
5673}
5674
5675impl ConcurrentInstanceState {
5676 pub fn pending_is_empty(&self) -> bool {
5677 self.pending.is_empty()
5678 }
5679}
5680
5681#[derive(Debug, Copy, Clone)]
5682pub(crate) enum CurrentThread {
5683 Guest(QualifiedThreadId),
5686 Host(TableId<HostTask>),
5688 DeferredHost(QualifiedThreadId),
5691 None,
5694}
5695
5696impl CurrentThread {
5697 fn guest(&self) -> Option<&QualifiedThreadId> {
5698 match self {
5699 Self::Guest(id) => Some(id),
5700 _ => None,
5701 }
5702 }
5703
5704 fn guest_task(&self) -> Option<TableId<GuestTask>> {
5705 match self {
5706 Self::Guest(id) => Some(id.task),
5707 _ => None,
5708 }
5709 }
5710
5711 fn is_none(&self) -> bool {
5712 matches!(self, Self::None)
5713 }
5714}
5715
5716impl From<QualifiedThreadId> for CurrentThread {
5717 fn from(id: QualifiedThreadId) -> Self {
5718 Self::Guest(id)
5719 }
5720}
5721
5722impl From<TableId<HostTask>> for CurrentThread {
5723 fn from(id: TableId<HostTask>) -> Self {
5724 Self::Host(id)
5725 }
5726}
5727
5728enum Priority {
5729 Switch,
5730 High,
5731 Low,
5732}
5733
5734pub struct ConcurrentState {
5736 unforced_current_thread: CurrentThread,
5742
5743 deferred_host_call_context: Option<CallContext>,
5749
5750 futures: AlwaysMut<Option<FuturesUnordered<HostTaskFuture>>>,
5755 table: AlwaysMut<ResourceTable>,
5757 switch_item: Option<WorkItem>,
5765 next_switch_item: Option<WorkItem>,
5771 high_priority: VecDeque<WorkItem>,
5773 low_priority: VecDeque<WorkItem>,
5775 suspend_reason: Option<SuspendReason>,
5779 worker: Option<StoreFiber<'static>>,
5783 worker_item: Option<WorkerItem>,
5785
5786 global_error_context_ref_counts:
5799 BTreeMap<TypeComponentGlobalErrorContextTableIndex, GlobalErrorContextRefCount>,
5800
5801 interesting_tasks: usize,
5814
5815 interesting_tasks_empty_waker: Option<Waker>,
5819
5820 ready_for_concurrent_call_waker: Option<Waker>,
5825
5826 event_loop_running: bool,
5828
5829 #[cfg(feature = "task-group-hook")]
5831 task_group_hook: Option<Box<dyn TaskGroupHook>>,
5832}
5833
5834impl Default for ConcurrentState {
5835 fn default() -> Self {
5836 Self {
5837 unforced_current_thread: CurrentThread::None,
5838 deferred_host_call_context: None,
5839 table: AlwaysMut::new(ResourceTable::new()),
5840 futures: AlwaysMut::new(Some(FuturesUnordered::new())),
5841 switch_item: None,
5842 next_switch_item: None,
5843 high_priority: VecDeque::new(),
5844 low_priority: VecDeque::new(),
5845 suspend_reason: None,
5846 worker: None,
5847 worker_item: None,
5848 global_error_context_ref_counts: BTreeMap::new(),
5849 interesting_tasks: 0,
5850 interesting_tasks_empty_waker: None,
5851 ready_for_concurrent_call_waker: None,
5852 event_loop_running: false,
5853 #[cfg(feature = "task-group-hook")]
5854 task_group_hook: None,
5855 }
5856 }
5857}
5858
5859impl ConcurrentState {
5860 pub(crate) fn take_fibers_and_futures(
5877 &mut self,
5878 fibers: &mut Vec<StoreFiber<'static>>,
5879 futures: &mut Vec<FuturesUnordered<HostTaskFuture>>,
5880 ) {
5881 let mut items = Vec::new();
5882 for (_, entry) in self.table.get_mut().iter_mut() {
5883 if let Some(set) = entry.downcast_mut::<WaitableSet>() {
5884 for mode in mem::take(&mut set.waiting).into_values() {
5885 match mode {
5886 WaitMode::Fiber(fiber) => {
5887 fibers.push(fiber);
5888 }
5889 WaitMode::Callback(_) => {}
5890 }
5891 }
5892 } else if let Some(thread) = entry.downcast_mut::<GuestThread>() {
5893 if let GuestThreadState::Suspended(fiber) | GuestThreadState::Ready { fiber, .. } =
5894 mem::replace(&mut thread.state, GuestThreadState::Completed)
5895 {
5896 fibers.push(fiber);
5897 }
5898 } else if let Some(item) = entry.downcast_mut::<Option<WorkItem>>() {
5899 if let Some(item) = item.take() {
5900 items.push(item);
5901 }
5902 }
5903 }
5904
5905 if let Some(fiber) = self.worker.take() {
5906 fibers.push(fiber);
5907 }
5908
5909 let mut handle_item = |item| match item {
5910 WorkItem::ResumeFiber { fiber, .. } => {
5911 fibers.push(fiber);
5912 }
5913 WorkItem::PushFuture(future) => {
5914 self.futures
5915 .get_mut()
5916 .as_mut()
5917 .unwrap()
5918 .push(future.into_inner());
5919 }
5920 WorkItem::ResumeThread { .. }
5921 | WorkItem::GuestCall { .. }
5922 | WorkItem::WorkerFunction(_) => {}
5923 };
5924
5925 for item in items {
5926 handle_item(item);
5927 }
5928 if let Some(item) = self.switch_item.take() {
5929 handle_item(item);
5930 }
5931 if let Some(item) = self.next_switch_item.take() {
5932 handle_item(item);
5933 }
5934 for item in mem::take(&mut self.high_priority) {
5935 handle_item(item);
5936 }
5937 for item in mem::take(&mut self.low_priority) {
5938 handle_item(item);
5939 }
5940
5941 if let Some(them) = self.futures.get_mut().take() {
5942 futures.push(them);
5943 }
5944 }
5945
5946 #[cfg(feature = "gc")]
5947 pub(crate) fn trace_fiber_roots(
5948 &mut self,
5949 modules: &ModuleRegistry,
5950 unwind: &dyn Unwind,
5951 gc_roots_list: &mut GcRootsList,
5952 ) {
5953 let ConcurrentState {
5954 table,
5955 worker,
5956 switch_item,
5957 next_switch_item,
5958 high_priority,
5959 low_priority,
5960
5961 futures: _,
5965
5966 worker_item: _,
5968 unforced_current_thread: _,
5969 deferred_host_call_context: _,
5970 suspend_reason: _,
5971 global_error_context_ref_counts: _,
5972 interesting_tasks: _,
5973 interesting_tasks_empty_waker: _,
5974 ready_for_concurrent_call_waker: _,
5975 event_loop_running: _,
5976 #[cfg(feature = "task-group-hook")]
5977 task_group_hook: _,
5978 } = self;
5979
5980 for (_, entry) in table.get_mut().iter_mut() {
5981 if let Some(set) = entry.downcast_mut::<WaitableSet>() {
5982 for mode in set.waiting.values_mut() {
5983 match mode {
5984 WaitMode::Fiber(fiber) => {
5985 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
5986 }
5987 WaitMode::Callback(_) => {}
5988 }
5989 }
5990 } else if let Some(thread) = entry.downcast_mut::<GuestThread>() {
5991 if let GuestThreadState::Suspended(fiber) | GuestThreadState::Ready { fiber, .. } =
5992 &mut thread.state
5993 {
5994 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
5995 }
5996 } else if let Some(Some(WorkItem::ResumeFiber { fiber, .. })) =
5997 entry.downcast_mut::<Option<WorkItem>>()
5998 {
5999 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6000 }
6001 }
6002
6003 if let Some(fiber) = worker {
6004 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6005 }
6006
6007 let mut handle_item = |item: &mut WorkItem| match item {
6008 WorkItem::ResumeFiber { fiber, .. } => {
6009 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6010 }
6011 WorkItem::PushFuture(_future) => {
6012 }
6015 WorkItem::ResumeThread { .. }
6016 | WorkItem::GuestCall { .. }
6017 | WorkItem::WorkerFunction(_) => {}
6018 };
6019
6020 if let Some(item) = switch_item {
6021 handle_item(item);
6022 }
6023 if let Some(item) = next_switch_item {
6024 handle_item(item);
6025 }
6026 for item in high_priority {
6027 handle_item(item);
6028 }
6029 for item in low_priority {
6030 handle_item(item);
6031 }
6032 }
6033
6034 fn push<V: Send + Sync + 'static>(
6035 &mut self,
6036 value: V,
6037 ) -> Result<TableId<V>, ResourceTableError> {
6038 self.table.get_mut().push(value).map(TableId::from)
6039 }
6040
6041 fn get_mut<V: 'static>(&mut self, id: TableId<V>) -> Result<&mut V, ResourceTableError> {
6042 self.table.get_mut().get_mut(&Resource::from(id))
6043 }
6044
6045 pub fn add_child<T: 'static, U: 'static>(
6046 &mut self,
6047 child: TableId<T>,
6048 parent: TableId<U>,
6049 ) -> Result<(), ResourceTableError> {
6050 self.table
6051 .get_mut()
6052 .add_child(Resource::from(child), Resource::from(parent))
6053 }
6054
6055 pub fn remove_child<T: 'static, U: 'static>(
6056 &mut self,
6057 child: TableId<T>,
6058 parent: TableId<U>,
6059 ) -> Result<(), ResourceTableError> {
6060 self.table
6061 .get_mut()
6062 .remove_child(Resource::from(child), Resource::from(parent))
6063 }
6064
6065 fn delete<V: 'static>(&mut self, id: TableId<V>) -> Result<V, ResourceTableError> {
6066 self.table.get_mut().delete(Resource::from(id))
6067 }
6068
6069 fn push_future(&mut self, future: HostTaskFuture) {
6070 self.push_high_priority(WorkItem::PushFuture(AlwaysMut::new(future)));
6077 }
6078
6079 fn set_switch_item(&mut self, item: WorkItem) -> Result<()> {
6080 log::trace!("set switch item: {item:?}");
6081
6082 if self.switch_item.is_some() {
6083 bail_bug!("switch item already set");
6084 }
6085
6086 self.switch_item = Some(item);
6087
6088 Ok(())
6089 }
6090
6091 fn take_next_switch_item(&mut self) -> Result<()> {
6092 if let Some(item) = self.next_switch_item.take() {
6093 self.set_switch_item(item)?;
6094 }
6095 Ok(())
6096 }
6097
6098 fn push_high_priority(&mut self, item: WorkItem) {
6099 log::trace!("push high priority: {item:?}");
6100 self.high_priority.push_front(item);
6101 }
6102
6103 fn push_low_priority(&mut self, item: WorkItem) {
6104 log::trace!("push low priority: {item:?}");
6105 self.low_priority.push_front(item);
6106 }
6107
6108 fn push_work_item(&mut self, item: WorkItem, priority: Priority) -> Result<()> {
6109 match priority {
6110 Priority::Switch => self.set_switch_item(item)?,
6111 Priority::High => self.push_high_priority(item),
6112 Priority::Low => self.push_low_priority(item),
6113 }
6114
6115 Ok(())
6116 }
6117
6118 fn promote_instance_local_thread_work_item(
6119 &mut self,
6120 current_instance: RuntimeInstance,
6121 ) -> Result<bool> {
6122 log::trace!("promote thread work items for {current_instance:?}");
6123
6124 self.promote_work_item_matching(|item: &WorkItem| {
6125 let result = match item {
6126 WorkItem::ResumeThread { instance, .. }
6127 | WorkItem::ResumeFiber { instance, .. }
6128 | WorkItem::GuestCall { instance, .. } => *instance == current_instance,
6129 _ => false,
6130 };
6131
6132 log::trace!("candidate {item:?}: {result}");
6133 result
6134 })
6135 }
6136
6137 fn promote_thread_work_item(&mut self, thread: QualifiedThreadId) -> Result<bool> {
6138 self.promote_work_item_matching(|item: &WorkItem| match item {
6139 WorkItem::ResumeThread {
6140 thread: item_thread,
6141 ..
6142 }
6143 | WorkItem::GuestCall {
6144 call:
6145 GuestCall {
6146 thread: item_thread,
6147 ..
6148 },
6149 ..
6150 } => *item_thread == thread,
6151 _ => false,
6152 })
6153 }
6154
6155 fn promote_work_item_matching<F>(&mut self, mut predicate: F) -> Result<bool>
6156 where
6157 F: FnMut(&WorkItem) -> bool,
6158 {
6159 for item in mem::take(&mut self.high_priority).into_iter().rev() {
6164 if self.switch_item.is_none() && predicate(&item) {
6165 self.set_switch_item(item)?;
6166 } else {
6167 self.push_high_priority(item);
6168 }
6169 }
6170
6171 if self.switch_item.is_none() {
6172 for item in mem::take(&mut self.low_priority).into_iter().rev() {
6173 if self.switch_item.is_none() && predicate(&item) {
6174 self.set_switch_item(item)?;
6175 } else {
6176 self.push_low_priority(item);
6177 }
6178 }
6179 }
6180
6181 Ok(self.switch_item.is_some())
6182 }
6183
6184 pub fn call_context(&mut self, task: Scope) -> Result<&mut CallContext> {
6187 match task {
6188 Scope::HostId(task) => {
6189 let task: TableId<HostTask> = TableId::new(task);
6190 Ok(&mut self.get_mut(task)?.call_context)
6191 }
6192 Scope::Id(task) => {
6193 let task: TableId<GuestTask> = TableId::new(task);
6194 Ok(&mut self.get_mut(task)?.call_context)
6195 }
6196 }
6197 }
6198
6199 pub(crate) fn deferred_host_call_context(&mut self) -> Option<&mut CallContext> {
6200 self.deferred_host_call_context.as_mut()
6201 }
6202
6203 fn futures_mut(&mut self) -> Result<&mut FuturesUnordered<HostTaskFuture>> {
6204 match self.futures.get_mut().as_mut() {
6205 Some(f) => Ok(f),
6206 None => bail_bug!("futures field of concurrent state is currently taken"),
6207 }
6208 }
6209
6210 pub(crate) fn table(&mut self) -> &mut ResourceTable {
6211 self.table.get_mut()
6212 }
6213
6214 fn debug_assert_deferred_host_invariant(&self) {
6215 debug_assert_eq!(
6216 self.deferred_host_call_context.is_some(),
6217 matches!(self.unforced_current_thread, CurrentThread::DeferredHost(_)),
6218 "a deferred host thread and call context must exist together",
6219 );
6220 }
6221
6222 fn materialize_host_task(&mut self) -> Result<CurrentThread> {
6223 self.debug_assert_deferred_host_invariant();
6224 let caller = match self.unforced_current_thread {
6225 CurrentThread::DeferredHost(caller) => caller,
6226 thread => return Ok(thread),
6227 };
6228
6229 let task = HostTask::new(self, HostTaskState::CalleeStarted, caller)?;
6231 let task = self.push(task)?;
6232 let call_context = self
6233 .deferred_host_call_context
6234 .take()
6235 .expect("deferred host call context should be present");
6236 self.get_mut(task)
6237 .expect("newly inserted host task should be present")
6238 .call_context = call_context;
6239 self.unforced_current_thread = CurrentThread::Host(task);
6240 self.debug_assert_deferred_host_invariant();
6241 log::trace!("new host task materialized {task:?}");
6242 Ok(CurrentThread::Host(task))
6243 }
6244
6245 fn materialize_current_host_task_id(&mut self) -> Result<Option<TableId<HostTask>>> {
6246 match self.materialize_host_task()? {
6247 CurrentThread::Host(id) => Ok(Some(id)),
6248 CurrentThread::None => Ok(None),
6249 CurrentThread::Guest(_) => {
6250 bail_bug!("tried to materialize a host task id from a guest thread")
6251 }
6252 CurrentThread::DeferredHost(_) => {
6253 bail_bug!(
6254 "current thread is a deferred host thread which should have been materialized"
6255 )
6256 }
6257 }
6258 }
6259
6260 pub(crate) fn materialize_current_scope(&mut self) -> Result<Scope> {
6261 match self.materialize_host_task()? {
6262 CurrentThread::Host(id) => Ok(Scope::HostId(id.rep())),
6263 _ => bail_bug!("current scope is not a deferred host scope"),
6264 }
6265 }
6266}
6267
6268fn for_any_lower<
6271 F: FnOnce(&mut dyn VMStore, &mut [MaybeUninit<ValRaw>]) -> Result<()> + Send + Sync,
6272>(
6273 fun: F,
6274) -> F {
6275 fun
6276}
6277
6278fn for_any_lift<
6280 F: FnOnce(&mut dyn VMStore, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>> + Send + Sync,
6281>(
6282 fun: F,
6283) -> F {
6284 fun
6285}
6286
6287fn check_ambient_store(id: StoreId) {
6288 let message = "\
6289 `Future`s which depend on asynchronous component tasks, streams, or \
6290 futures to complete may only be polled from the event loop of the \
6291 store to which they belong. Please use \
6292 `StoreContextMut::{run_concurrent,spawn}` to poll or await them.\
6293 ";
6294 tls::try_get(|store| {
6295 let matched = match store {
6296 tls::TryGet::Some(store) => store.id() == id,
6297 tls::TryGet::Taken | tls::TryGet::None => false,
6298 };
6299
6300 if !matched {
6301 panic!("{message}")
6302 }
6303 });
6304}
6305
6306fn unpack_callback_code(code: u32) -> (u32, u32) {
6307 (code & 0xF, code >> 4)
6308}
6309
6310struct WaitableCheckParams {
6314 set: TableId<WaitableSet>,
6315 options: OptionsIndex,
6316 payload: u32,
6317}
6318
6319enum WaitableCheck {
6322 Wait,
6323 Poll,
6324}
6325
6326pub(crate) struct PreparedCall<R> {
6328 handle: Func,
6330 thread: QualifiedThreadId,
6332 param_count: usize,
6334 rx: oneshot::Receiver<LiftedResult>,
6337 runtime_instance: RuntimeInstance,
6339 _phantom: PhantomData<R>,
6340}
6341
6342impl<R> PreparedCall<R> {
6343 pub(crate) fn task_id(&self) -> TaskId {
6345 TaskId {
6346 task: self.thread.task,
6347 runtime_instance: self.runtime_instance,
6348 }
6349 }
6350}
6351
6352pub(crate) struct TaskId {
6354 task: TableId<GuestTask>,
6355 runtime_instance: RuntimeInstance,
6356}
6357
6358impl TaskId {
6359 pub(crate) fn host_future_dropped(&self, store: &mut StoreOpaque) -> Result<()> {
6365 let task = store.concurrent_state_mut()?.get_mut(self.task)?;
6366 let delete = if !task.already_lowered_parameters() {
6367 store.cancel_guest_subtask_without_lowered_parameters(
6368 self.runtime_instance,
6369 self.task,
6370 )?;
6371 true
6372 } else {
6373 task.host_future_state = HostFutureState::Dropped;
6374 task.ready_to_delete()
6375 };
6376 if delete {
6377 Waitable::Guest(self.task).delete_from(store)?
6378 }
6379 Ok(())
6380 }
6381}
6382
6383pub(crate) fn prepare_call<T, R>(
6389 mut store: StoreContextMut<T>,
6390 handle: Func,
6391 param_count: usize,
6392 host_future_present: bool,
6393 lower_params: impl FnOnce(StoreContextMut<T>, &mut [MaybeUninit<ValRaw>]) -> Result<()>
6394 + Send
6395 + Sync
6396 + 'static,
6397 lift_result: impl FnOnce(&mut StoreOpaque, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>>
6398 + Send
6399 + Sync
6400 + 'static,
6401) -> Result<PreparedCall<R>> {
6402 if !store.0.may_enter() {
6403 bail!(Trap::CannotEnterComponent);
6404 }
6405
6406 let (options, _flags, ty, raw_options) = handle.abi_info(store.0);
6407
6408 let instance = handle.instance().id().get(store.0);
6409 let options = &instance.component().env_component().options[options];
6410 let ty = &instance.component().types()[ty];
6411 let async_typed = ty.async_;
6412 let async_lifted = raw_options.async_;
6413 let task_return_type = ty.results;
6414 let component_instance = raw_options.instance;
6415 let callback = options.callback.map(|i| instance.runtime_callback(i));
6416 let memory = options
6417 .memory()
6418 .map(|i| instance.runtime_memory(i))
6419 .map(SendSyncPtr::new);
6420 let string_encoding = options.string_encoding;
6421 let token = StoreToken::new(store.as_context_mut());
6422 let caller = store.0.materialize_host_task_id()?;
6423 let state = store.0.concurrent_state_mut()?;
6424
6425 let (tx, rx) = oneshot::channel();
6426
6427 let instance = handle.instance().runtime_instance(component_instance);
6428 let thread = GuestTask::new(
6429 state,
6430 Box::new(for_any_lower(move |store, params| {
6431 lower_params(token.as_context_mut(store), params)
6432 })),
6433 LiftResult {
6434 lift: Box::new(for_any_lift(move |store, result| {
6435 lift_result(store, result)
6436 })),
6437 ty: task_return_type,
6438 memory,
6439 string_encoding,
6440 },
6441 Caller::Host {
6442 tx: Some(tx),
6443 host_future_present,
6444 caller,
6445 },
6446 callback.map(|callback| {
6447 let callback = SendSyncPtr::new(callback);
6448 let instance = handle.instance();
6449 Box::new(move |store: &mut dyn VMStore, event, handle| {
6450 let store = token.as_context_mut(store);
6451 unsafe { instance.call_callback(store, callback, event, handle) }
6454 }) as CallbackFn
6455 }),
6456 instance,
6457 async_typed,
6458 async_lifted,
6459 )?;
6460
6461 Ok(PreparedCall {
6462 handle,
6463 thread,
6464 param_count,
6465 runtime_instance: instance,
6466 rx,
6467 _phantom: PhantomData,
6468 })
6469}
6470
6471pub(crate) struct StagedCall<R> {
6472 store: StoreId,
6473 rx: oneshot::Receiver<LiftedResult>,
6474 _marker: PhantomData<fn() -> R>,
6475 group: TaskGroupId,
6476}
6477
6478impl<R> StagedCall<R> {
6479 pub(crate) fn new<T: 'static>(
6486 mut store: StoreContextMut<T>,
6487 prepared: PreparedCall<R>,
6488 ) -> Result<StagedCall<R>> {
6489 let PreparedCall {
6490 handle,
6491 thread,
6492 param_count,
6493 rx,
6494 ..
6495 } = prepared;
6496
6497 stage_call0(store.as_context_mut(), handle, thread, param_count)?;
6498
6499 Ok(StagedCall {
6500 store: store.0.id(),
6501 rx,
6502 _marker: PhantomData,
6503 group: store.0.concurrent_state_mut()?.get_mut(thread.task)?.group,
6504 })
6505 }
6506}
6507
6508impl<R> Future for StagedCall<R>
6509where
6510 R: 'static,
6511{
6512 type Output = Result<R>;
6513
6514 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
6515 check_ambient_store(self.store);
6516 Pin::new(&mut self.rx).poll(cx).map(|result| match result {
6517 Ok(r) => match r.downcast() {
6518 Ok(r) => Ok(*r),
6519 Err(_) => bail_bug!("wrong type of value produced"),
6520 },
6521 Err(oneshot::Canceled) => bail_bug!("channel erroneously dropped"),
6522 })
6523 }
6524}
6525
6526fn stage_call0<T: 'static>(
6529 store: StoreContextMut<T>,
6530 handle: Func,
6531 guest_thread: QualifiedThreadId,
6532 param_count: usize,
6533) -> Result<()> {
6534 let (_options, _, _ty, raw_options) = handle.abi_info(store.0);
6535 let is_concurrent = raw_options.async_;
6536 let callback = raw_options.callback;
6537 let instance = handle.instance();
6538 let callee = handle.lifted_core_func(store.0);
6539 let post_return = raw_options
6540 .post_return
6541 .map(|i| instance.id().get(store.0).runtime_post_return(i));
6542 let callback = callback.map(|i| {
6543 let instance = instance.id().get(store.0);
6544 SendSyncPtr::new(instance.runtime_callback(i))
6545 });
6546
6547 log::trace!("queueing call {guest_thread:?}");
6548
6549 unsafe {
6553 instance.stage_call(
6554 store,
6555 guest_thread,
6556 SendSyncPtr::new(callee),
6557 param_count,
6558 1,
6559 is_concurrent,
6560 callback,
6561 post_return.map(SendSyncPtr::new),
6562 true,
6563 )
6564 }
6565}
6566
6567#[cfg(all(test, feature = "cranelift", feature = "wat"))]
6568mod tests {
6569 use super::*;
6570 use crate::component::{Component, Linker};
6571 use crate::store::AsStoreOpaque;
6572 use crate::{Config, Engine};
6573
6574 fn host_subtask(
6575 state: HostTaskState,
6576 event: Option<Event>,
6577 ) -> Result<(Store<()>, Instance, TableId<HostTask>, u32)> {
6578 let mut config = Config::new();
6579 config.wasm_component_model_async(true);
6580 let engine = Engine::new(&config)?;
6581 let component = Component::new(&engine, "(component)")?;
6582 let mut store = Store::new(&engine, ());
6583 let instance = Linker::new(&engine).instantiate(&mut store, &component)?;
6584 let store_opaque = store.as_store_opaque();
6585 let concurrent_state = store_opaque.concurrent_state_mut()?;
6586 let group = concurrent_state.make_task_group()?;
6589 let task = concurrent_state.push(HostTask {
6590 common: WaitableCommon::default(),
6591 call_context: CallContext::default(),
6592 state,
6593 group,
6594 })?;
6595 let handle = store_opaque
6596 .instance_state(instance.runtime_instance(RuntimeComponentInstanceIndex::from_u32(0)))
6597 .handle_table()
6598 .subtask_insert_host(task.rep())?;
6599 let common = &mut store_opaque.concurrent_state_mut()?.get_mut(task)?.common;
6600 common.handle = Some(handle);
6601 common.event = event;
6602 Ok((store, instance, task, handle))
6603 }
6604
6605 #[test]
6606 fn host_subtask_drop_during_cancellation() -> Result<()> {
6607 for abort_completed in [false, true] {
6608 let (handle, future) = JoinHandle::run(future::pending::<()>());
6609 let mut future = pin!(future);
6610 let (mut store, instance, task, handle) =
6611 host_subtask(HostTaskState::CalleeRunning(handle), None)?;
6612 let store = store.as_store_opaque();
6613 let caller = RuntimeComponentInstanceIndex::from_u32(0);
6614 assert_eq!(
6615 instance.subtask_cancel(store, caller, true, handle)?,
6616 BLOCKED
6617 );
6618 if abort_completed {
6619 assert!(matches!(
6622 future
6623 .as_mut()
6624 .poll(&mut Context::from_waker(Waker::noop())),
6625 Poll::Ready(None),
6626 ));
6627 }
6628 for async_ in [false, true] {
6629 let err = instance
6630 .subtask_cancel(store, caller, async_, handle)
6631 .unwrap_err();
6632 assert_eq!(err.downcast::<Trap>()?, Trap::SubtaskCancelAfterTerminal);
6633 }
6634 let err = instance.subtask_drop(store, caller, handle).unwrap_err();
6635 assert_eq!(err.downcast::<Trap>()?, Trap::SubtaskDropNotResolved);
6636 assert!(store.concurrent_state_mut()?.get_mut(task).is_ok());
6637 }
6638 Ok(())
6639 }
6640
6641 #[test]
6642 fn host_subtask_cancel_after_completion() -> Result<()> {
6643 for async_ in [false, true] {
6644 let (mut store, instance, task, handle) = host_subtask(
6645 HostTaskState::CalleeDone { cancelled: false },
6646 Some(Event::Subtask {
6647 status: Status::Returned,
6648 }),
6649 )?;
6650 let store = store.as_store_opaque();
6651 let caller = RuntimeComponentInstanceIndex::from_u32(0);
6652 assert_eq!(
6653 instance.subtask_cancel(store, caller, async_, handle)?,
6654 Status::Returned as u32,
6655 );
6656 let err = instance
6657 .subtask_cancel(store, caller, async_, handle)
6658 .unwrap_err();
6659 assert_eq!(err.downcast::<Trap>()?, Trap::SubtaskCancelAfterTerminal);
6660 instance.subtask_drop(store, caller, handle)?;
6661 assert!(store.concurrent_state_mut()?.get_mut(task).is_err());
6662 }
6663 Ok(())
6664 }
6665
6666 #[test]
6667 fn host_subtask_drop_requires_terminal_event_delivery() -> Result<()> {
6668 for (cancelled, status) in [
6669 (false, Status::Returned),
6670 (true, Status::Returned),
6671 (true, Status::ReturnCancelled),
6672 ] {
6673 for delivered in [false, true] {
6674 let event = if delivered {
6675 None
6676 } else {
6677 Some(Event::Subtask { status })
6678 };
6679 let (mut store, instance, task, handle) =
6680 host_subtask(HostTaskState::CalleeDone { cancelled }, event)?;
6681 let store = store.as_store_opaque();
6682 let result = instance.subtask_drop(
6683 store,
6684 RuntimeComponentInstanceIndex::from_u32(0),
6685 handle,
6686 );
6687 if delivered {
6688 result?;
6689 assert!(store.concurrent_state_mut()?.get_mut(task).is_err());
6690 } else {
6691 let err = result.unwrap_err();
6692 assert_eq!(err.downcast::<Trap>()?, Trap::SubtaskDropNotResolved);
6693 assert!(store.concurrent_state_mut()?.get_mut(task).is_ok());
6694 }
6695 }
6696 }
6697 Ok(())
6698 }
6699}