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::{
68 AlwaysMut, SendSyncPtr, UncaughtException, VMFuncRef, VMLazyThread, VMMemoryDefinition, VMStore,
69};
70use crate::{
71 AsContext, AsContextMut, FuncType, Result, StoreContext, StoreContextMut, ValRaw, ValType, bail,
72};
73use crate::{Instance as ModuleInstance, bail_bug};
74use alloc::borrow::ToOwned;
75use alloc::collections::{BTreeMap, BTreeSet, VecDeque};
76use core::any::Any;
77use core::cell::UnsafeCell;
78use core::fmt;
79use core::future;
80use core::future::Future;
81use core::marker::PhantomData;
82use core::mem::{self, ManuallyDrop, MaybeUninit};
83use core::ops::DerefMut;
84use core::pin::{Pin, pin};
85use core::ptr::{self, NonNull};
86use core::task::{Context, Poll, Waker};
87use futures::channel::oneshot;
88use futures::stream::{FuturesUnordered, StreamExt};
89use futures_and_streams::{FlatAbi, ReturnCode, TransmitHandle, TransmitIndex};
90use table::{TableDebug, TableId};
91use wasmtime_environ::component::{
92 CanonicalAbiInfo, CanonicalOptions, CanonicalOptionsDataModel, MAX_FLAT_PARAMS,
93 MAX_FLAT_RESULTS, OptionsIndex, PREPARE_ASYNC_NO_RESULT, PREPARE_ASYNC_WITH_RESULT,
94 RuntimeComponentInstanceIndex, RuntimeTableIndex, StringEncoding,
95 TypeComponentGlobalErrorContextTableIndex, TypeComponentLocalErrorContextTableIndex,
96 TypeFuncIndex, TypeFutureTableIndex, TypeStreamTableIndex, TypeTupleIndex,
97};
98use wasmtime_environ::packed_option::ReservedValue;
99use wasmtime_environ::{NUM_COMPONENT_CONTEXT_SLOTS, Trap};
100#[cfg(feature = "gc")]
101use wasmtime_unwinder::Unwind;
102
103pub use abort::JoinHandle;
104pub use func::{FuncCallConcurrent, TypedFuncCallConcurrent};
105pub use future_stream_any::{FutureAny, StreamAny};
106pub use futures_and_streams::{
107 Destination, DirectDestination, DirectSource, ErrorContext, FutureConsumer, FutureProducer,
108 FutureReader, GuardedFutureReader, GuardedStreamReader, ReadBuffer, Source, StreamConsumer,
109 StreamProducer, StreamReader, StreamResult, VecBuffer, WriteBuffer,
110};
111pub(crate) use futures_and_streams::{ResourcePair, lower_error_context_to_index};
112#[cfg(feature = "task-group-hook")]
113pub use task_group_hook::TaskGroupHook;
114pub use task_group_hook::TaskGroupId;
115
116mod abort;
117mod error_contexts;
118mod func;
119mod future_stream_any;
120mod futures_and_streams;
121pub(crate) mod table;
122#[cfg(feature = "task-group-hook")]
123mod task_group_hook;
124#[cfg(not(feature = "task-group-hook"))]
125mod task_group_hook_disabled;
126#[cfg(not(feature = "task-group-hook"))]
127use task_group_hook_disabled as task_group_hook;
128pub(crate) mod tls;
129
130const BLOCKED: u32 = 0xffff_ffff;
133
134#[derive(Clone, Copy, Eq, PartialEq, Debug)]
136pub enum Status {
137 Starting = 0,
138 Started = 1,
139 Returned = 2,
140 StartCancelled = 3,
141 ReturnCancelled = 4,
142}
143
144impl Status {
145 pub fn pack(self, waitable: Option<u32>) -> u32 {
151 assert!(matches!(self, Status::Returned) == waitable.is_none());
152 let waitable = waitable.unwrap_or(0);
153 assert!(waitable < (1 << 28));
154 (waitable << 4) | (self as u32)
155 }
156}
157
158#[derive(Clone, Copy, Debug)]
161enum Event {
162 None,
163 Subtask {
164 status: Status,
165 },
166 StreamRead {
167 code: ReturnCode,
168 pending: Option<(TypeStreamTableIndex, u32)>,
169 },
170 StreamWrite {
171 code: ReturnCode,
172 pending: Option<(TypeStreamTableIndex, u32)>,
173 },
174 FutureRead {
175 code: ReturnCode,
176 pending: Option<(TypeFutureTableIndex, u32)>,
177 },
178 FutureWrite {
179 code: ReturnCode,
180 pending: Option<(TypeFutureTableIndex, u32)>,
181 },
182 Cancelled,
183}
184
185impl Event {
186 fn parts(self) -> (u32, u32) {
191 const EVENT_NONE: u32 = 0;
192 const EVENT_SUBTASK: u32 = 1;
193 const EVENT_STREAM_READ: u32 = 2;
194 const EVENT_STREAM_WRITE: u32 = 3;
195 const EVENT_FUTURE_READ: u32 = 4;
196 const EVENT_FUTURE_WRITE: u32 = 5;
197 const EVENT_CANCELLED: u32 = 6;
198 match self {
199 Event::None => (EVENT_NONE, 0),
200 Event::Cancelled => (EVENT_CANCELLED, 0),
201 Event::Subtask { status } => (EVENT_SUBTASK, status as u32),
202 Event::StreamRead { code, .. } => (EVENT_STREAM_READ, code.encode()),
203 Event::StreamWrite { code, .. } => (EVENT_STREAM_WRITE, code.encode()),
204 Event::FutureRead { code, .. } => (EVENT_FUTURE_READ, code.encode()),
205 Event::FutureWrite { code, .. } => (EVENT_FUTURE_WRITE, code.encode()),
206 }
207 }
208}
209
210mod callback_code {
212 pub const EXIT: u32 = 0;
213 pub const YIELD: u32 = 1;
214 pub const WAIT: u32 = 2;
215}
216
217const START_FLAG_ASYNC_CALLEE: u32 = wasmtime_environ::component::START_FLAG_ASYNC_CALLEE as u32;
221
222pub struct Access<'a, T: 'static, D: HasData + ?Sized = HasSelf<T>> {
228 store: StoreContextMut<'a, T>,
229 get_data: fn(&mut T) -> D::Data<'_>,
230}
231
232impl<'a, T, D> Access<'a, T, D>
233where
234 D: HasData + ?Sized,
235 T: 'static,
236{
237 pub fn new(store: StoreContextMut<'a, T>, get_data: fn(&mut T) -> D::Data<'_>) -> Self {
239 Self { store, get_data }
240 }
241
242 pub fn data_mut(&mut self) -> &mut T {
244 self.store.data_mut()
245 }
246
247 pub fn get(&mut self) -> D::Data<'_> {
249 (self.get_data)(self.data_mut())
250 }
251
252 pub fn spawn(&mut self, task: impl for<'fut> AccessorTask<'fut, T, D>) -> Result<JoinHandle>
256 where
257 T: 'static,
258 {
259 let accessor = Accessor {
260 get_data: self.get_data,
261 token: StoreToken::new(self.store.as_context_mut()),
262 };
263 self.store
264 .as_context_mut()
265 .spawn_with_accessor(accessor, task)
266 }
267
268 pub fn getter(&self) -> fn(&mut T) -> D::Data<'_> {
271 self.get_data
272 }
273}
274
275impl<'a, T, D> AsContext for Access<'a, T, D>
276where
277 D: HasData + ?Sized,
278 T: 'static,
279{
280 type Data = T;
281
282 fn as_context(&self) -> StoreContext<'_, T> {
283 self.store.as_context()
284 }
285}
286
287impl<'a, T, D> AsContextMut for Access<'a, T, D>
288where
289 D: HasData + ?Sized,
290 T: 'static,
291{
292 fn as_context_mut(&mut self) -> StoreContextMut<'_, T> {
293 self.store.as_context_mut()
294 }
295}
296
297pub struct Accessor<T: 'static, D = HasSelf<T>>
357where
358 D: HasData + ?Sized,
359{
360 token: StoreToken<T>,
361 get_data: fn(&mut T) -> D::Data<'_>,
362}
363
364pub trait AsAccessor {
381 type Data: 'static;
383
384 type AccessorData: HasData + ?Sized;
387
388 fn as_accessor(&self) -> &Accessor<Self::Data, Self::AccessorData>;
390}
391
392impl<T: AsAccessor + ?Sized> AsAccessor for &T {
393 type Data = T::Data;
394 type AccessorData = T::AccessorData;
395
396 fn as_accessor(&self) -> &Accessor<Self::Data, Self::AccessorData> {
397 T::as_accessor(self)
398 }
399}
400
401impl<T, D: HasData + ?Sized> AsAccessor for Accessor<T, D> {
402 type Data = T;
403 type AccessorData = D;
404
405 fn as_accessor(&self) -> &Accessor<T, D> {
406 self
407 }
408}
409
410const _: () = {
433 const fn assert<T: Send + Sync>() {}
434 assert::<Accessor<UnsafeCell<u32>>>();
435};
436
437impl<T> Accessor<T> {
438 pub(crate) fn new(token: StoreToken<T>) -> Self {
447 Self {
448 token,
449 get_data: |x| x,
450 }
451 }
452}
453
454impl<T, D> Accessor<T, D>
455where
456 D: HasData + ?Sized,
457{
458 pub fn with<R>(&self, fun: impl FnOnce(Access<'_, T, D>) -> R) -> R {
476 tls::get(|vmstore| {
477 fun(Access {
478 store: self.token.as_context_mut(vmstore),
479 get_data: self.get_data,
480 })
481 })
482 }
483
484 pub fn getter(&self) -> fn(&mut T) -> D::Data<'_> {
487 self.get_data
488 }
489
490 pub fn with_getter<D2: HasData>(
507 &self,
508 get_data: fn(&mut T) -> D2::Data<'_>,
509 ) -> Accessor<T, D2> {
510 Accessor {
511 token: self.token,
512 get_data,
513 }
514 }
515
516 pub fn spawn(&self, task: impl for<'fut> AccessorTask<'fut, T, D>) -> Result<JoinHandle>
532 where
533 T: 'static,
534 {
535 let accessor = self.clone_for_spawn();
536 self.with(|mut access| access.as_context_mut().spawn_with_accessor(accessor, task))
537 }
538
539 fn clone_for_spawn(&self) -> Self {
540 Self {
541 token: self.token,
542 get_data: self.get_data,
543 }
544 }
545
546 pub fn poll_no_interesting_tasks(&self, cx: &mut Context<'_>) -> Poll<()> {
582 self.with(|mut access| {
583 let store = access.as_context_mut().0;
584 let state = store.concurrent_state_mut_without_forcing_current_thread();
585 if state.interesting_tasks == 0 {
586 Poll::Ready(())
587 } else {
588 state.interesting_tasks_empty_waker = Some(cx.waker().clone());
589 Poll::Pending
590 }
591 })
592 }
593
594 pub fn poll_ready_for_concurrent_call(&self, func: Func, cx: &mut Context<'_>) -> Poll<()> {
611 self.with(|mut access| {
612 let store = access.as_context_mut().0;
613 let (_, _, _, raw_options) = func.abi_info(store);
614 let instance = func.instance().runtime_instance(raw_options.instance);
615 let state = store.instance_state(instance).concurrent_state();
616 if state.backpressure == 0 {
617 Poll::Ready(())
618 } else {
619 store
620 .concurrent_state_mut_without_forcing_current_thread()
621 .ready_for_concurrent_call_waker = Some(cx.waker().clone());
622 Poll::Pending
623 }
624 })
625 }
626}
627
628pub trait AccessorTask<'fut, T, D = HasSelf<T>>:
650 AsyncFnOnce(&Accessor<T, D>) -> Result<()> + Send + 'static
651where
652 D: HasData + ?Sized,
653{
654 fn run(self, accessor: &'fut Accessor<T, D>) -> impl Future<Output = Result<()>> + Send + 'fut;
656}
657
658impl<'fut, F, Fut, T, D> AccessorTask<'fut, T, D> for F
659where
660 T: 'static,
661 F: AsyncFnOnce(&Accessor<T, D>) -> Result<()>,
662 F: FnOnce(&'fut Accessor<T, D>) -> Fut + Send + 'static,
663 Fut: Future<Output = Result<()>> + Send + 'fut,
664 D: HasData,
665{
666 fn run(self, accessor: &'fut Accessor<T, D>) -> impl Future<Output = Result<()>> + Send + 'fut {
667 (self)(accessor)
668 }
669}
670
671enum CallerInfo {
674 Async {
676 params: Vec<ValRaw>,
677 has_result: bool,
678 },
679 Sync {
681 params: Vec<ValRaw>,
682 result_count: u32,
683 },
684}
685
686enum WaitMode {
688 Fiber(StoreFiber<'static>),
690 Callback(Instance),
693}
694
695impl fmt::Debug for WaitMode {
696 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
697 match self {
698 Self::Fiber(_) => f.debug_tuple("Fiber").finish(),
699 Self::Callback(instance) => f.debug_tuple("Callback").field(instance).finish(),
700 }
701 }
702}
703
704#[derive(Debug)]
706enum SuspendReason {
707 Waiting {
710 set: TableId<WaitableSet>,
711 thread: QualifiedThreadId,
712 },
713 YieldingToSubtask { thread: QualifiedThreadId },
716 NeedWork,
719 Yielding { thread: QualifiedThreadId },
722 ExplicitlySuspending { thread: QualifiedThreadId },
725}
726
727enum GuestCallKind {
729 DeliverEvent {
732 instance: Instance,
734 set: Option<TableId<WaitableSet>>,
739 },
740 StartImplicit(Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send + Sync>),
746 StartExplicit(Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send + Sync>),
747}
748
749impl fmt::Debug for GuestCallKind {
750 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
751 match self {
752 Self::DeliverEvent { instance, set } => f
753 .debug_struct("DeliverEvent")
754 .field("instance", instance)
755 .field("set", set)
756 .finish(),
757 Self::StartImplicit(_) => f.debug_tuple("StartImplicit").finish(),
758 Self::StartExplicit(_) => f.debug_tuple("StartExplicit").finish(),
759 }
760 }
761}
762
763#[derive(Copy, Clone, Debug)]
765pub enum SuspensionTarget {
766 Resume(u32),
767 Promote(u32),
768 None,
769}
770
771#[derive(Copy, Clone, Debug)]
773pub enum ResumeThread {
774 Promote,
775 Resume,
776 ResumeLater,
777}
778
779#[derive(Debug)]
781struct GuestCall {
782 thread: QualifiedThreadId,
783 kind: GuestCallKind,
784}
785
786impl GuestCall {
787 fn is_ready(&self, store: &mut StoreOpaque) -> Result<bool> {
797 let task = store.concurrent_state_mut()?.get_mut(self.thread.task)?;
798 let async_typed = task.async_typed;
799 let instance = task.instance;
800 let state = store.instance_state(instance).concurrent_state();
801
802 let ready = match &self.kind {
803 GuestCallKind::DeliverEvent { .. } => !state.do_not_enter,
804 GuestCallKind::StartImplicit(_) => {
805 !async_typed || !(state.do_not_enter || state.backpressure > 0)
806 }
807 GuestCallKind::StartExplicit(_) => true,
808 };
809 log::trace!(
810 "call {self:?} ready? {ready} (do_not_enter: {}; backpressure: {})",
811 state.do_not_enter,
812 state.backpressure
813 );
814 Ok(ready)
815 }
816}
817
818enum WorkerItem {
820 GuestCall(GuestCall),
821 Function(AlwaysMut<Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send>>),
822}
823
824enum WorkItem {
827 PushFuture(AlwaysMut<HostTaskFuture>),
829 ResumeFiber {
831 instance: RuntimeInstance,
832 thread: QualifiedThreadId,
833 fiber: StoreFiber<'static>,
834 },
835 ResumeThread {
837 instance: RuntimeInstance,
838 thread: QualifiedThreadId,
839 },
840 GuestCall {
842 instance: RuntimeInstance,
843 call: GuestCall,
844 },
845 WorkerFunction(AlwaysMut<Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send>>),
847}
848
849impl fmt::Debug for WorkItem {
850 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
851 match self {
852 Self::PushFuture(_) => f.debug_tuple("PushFuture").finish(),
853 Self::ResumeFiber {
854 instance, thread, ..
855 } => f
856 .debug_struct("ResumeFiber")
857 .field("instance", instance)
858 .field("thread", thread)
859 .finish(),
860 Self::ResumeThread { instance, thread } => f
861 .debug_struct("ResumeThread")
862 .field("instance", instance)
863 .field("thread", thread)
864 .finish(),
865 Self::GuestCall { instance, call } => f
866 .debug_struct("GuestCall")
867 .field("instance", instance)
868 .field("call", call)
869 .finish(),
870 Self::WorkerFunction(_) => f.debug_tuple("WorkerFunction").finish(),
871 }
872 }
873}
874
875#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
877pub(crate) enum WaitResult {
878 Cancelled,
879 Completed,
880}
881
882pub(crate) fn poll_and_block<R: Send + Sync + 'static>(
890 store: &mut dyn VMStore,
891 host_task: EnteredHostTask,
892 future: impl Future<Output = Result<R>> + Send + 'static,
893) -> Result<R> {
894 let mut future = Box::pin(future);
901 let poll = tls::set(store, || {
902 future
903 .as_mut()
904 .poll(&mut Context::from_waker(&Waker::noop()))
905 });
906
907 let caller = match host_task {
908 Some(caller) => caller,
909 None => bail_bug!("host task wasn't created but should have been"),
910 };
911
912 let task = match poll {
913 Poll::Ready(result) => return result,
915
916 Poll::Pending => {
921 let Some(task) = store.materialize_host_task_id()? else {
922 bail_bug!("current thread is not a host thread")
923 };
924
925 let future = Box::pin(async move {
928 let result = run_with_host_task_set(task, future).await??;
929 tls::get(move |store| {
930 let state = store.concurrent_state_mut()?;
931 let host_state = &mut state.get_mut(task)?.state;
932 assert!(matches!(host_state, HostTaskState::CalleeStarted));
933 *host_state = HostTaskState::CalleeFinished(Box::new(result));
934
935 Waitable::Host(task).set_event(
936 state,
937 Some(Event::Subtask {
938 status: Status::Returned,
939 }),
940 )?;
941
942 Ok(())
943 })
944 }) as HostTaskFuture;
945
946 let caller_instance = store.concurrent_state_mut()?.get_mut(caller.task)?.instance;
947 store.switch_or_trap_if_may_not_suspend(caller_instance)?;
948
949 let state = store.concurrent_state_mut()?;
950 state.push_future(future);
951
952 let set = state.get_mut(caller.thread)?.sync_call_set;
953 Waitable::Host(task).join(state, Some(set))?;
954
955 store.suspend(SuspendReason::Waiting {
956 set,
957 thread: caller,
958 })?;
959
960 Waitable::Host(task).join(store.concurrent_state_mut()?, None)?;
964 task
965 }
966 };
967
968 let host_state = &mut store.concurrent_state_mut()?.get_mut(task)?.state;
970 match mem::replace(host_state, HostTaskState::CalleeDone { cancelled: false }) {
971 HostTaskState::CalleeFinished(result) => Ok(match result.downcast() {
972 Ok(result) => *result,
973 Err(_) => bail_bug!("host task finished with wrong type of result"),
974 }),
975 _ => bail_bug!("unexpected host task state after completion"),
976 }
977}
978
979fn handle_guest_call(store: &mut dyn VMStore, call: GuestCall) -> Result<()> {
981 match call.kind {
982 GuestCallKind::DeliverEvent { instance, set } => {
983 if let Some(set) = set {
986 store.concurrent_state_mut()?.get_mut(set)?.stop_waiting()?;
987 }
988 let (event, waitable) = match instance.get_event(store, call.thread.task, set, true)? {
989 Some(pair) => pair,
990 None => match set {
991 Some(set) => {
996 log::trace!(
997 "event for {:?} on {set:?} no longer present; waiting again",
998 call.thread
999 );
1000 let state = store.concurrent_state_mut()?;
1001 return instance.wait_with_callback(state, call.thread, set);
1002 }
1003 None => (Event::None, None),
1008 },
1009 };
1010 let state = store.concurrent_state_mut()?;
1011 let task = state.get_mut(call.thread.task)?;
1012 let runtime_instance = task.instance;
1013 let handle = waitable.map(|(_, v)| v).unwrap_or(0);
1014
1015 log::trace!(
1016 "use callback to deliver event {event:?} to {:?} for {waitable:?}",
1017 call.thread,
1018 );
1019
1020 let old_thread = store.set_thread(call.thread)?;
1021 log::trace!(
1022 "GuestCallKind::DeliverEvent: replaced {old_thread:?} with {:?} as current thread",
1023 call.thread
1024 );
1025
1026 store.enter_instance(runtime_instance);
1027
1028 let Some(callback) = store
1029 .concurrent_state_mut()?
1030 .get_mut(call.thread.task)?
1031 .callback
1032 .take()
1033 else {
1034 bail_bug!("guest task callback field not present")
1035 };
1036
1037 let code = callback(store, event, handle)?;
1038
1039 store
1040 .concurrent_state_mut()?
1041 .get_mut(call.thread.task)?
1042 .callback = Some(callback);
1043
1044 store.exit_instance(runtime_instance)?;
1045
1046 store.set_thread(old_thread)?;
1047
1048 instance.handle_callback_code(store, call.thread, runtime_instance.index, code)?;
1049
1050 log::trace!("GuestCallKind::DeliverEvent: restored {old_thread:?} as current thread");
1051 }
1052 GuestCallKind::StartImplicit(fun) => {
1053 fun(store)?;
1054 }
1055 GuestCallKind::StartExplicit(fun) => {
1056 fun(store)?;
1057 }
1058 }
1059
1060 Ok(())
1061}
1062
1063impl<T> Store<T> {
1064 pub async fn run_concurrent<R>(&mut self, fun: impl AsyncFnOnce(&Accessor<T>) -> R) -> Result<R>
1066 where
1067 T: Send + 'static,
1068 {
1069 ensure!(
1070 self.as_context().0.concurrency_support(),
1071 "cannot use `run_concurrent` when Config::concurrency_support disabled",
1072 );
1073 self.as_context_mut().run_concurrent(fun).await
1074 }
1075
1076 #[doc(hidden)]
1077 pub fn assert_concurrent_state_empty(&mut self) {
1078 self.as_context_mut().assert_concurrent_state_empty();
1079 }
1080
1081 #[doc(hidden)]
1082 pub fn concurrent_state_table_size(&mut self) -> usize {
1083 self.as_context_mut().concurrent_state_table_size()
1084 }
1085
1086 pub fn spawn(
1088 &mut self,
1089 task: impl for<'fut> AccessorTask<'fut, T, HasSelf<T>>,
1090 ) -> Result<JoinHandle>
1091 where
1092 T: 'static,
1093 {
1094 self.as_context_mut().spawn(task)
1095 }
1096}
1097
1098impl<T> StoreContextMut<'_, T> {
1099 #[doc(hidden)]
1110 pub fn assert_concurrent_state_empty(self) {
1111 let store = self.0;
1112 store
1113 .store_data_mut()
1114 .components
1115 .assert_instance_states_empty();
1116 let state = store.concurrent_state_mut().unwrap();
1117 assert!(
1118 state.table.get_mut().is_empty(),
1119 "non-empty table: {:?}",
1120 state.table.get_mut()
1121 );
1122 assert!(state.switch_item.is_none());
1123 assert!(state.next_switch_item.is_none());
1124 assert!(state.high_priority.is_empty());
1125 assert!(state.low_priority.is_empty());
1126 assert!(state.unforced_current_thread.is_none());
1127 assert!(state.deferred_host_call_context.is_none());
1128 assert!(state.futures_mut().unwrap().is_empty());
1129 assert!(state.global_error_context_ref_counts.is_empty());
1130 }
1131
1132 #[doc(hidden)]
1137 pub fn concurrent_state_table_size(&mut self) -> usize {
1138 self.0
1139 .concurrent_state_mut()
1140 .unwrap()
1141 .table
1142 .get_mut()
1143 .iter_mut()
1144 .count()
1145 }
1146
1147 pub fn spawn(mut self, task: impl for<'fut> AccessorTask<'fut, T>) -> Result<JoinHandle>
1157 where
1158 T: 'static,
1159 {
1160 let accessor = Accessor::new(StoreToken::new(self.as_context_mut()));
1161 self.spawn_with_accessor(accessor, task)
1162 }
1163
1164 fn spawn_with_accessor<D>(
1167 self,
1168 accessor: Accessor<T, D>,
1169 task: impl for<'fut> AccessorTask<'fut, T, D>,
1170 ) -> Result<JoinHandle>
1171 where
1172 T: 'static,
1173 D: HasData + ?Sized,
1174 {
1175 let (handle, future) = JoinHandle::run(async move { task.run(&accessor).await });
1179 self.0
1180 .concurrent_state_mut()?
1181 .push_future(Box::pin(async move { future.await.unwrap_or(Ok(())) }));
1182 Ok(handle)
1183 }
1184
1185 pub async fn run_concurrent<R>(self, fun: impl AsyncFnOnce(&Accessor<T>) -> R) -> Result<R>
1269 where
1270 T: Send + 'static,
1271 {
1272 ensure!(
1273 self.0.concurrency_support(),
1274 "cannot use `run_concurrent` when Config::concurrency_support disabled",
1275 );
1276 self.do_run_concurrent(fun, false).await
1277 }
1278
1279 pub(super) async fn run_concurrent_trap_on_idle<R>(
1280 self,
1281 fun: impl AsyncFnOnce(&Accessor<T>) -> R,
1282 ) -> Result<R> {
1283 self.do_run_concurrent(fun, true).await
1284 }
1285
1286 async fn do_run_concurrent<R>(
1287 mut self,
1288 fun: impl AsyncFnOnce(&Accessor<T>) -> R,
1289 trap_on_idle: bool,
1290 ) -> Result<R> {
1291 debug_assert!(self.0.concurrency_support());
1292 let already_running = self
1293 .0
1294 .concurrent_state_mut_already_forced_current_thread()
1295 .event_loop_running;
1296 if already_running {
1297 bail!("Recursive `StoreContextMut::run_concurrent` calls not supported")
1298 }
1299 let token = StoreToken::new(self.as_context_mut());
1300
1301 struct Dropper<'a, T: 'static, V> {
1302 store: StoreContextMut<'a, T>,
1303 value: ManuallyDrop<V>,
1304 }
1305
1306 impl<'a, T, V> Drop for Dropper<'a, T, V> {
1307 fn drop(&mut self) {
1308 self.store
1309 .0
1310 .concurrent_state_mut_already_forced_current_thread()
1311 .event_loop_running = false;
1312
1313 tls::set(self.store.0, || {
1314 unsafe { ManuallyDrop::drop(&mut self.value) }
1319 });
1320 }
1321 }
1322
1323 let accessor = &Accessor::new(token);
1324 self.0
1325 .concurrent_state_mut_already_forced_current_thread()
1326 .event_loop_running = true;
1327 let dropper = &mut Dropper {
1328 store: self,
1329 value: ManuallyDrop::new(fun(accessor)),
1330 };
1331 let future = unsafe { Pin::new_unchecked(dropper.value.deref_mut()) };
1333
1334 let result = dropper
1335 .store
1336 .as_context_mut()
1337 .poll_until(future, trap_on_idle)
1338 .await;
1339
1340 if result.is_err() {
1341 dropper.store.0.set_trapped();
1342 }
1343
1344 result
1345 }
1346
1347 async fn poll_until<R>(
1353 mut self,
1354 mut future: Pin<&mut impl Future<Output = R>>,
1355 trap_on_idle: bool,
1356 ) -> Result<R> {
1357 struct Reset<'a, T: 'static> {
1358 store: StoreContextMut<'a, T>,
1359 futures: Option<FuturesUnordered<HostTaskFuture>>,
1360 }
1361
1362 impl<'a, T> Drop for Reset<'a, T> {
1363 fn drop(&mut self) {
1364 if let Some(futures) = self.futures.take() {
1365 *self
1366 .store
1367 .0
1368 .concurrent_state_mut_already_forced_current_thread()
1369 .futures
1370 .get_mut() = Some(futures);
1371 }
1372 }
1373 }
1374
1375 const MAX_TURNS_WITHOUT_YIELD: usize = 128;
1379 let mut turns_without_yield = 0;
1380
1381 loop {
1382 let futures = self.0.concurrent_state_mut()?.futures.get_mut().take();
1386 let mut reset = Reset {
1387 store: self.as_context_mut(),
1388 futures,
1389 };
1390 let mut next = match reset.futures.as_mut() {
1391 Some(f) => pin!(f.next()),
1392 None => bail_bug!("concurrent state missing futures field"),
1393 };
1394
1395 enum PollResult<R> {
1396 Complete(R),
1397 ProcessWork {
1398 ready: Option<WorkItem>,
1399 low_priority: bool,
1400 },
1401 }
1402
1403 let result = future::poll_fn(|cx| {
1404 if let Poll::Ready(value) = tls::set(reset.store.0, || future.as_mut().poll(cx)) {
1407 return Poll::Ready(Ok(PollResult::Complete(value)));
1408 }
1409
1410 if reset.store.0.trapped() {
1419 return Poll::Ready(Err(Trap::CannotEnterComponent.into()));
1420 }
1421
1422 let next = match tls::set(reset.store.0, || next.as_mut().poll(cx)) {
1426 Poll::Ready(Some(output)) => {
1427 match output {
1428 Err(e) => return Poll::Ready(Err(e)),
1429 Ok(()) => {}
1430 }
1431 Poll::Ready(true)
1432 }
1433 Poll::Ready(None) => Poll::Ready(false),
1434 Poll::Pending => Poll::Pending,
1435 };
1436
1437 let state = reset.store.0.concurrent_state_mut()?;
1452 let mut ready = state.switch_item.take();
1453 let mut low_priority = false;
1454 if ready.is_none() {
1455 ready = state.high_priority.pop_back();
1456 if ready.is_none() {
1457 ready = state.low_priority.pop_back();
1458 low_priority = true;
1459 }
1460 }
1461 if ready.is_some() {
1462 return Poll::Ready(Ok(PollResult::ProcessWork {
1463 ready,
1464 low_priority,
1465 }));
1466 }
1467
1468 return match next {
1472 Poll::Ready(true) => {
1473 Poll::Ready(Ok(PollResult::ProcessWork {
1479 ready: None,
1480 low_priority: false,
1481 }))
1482 }
1483 Poll::Ready(false) => {
1484 if let Poll::Ready(value) =
1488 tls::set(reset.store.0, || future.as_mut().poll(cx))
1489 {
1490 Poll::Ready(Ok(PollResult::Complete(value)))
1491 } else {
1492 if trap_on_idle {
1498 Poll::Ready(Err(if reset.store.0.any_may_not_suspend()? {
1505 Trap::CannotBlockSyncTask.into()
1506 } else {
1507 Trap::AsyncDeadlock.into()
1509 }))
1510 } else {
1511 Poll::Pending
1515 }
1516 }
1517 }
1518 Poll::Pending => Poll::Pending,
1523 };
1524 })
1525 .await;
1526
1527 drop(reset);
1531
1532 match result? {
1533 PollResult::Complete(value) => break Ok(value),
1536 PollResult::ProcessWork {
1539 ready,
1540 low_priority,
1541 } => {
1542 struct Dispose<'a, T: 'static> {
1543 store: StoreContextMut<'a, T>,
1544 ready: Option<WorkItem>,
1545 }
1546
1547 impl<'a, T> Drop for Dispose<'a, T> {
1548 fn drop(&mut self) {
1549 if let Some(item) = self.ready.take() {
1550 match item {
1551 WorkItem::ResumeFiber { mut fiber, .. } => {
1552 fiber.dispose(self.store.0);
1553 }
1554 WorkItem::PushFuture(future) => {
1555 tls::set(self.store.0, move || drop(future))
1556 }
1557 _ => {}
1558 }
1559 }
1560 }
1561 }
1562
1563 let mut dispose = Dispose {
1564 store: self.as_context_mut(),
1565 ready,
1566 };
1567
1568 if low_priority {
1590 dispose.store.0.yield_now().await;
1591 turns_without_yield = 0;
1592 }
1593
1594 if let Some(item) = dispose.ready.take() {
1595 dispose
1596 .store
1597 .as_context_mut()
1598 .handle_work_item(item)
1599 .await?;
1600 }
1601
1602 turns_without_yield += 1;
1603 if turns_without_yield == MAX_TURNS_WITHOUT_YIELD {
1604 turns_without_yield = 0;
1605 dispose.store.0.yield_now().await;
1606 }
1607 }
1608 }
1609 }
1610 }
1611
1612 async fn handle_work_item(self, item: WorkItem) -> Result<()> {
1614 log::trace!("handle work item {item:?}");
1615 match item {
1616 WorkItem::PushFuture(future) => {
1617 self.0
1618 .concurrent_state_mut()?
1619 .futures_mut()?
1620 .push(future.into_inner());
1621 }
1622 WorkItem::ResumeFiber { fiber, .. } => {
1623 self.0.resume_fiber(fiber).await?;
1624 }
1625 WorkItem::ResumeThread { thread, .. } => {
1626 if let GuestThreadState::Ready { fiber, .. } = mem::replace(
1627 &mut self.0.concurrent_state_mut()?.get_mut(thread.thread)?.state,
1628 GuestThreadState::Running,
1629 ) {
1630 self.0.resume_fiber(fiber).await?;
1631 } else {
1632 bail_bug!("cannot resume non-pending thread {thread:?}");
1633 }
1634 }
1635 WorkItem::GuestCall { call, .. } => {
1636 if call.is_ready(self.0)? {
1637 self.0
1638 .concurrent_state_mut()?
1639 .get_mut(call.thread.thread)?
1640 .wake_on_cancel = WakeOnCancel::None;
1641 self.run_on_worker(WorkerItem::GuestCall(call)).await?;
1642 } else {
1643 let state = self.0.concurrent_state_mut()?;
1644 let task = state.get_mut(call.thread.task)?;
1645 if !task.starting_sent {
1646 task.starting_sent = true;
1647 if let GuestCallKind::StartImplicit(_) = &call.kind {
1648 Waitable::Guest(call.thread.task).set_event(
1649 state,
1650 Some(Event::Subtask {
1651 status: Status::Starting,
1652 }),
1653 )?;
1654 }
1655 }
1656
1657 let instance = state.get_mut(call.thread.task)?.instance;
1658 self.0
1659 .instance_state(instance)
1660 .concurrent_state()
1661 .pending
1662 .insert(call.thread, call.kind);
1663
1664 self.0.concurrent_state_mut()?.take_next_switch_item()?;
1668 }
1669 }
1670 WorkItem::WorkerFunction(fun) => {
1671 self.run_on_worker(WorkerItem::Function(fun)).await?;
1672 }
1673 }
1674
1675 Ok(())
1676 }
1677
1678 async fn run_on_worker(self, item: WorkerItem) -> Result<()> {
1680 let worker = if let Some(fiber) = self.0.concurrent_state_mut()?.worker.take() {
1681 fiber
1682 } else {
1683 unsafe {
1702 fiber::make_fiber_unchecked(self.0, move |store| {
1703 loop {
1704 let Some(item) = store.concurrent_state_mut()?.worker_item.take() else {
1705 bail_bug!("worker_item not present when resuming fiber")
1706 };
1707 match item {
1708 WorkerItem::GuestCall(call) => handle_guest_call(store, call)?,
1709 WorkerItem::Function(fun) => fun.into_inner()(store)?,
1710 }
1711
1712 store.suspend(SuspendReason::NeedWork)?;
1713 }
1714 })?
1715 }
1716 };
1717
1718 let worker_item = &mut self.0.concurrent_state_mut()?.worker_item;
1719 assert!(worker_item.is_none());
1720 *worker_item = Some(item);
1721
1722 self.0.resume_fiber(worker).await
1723 }
1724
1725 pub(crate) fn wrap_call<F, R>(self, closure: F) -> impl Future<Output = Result<R>> + 'static
1730 where
1731 T: 'static,
1732 F: FnOnce(&Accessor<T>) -> Pin<Box<dyn Future<Output = Result<R>> + Send + '_>>
1733 + Send
1734 + Sync
1735 + 'static,
1736 R: Send + Sync + 'static,
1737 {
1738 let token = StoreToken::new(self);
1739 async move {
1740 let mut accessor = Accessor::new(token);
1741 closure(&mut accessor).await
1742 }
1743 }
1744
1745 pub(crate) async fn start_instance(
1746 &mut self,
1747 instance: ModuleInstance,
1748 callee: Option<RuntimeInstance>,
1749 ) -> Result<ModuleInstance> {
1750 let (tx, rx) = oneshot::channel();
1751 let token = StoreToken::new(self.as_context_mut());
1752 self.0.queue_task(move |store| {
1753 _ = tx.send(
1754 super::instance::start_raw(&mut token.as_context_mut(store), instance, callee)
1755 .map(|()| instance),
1756 );
1757 Ok(())
1758 })?;
1759 self.as_context_mut()
1760 .run_concurrent_trap_on_idle(async |_| {
1761 rx.await
1762 .map_err(|_| format_err!("oneshot channel canceled"))
1763 })
1764 .await??
1765 }
1766}
1767
1768pub type EnteredHostTask = Option<QualifiedThreadId>;
1775
1776impl StoreOpaque {
1777 #[inline]
1781 pub(crate) fn current_thread(&mut self) -> Result<CurrentThread> {
1782 if !self.concurrency_support() {
1784 return Ok(CurrentThread::None);
1785 }
1786
1787 if !self
1790 .vm_store_context_mut()
1791 .current_thread_mut()
1792 .is_deferred()
1793 {
1794 return Ok(self
1795 .concurrent_state_mut_already_forced_current_thread()
1796 .unforced_current_thread);
1797 }
1798
1799 self.force_deferred_current_thread()
1800 }
1801
1802 #[cold]
1805 fn force_deferred_current_thread(&mut self) -> Result<CurrentThread> {
1806 let state = self.concurrent_state_mut_without_forcing_current_thread();
1815 let id = match state.unforced_current_thread.guest_task() {
1816 Some(task) => state.get_mut(task)?.instance.instance,
1817 None => bail_bug!("deferred component-model thread with non-guest base"),
1818 };
1819
1820 let mut frames = Vec::new();
1823 let mut cur = *self.vm_store_context_mut().current_thread_mut();
1824 while let Some(ptr) = cur.as_deferred() {
1825 let deferred = unsafe { ptr.as_non_null().as_ref() };
1830 frames.push((
1831 deferred.callee_async != 0,
1832 deferred.callee_instance,
1833 deferred.saved_context,
1834 ));
1835 cur = deferred.parent;
1836 }
1837
1838 *self.vm_store_context_mut().current_thread_mut() = VMLazyThread::forced();
1842
1843 let current_context = *self.vm_store_context_mut().component_context_mut();
1846
1847 for (callee_async, callee_instance, saved_context) in frames.into_iter().rev() {
1851 *self.vm_store_context_mut().component_context_mut() = saved_context;
1855 let callee = RuntimeInstance {
1856 instance: id,
1857 index: RuntimeComponentInstanceIndex::from_u32(callee_instance),
1858 };
1859 self.enter_guest_sync_call(callee_async, callee)?;
1860 }
1861
1862 *self.vm_store_context_mut().component_context_mut() = current_context;
1864
1865 Ok(self
1866 .concurrent_state_mut_without_forcing_current_thread()
1867 .unforced_current_thread)
1868 }
1869
1870 fn current_guest_thread(&mut self) -> Result<QualifiedThreadId> {
1871 match self.current_thread()?.guest() {
1872 Some(id) => Ok(*id),
1873 None => bail_bug!("current thread is not a guest thread"),
1874 }
1875 }
1876
1877 pub(crate) fn current_materialized_host_task(&mut self) -> Result<Option<TableId<HostTask>>> {
1881 match self.current_thread()? {
1882 CurrentThread::Host(id) => Ok(Some(id)),
1883 CurrentThread::DeferredHost(_) | CurrentThread::None => Ok(None),
1884 _ => bail_bug!("current thread is not a host thread"),
1885 }
1886 }
1887
1888 fn materialize_host_task_id(&mut self) -> Result<Option<TableId<HostTask>>> {
1891 Ok(self
1892 .concurrent_state_mut()?
1893 .materialize_current_host_task_id()?)
1894 }
1895
1896 fn enter_sync_call(&mut self, callee: RuntimeInstance) -> Result<()> {
1897 log::trace!("enter sync-typed call {callee:?}");
1898 let state = self.instance_state(callee).concurrent_state();
1899 let old_do_not_suspend = state.do_not_suspend;
1900 state.do_not_suspend = true;
1901
1902 let thread = self.current_guest_thread()?;
1903 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
1904 if thread.old_do_not_suspend.is_some() {
1905 bail_bug!("current thread already has `old_do_not_suspend` value");
1906 }
1907
1908 thread.old_do_not_suspend = Some(old_do_not_suspend);
1909
1910 Ok(())
1911 }
1912
1913 fn exit_sync_call(&mut self, callee: RuntimeInstance) -> Result<()> {
1914 log::trace!("exit sync-typed call {callee:?}");
1915 let thread = self.current_guest_thread()?;
1916 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
1917 let Some(old_do_not_suspend) = thread.old_do_not_suspend.take() else {
1918 bail_bug!("current thread missing `old_do_not_suspend` value");
1919 };
1920 let state = self.instance_state(callee).concurrent_state();
1921 state.do_not_suspend = old_do_not_suspend;
1922 Ok(())
1923 }
1924
1925 pub(crate) fn enter_guest_sync_call(
1937 &mut self,
1938 callee_async_typed: bool,
1939 callee: RuntimeInstance,
1940 ) -> Result<()> {
1941 log::trace!("enter sync-lifted call {callee:?}");
1942 if !self.concurrency_support() {
1943 return self.enter_call_not_concurrent();
1944 }
1945
1946 let thread = self.current_thread()?;
1947 let caller = if let Some(thread) = thread.guest() {
1948 Caller::Guest { thread: *thread }
1949 } else {
1950 Caller::Host {
1951 tx: None,
1952 host_future_present: false,
1953 caller: self.materialize_host_task_id()?,
1954 }
1955 };
1956
1957 let state = self.concurrent_state_mut()?;
1958 let guest_thread = GuestTask::new(
1959 state,
1960 Box::new(move |_, _| bail_bug!("cannot lower params in sync call")),
1961 LiftResult {
1962 lift: Box::new(move |_, _| bail_bug!("cannot lift result in sync call")),
1963 ty: TypeTupleIndex::reserved_value(),
1964 memory: None,
1965 string_encoding: StringEncoding::Utf8,
1966 },
1967 caller,
1968 None,
1969 callee,
1970 callee_async_typed,
1971 false,
1972 )?;
1973
1974 Instance::from_wasmtime(self, callee.instance).add_guest_thread_to_instance_table(
1975 guest_thread.thread,
1976 self,
1977 callee.index,
1978 )?;
1979 self.set_thread(guest_thread)?;
1980
1981 if !callee_async_typed {
1982 self.enter_sync_call(callee)?;
1983 }
1984
1985 Ok(())
1986 }
1987
1988 pub(crate) fn exit_guest_sync_call(&mut self) -> Result<()> {
1996 if !self.concurrency_support() {
1997 return Ok(self.exit_call_not_concurrent());
1998 }
1999
2000 let thread = match self.current_thread()?.guest() {
2001 Some(t) => *t,
2002 None => bail_bug!("expected task when exiting"),
2003 };
2004 let task = self.concurrent_state_mut()?.get_mut(thread.task)?;
2005 let instance = task.instance;
2006
2007 let caller = match &task.caller {
2008 &Caller::Guest { thread } => thread.into(),
2009 &Caller::Host { caller, .. } => caller
2010 .map(CurrentThread::Host)
2011 .unwrap_or(CurrentThread::None),
2012 };
2013 task.lift_result = None;
2014 task.exited = true;
2015 let async_typed = task.async_typed;
2016
2017 if !async_typed {
2018 self.exit_sync_call(instance)?;
2019 }
2020
2021 self.set_thread(caller)?;
2022
2023 log::trace!("exit sync-lifted call {instance:?}");
2024
2025 if async_typed {
2026 self.switch_or_trap_if_may_not_suspend(instance)?;
2031 }
2032
2033 self.cleanup_thread(thread, instance, CleanupTask::Yes)?;
2034
2035 Ok(())
2036 }
2037
2038 pub(crate) fn host_task_create(&mut self) -> Result<EnteredHostTask> {
2045 if !self.concurrency_support() {
2046 self.enter_call_not_concurrent()?;
2047 return Ok(None);
2048 }
2049 let caller = self.current_guest_thread()?;
2050 log::trace!("new deferred host task with caller {caller:?}");
2051
2052 self.set_thread(CurrentThread::DeferredHost(caller))?;
2053 let state = self.concurrent_state_mut()?;
2054 debug_assert!(state.deferred_host_call_context.is_none());
2055 state.deferred_host_call_context = Some(CallContext::default());
2056 state.debug_assert_deferred_host_invariant();
2057 Ok(Some(caller))
2058 }
2059
2060 pub(crate) fn host_task_delete(
2067 &mut self,
2068 original_task: EnteredHostTask,
2069 materialized_task: Option<TableId<HostTask>>,
2070 ) -> Result<()> {
2071 match original_task {
2072 Some(caller) => {
2073 self.set_thread(caller)?;
2074 if materialized_task.is_none() {
2075 let state = self.concurrent_state_mut()?;
2076 let context = state
2077 .deferred_host_call_context
2078 .take()
2079 .expect("deferred host call context should be present");
2080 debug_assert!(context.is_empty());
2081 state.debug_assert_deferred_host_invariant();
2082 }
2083 log::trace!(
2084 "delete host task with caller {original_task:?} and materialized as {materialized_task:?}"
2085 );
2086 if let Some(task) = materialized_task {
2087 Waitable::Host(task).delete_from(self)?;
2088 }
2089 }
2090 None => {
2091 debug_assert!(materialized_task.is_none());
2092 self.exit_call_not_concurrent();
2093 }
2094 }
2095 Ok(())
2096 }
2097
2098 fn instance_state(&mut self, instance: RuntimeInstance) -> &mut InstanceState {
2101 self.component_instance_mut(instance.instance)
2102 .instance_state(instance.index)
2103 }
2104
2105 pub(crate) fn set_thread(&mut self, thread: impl Into<CurrentThread>) -> Result<CurrentThread> {
2111 let thread = thread.into();
2112 let state = self.concurrent_state_mut()?;
2113 state.debug_assert_deferred_host_invariant();
2114 let old_thread = mem::replace(&mut state.unforced_current_thread, thread);
2115
2116 state.handle_thread_switch(old_thread, thread)?;
2117
2118 if let Some(old_thread) = old_thread.guest() {
2126 let old_context = *self.vm_store_context_mut().component_context_mut();
2127 self.concurrent_state_mut()?
2128 .get_mut(old_thread.thread)?
2129 .context = old_context;
2130 }
2131 if cfg!(debug_assertions) {
2132 *self.vm_store_context_mut().component_context_mut() =
2133 [u32::MAX; NUM_COMPONENT_CONTEXT_SLOTS];
2134 }
2135 if let Some(thread) = thread.guest() {
2136 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
2137 let context = thread.context;
2138 if cfg!(debug_assertions) {
2139 thread.context = [u32::MAX; NUM_COMPONENT_CONTEXT_SLOTS];
2140 }
2141 *self.vm_store_context_mut().component_context_mut() = context;
2142 }
2143
2144 *self.vm_store_context_mut().current_thread_mut() = if thread.is_none() {
2146 VMLazyThread::none()
2147 } else {
2148 VMLazyThread::forced()
2149 };
2150
2151 Ok(old_thread)
2152 }
2153
2154 fn switch_or_trap_if_may_not_suspend(&mut self, instance: RuntimeInstance) -> Result<()> {
2156 if self.switch_if_may_not_suspend(instance)? {
2157 Ok(())
2158 } else {
2159 Err(Trap::CannotBlockSyncTask.into())
2160 }
2161 }
2162
2163 fn switch_if_may_not_suspend(&mut self, instance: RuntimeInstance) -> Result<bool> {
2167 self.concurrent_state_mut()?;
2171
2172 Ok(!self.concurrency_support()
2173 || !self
2174 .instance_state(instance)
2175 .concurrent_state()
2176 .do_not_suspend
2177 || self
2178 .concurrent_state_mut()?
2179 .promote_instance_local_thread_work_item(instance)?)
2180 }
2181
2182 fn enter_instance(&mut self, instance: RuntimeInstance) {
2186 log::trace!("enter {instance:?}");
2187 self.instance_state(instance)
2188 .concurrent_state()
2189 .do_not_enter = true;
2190 }
2191
2192 fn exit_instance(&mut self, instance: RuntimeInstance) -> Result<()> {
2196 log::trace!("exit {instance:?}");
2197 self.instance_state(instance)
2198 .concurrent_state()
2199 .do_not_enter = false;
2200 self.partition_pending(instance)
2201 }
2202
2203 fn partition_pending(&mut self, instance: RuntimeInstance) -> Result<()> {
2211 for (thread, kind) in
2212 mem::take(&mut self.instance_state(instance).concurrent_state().pending).into_iter()
2213 {
2214 let call = GuestCall { thread, kind };
2215 if call.is_ready(self)? {
2216 self.concurrent_state_mut()?
2217 .push_high_priority(WorkItem::GuestCall { instance, call });
2218 } else {
2219 self.instance_state(instance)
2220 .concurrent_state()
2221 .pending
2222 .insert(call.thread, call.kind);
2223 }
2224 }
2225
2226 if let Some(waker) = self
2227 .concurrent_state_mut()?
2228 .ready_for_concurrent_call_waker
2229 .take()
2230 {
2231 waker.wake();
2232 }
2233
2234 Ok(())
2235 }
2236
2237 pub(crate) fn backpressure_modify(
2239 &mut self,
2240 caller_instance: RuntimeInstance,
2241 modify: impl FnOnce(u16) -> Option<u16>,
2242 ) -> Result<()> {
2243 let state = self.instance_state(caller_instance).concurrent_state();
2244 let old = state.backpressure;
2245 let new = modify(old).ok_or_else(|| Trap::BackpressureOverflow)?;
2246 state.backpressure = new;
2247
2248 if old > 0 && new == 0 {
2249 self.partition_pending(caller_instance)?;
2252 }
2253
2254 Ok(())
2255 }
2256
2257 async fn resume_fiber(&mut self, fiber: StoreFiber<'static>) -> Result<()> {
2260 let old_thread = self.current_thread()?;
2261 log::trace!("resume_fiber: save current thread {old_thread:?}");
2262
2263 let fiber = fiber::resolve_or_release(self, fiber).await?;
2264
2265 self.set_thread(old_thread)?;
2266
2267 let state = self.concurrent_state_mut()?;
2268
2269 if let Some(ot) = old_thread.guest() {
2270 state.get_mut(ot.thread)?.state = GuestThreadState::Running;
2271 }
2272 log::trace!("resume_fiber: restore current thread {old_thread:?}");
2273
2274 if let Some(mut fiber) = fiber {
2275 log::trace!("resume_fiber: suspend reason {:?}", &state.suspend_reason);
2276 let reason = match state.suspend_reason.take() {
2278 Some(r) => r,
2279 None => bail_bug!("suspend reason missing when resuming fiber"),
2280 };
2281 match reason {
2282 SuspendReason::NeedWork => {
2283 if state.worker.is_none() {
2284 state.worker = Some(fiber);
2285 } else {
2286 fiber.dispose(self);
2287 }
2288 }
2289 SuspendReason::Yielding { thread } => {
2290 state.get_mut(thread.thread)?.state = GuestThreadState::Ready { fiber };
2291 let instance = state.get_mut(thread.task)?.instance;
2292 state.push_low_priority(WorkItem::ResumeThread { instance, thread });
2293 }
2294 SuspendReason::ExplicitlySuspending { thread } => {
2295 state.get_mut(thread.thread)?.state = GuestThreadState::Suspended(fiber);
2296 }
2297 SuspendReason::Waiting { set, thread } => {
2298 let old = state
2299 .get_mut(set)?
2300 .waiting
2301 .insert(thread, WaitMode::Fiber(fiber));
2302 assert!(old.is_none());
2303 }
2304 SuspendReason::YieldingToSubtask { thread } => {
2305 let item = WorkItem::ResumeFiber {
2314 instance: state.get_mut(thread.task)?.instance,
2315 thread,
2316 fiber,
2317 };
2318
2319 if state.next_switch_item.replace(item).is_some() {
2320 bail_bug!(
2323 "`ConcurrentState::next_switch_item` was already `Some(_)` when \
2324 a thread wanted to wait on a subtask"
2325 );
2326 }
2327 }
2328 };
2329 } else {
2330 log::trace!("resume_fiber: fiber has exited");
2331 }
2332
2333 Ok(())
2334 }
2335
2336 fn suspend(&mut self, reason: SuspendReason) -> Result<()> {
2342 log::trace!("suspend fiber: {reason:?}");
2343
2344 let state = self.concurrent_state_mut()?;
2345
2346 let (save_and_restore_thread, save_and_restore_next_switch_item) = match &reason {
2353 SuspendReason::Yielding { .. }
2354 | SuspendReason::Waiting { .. }
2355 | SuspendReason::ExplicitlySuspending { .. } => {
2356 if state.switch_item.is_none() {
2359 state.take_next_switch_item()?;
2360 }
2361
2362 (true, false)
2363 }
2364 SuspendReason::YieldingToSubtask { .. } => (true, true),
2365 SuspendReason::NeedWork => (false, false),
2366 };
2367
2368 let old_next_switch_item = if save_and_restore_next_switch_item {
2369 let item = state.next_switch_item.take();
2370 Some(state.push(item)?)
2374 } else {
2375 None
2376 };
2377
2378 let old_guest_thread = if save_and_restore_thread {
2379 self.current_thread()?
2380 } else {
2381 CurrentThread::None
2382 };
2383
2384 let waiting_set = match &reason {
2385 SuspendReason::Waiting { set, .. } => Some(*set),
2386 _ => None,
2387 };
2388
2389 let suspend_reason = &mut self.concurrent_state_mut()?.suspend_reason;
2390 assert!(suspend_reason.is_none());
2391 *suspend_reason = Some(reason);
2392
2393 if !self.fiber_async_state_mut().can_block() {
2396 return Err(format_err!("future dropped"));
2397 }
2398
2399 if let Some(set) = waiting_set {
2402 self.concurrent_state_mut()?.get_mut(set)?.num_waiting += 1;
2403 }
2404
2405 self.with_blocking(|_, cx| cx.suspend(StoreFiberYield::ReleaseStore))?;
2406
2407 if let Some(set) = waiting_set {
2408 self.concurrent_state_mut()?.get_mut(set)?.stop_waiting()?;
2409 }
2410
2411 if save_and_restore_thread {
2412 self.set_thread(old_guest_thread)?;
2413 }
2414
2415 if let Some(item) = old_next_switch_item {
2416 let state = self.concurrent_state_mut()?;
2417 state.next_switch_item = state.delete(item)?;
2418 }
2419
2420 Ok(())
2421 }
2422
2423 fn wait_for_event(
2424 &mut self,
2425 caller_instance: RuntimeInstance,
2426 waitable: Waitable,
2427 ) -> Result<()> {
2428 let caller = self.current_guest_thread()?;
2429 let state = self.concurrent_state_mut()?;
2430
2431 waitable.trap_if_in_waitable_set(state)?;
2432
2433 let set = state.get_mut(caller.thread)?.sync_call_set;
2434 waitable.join(state, Some(set))?;
2435
2436 self.switch_or_trap_if_may_not_suspend(caller_instance)?;
2437
2438 self.suspend(SuspendReason::Waiting {
2439 set,
2440 thread: caller,
2441 })?;
2442 let state = self.concurrent_state_mut()?;
2443
2444 waitable.join(state, None)
2445 }
2446
2447 fn cleanup_thread(
2469 &mut self,
2470 guest_thread: QualifiedThreadId,
2471 runtime_instance: RuntimeInstance,
2472 cleanup_task: CleanupTask,
2473 ) -> Result<()> {
2474 let state = self.concurrent_state_mut()?;
2475 state.take_next_switch_item()?;
2478 let thread_data = state.get_mut(guest_thread.thread)?;
2479 let sync_call_set = thread_data.sync_call_set;
2480 if let Some(guest_id) = thread_data.instance_rep {
2481 self.instance_state(runtime_instance)
2482 .thread_handle_table()
2483 .guest_thread_remove(guest_id)?;
2484 }
2485 let state = self.concurrent_state_mut()?;
2486
2487 for waitable in mem::take(&mut state.get_mut(sync_call_set)?.ready) {
2489 if let Some(Event::Subtask {
2490 status: Status::Returned | Status::ReturnCancelled,
2491 }) = waitable.common(self.concurrent_state_mut()?)?.event
2492 {
2493 waitable.delete_from(self)?;
2494 }
2495 }
2496
2497 let state = self.concurrent_state_mut()?;
2498 state.delete(guest_thread.thread)?;
2499 state.delete(sync_call_set)?;
2500 let task = state.get_mut(guest_thread.task)?;
2501 task.threads.remove(&guest_thread.thread);
2502
2503 if task.threads.is_empty() && !task.returned_or_cancelled() {
2504 bail!(Trap::NoAsyncResult);
2505 }
2506 let ready_to_delete = task.ready_to_delete();
2507
2508 if !task.decremented_interesting_task_count && task.exited && task.returned_or_cancelled() {
2509 task.decremented_interesting_task_count = true;
2510
2511 debug_assert!(state.interesting_tasks > 0);
2512 state.interesting_tasks -= 1;
2513 if state.interesting_tasks == 0
2514 && let Some(waker) = state.interesting_tasks_empty_waker.take()
2515 {
2516 waker.wake();
2517 }
2518 }
2519
2520 match cleanup_task {
2521 CleanupTask::Yes => {
2522 if ready_to_delete {
2523 Waitable::Guest(guest_thread.task).delete_from(self)?;
2524 }
2525 }
2526 CleanupTask::No => {}
2527 }
2528
2529 Ok(())
2530 }
2531
2532 fn cancel_guest_subtask_without_lowered_parameters(
2545 &mut self,
2546 caller_instance: RuntimeInstance,
2547 guest_task: TableId<GuestTask>,
2548 ) -> Result<()> {
2549 let concurrent_state = self.concurrent_state_mut()?;
2550 let task = concurrent_state.get_mut(guest_task)?;
2551 assert!(!task.already_lowered_parameters());
2552 task.lower_params = None;
2556 task.lift_result = None;
2557 task.exited = true;
2558 let instance = task.instance;
2559
2560 assert_eq!(1, task.threads.len());
2563 let thread = *task.threads.iter().next().unwrap();
2564 self.cleanup_thread(
2565 QualifiedThreadId {
2566 task: guest_task,
2567 thread,
2568 },
2569 caller_instance,
2570 CleanupTask::No,
2571 )?;
2572
2573 let pending = &mut self.instance_state(instance).concurrent_state().pending;
2575 let pending_count = pending.len();
2576 pending.retain(|thread, _| thread.task != guest_task);
2577 if pending.len() == pending_count {
2579 bail!(Trap::SubtaskCancelAfterTerminal);
2580 }
2581 Ok(())
2582 }
2583
2584 pub(crate) fn current_scope(&mut self) -> Result<Option<CurrentScope>> {
2587 if !self.concurrency_support() {
2588 return Ok(self
2589 .current_scope_id_not_concurrent()?
2590 .map(|id| CurrentScope::Id(Scope::Id(id))));
2591 }
2592
2593 Ok(match self.current_thread()? {
2594 CurrentThread::Guest(id) => Some(CurrentScope::Id(Scope::Id(id.task.rep()))),
2595 CurrentThread::Host(id) => Some(CurrentScope::Id(Scope::HostId(id.rep()))),
2596 CurrentThread::DeferredHost(_) => Some(CurrentScope::DeferredHost),
2597 CurrentThread::None => return Ok(None),
2598 })
2599 }
2600
2601 pub(crate) fn queue_task(
2602 &mut self,
2603 task: impl FnOnce(&mut dyn VMStore) -> Result<()> + Send + 'static,
2604 ) -> Result<()> {
2605 self.concurrent_state_mut()?
2606 .push_high_priority(WorkItem::WorkerFunction(AlwaysMut::new(Box::new(task))));
2607 Ok(())
2608 }
2609
2610 fn any_may_not_suspend(&mut self) -> Result<bool> {
2619 Ok(self
2627 .concurrent_state_mut()?
2628 .table
2629 .get_mut()
2630 .iter_mut()
2631 .filter_map(|(_, entry)| {
2632 if let Some(task) = entry.downcast_ref::<GuestTask>() {
2633 Some(task.instance)
2634 } else {
2635 None
2636 }
2637 })
2638 .collect::<Vec<_>>()
2639 .into_iter()
2640 .any(|instance| {
2641 self.instance_state(instance)
2642 .concurrent_state()
2643 .do_not_suspend
2644 }))
2645 }
2646}
2647
2648enum CleanupTask {
2649 Yes,
2650 No,
2651}
2652
2653impl Instance {
2654 fn get_event(
2657 self,
2658 store: &mut StoreOpaque,
2659 guest_task: TableId<GuestTask>,
2660 set: Option<TableId<WaitableSet>>,
2661 cancellable: bool,
2662 ) -> Result<Option<(Event, Option<(Waitable, u32)>)>> {
2663 let state = store.concurrent_state_mut()?;
2664
2665 let task = state.get_mut(guest_task)?;
2666 let event = &mut task.event;
2667 if let Some(ev) = event
2668 && (cancellable || !matches!(ev, Event::Cancelled))
2669 {
2670 log::trace!("deliver event {ev:?} to {guest_task:?}");
2671
2672 if matches!(ev, Event::Cancelled) {
2673 task.cancel_request_delivered = true;
2674 }
2675
2676 let ev = *ev;
2677 *event = None;
2678 return Ok(Some((ev, None)));
2679 }
2680
2681 let set = match set {
2682 Some(set) => set,
2683 None => return Ok(None),
2684 };
2685 let waitable = match state.get_mut(set)?.ready.pop_first() {
2686 Some(v) => v,
2687 None => return Ok(None),
2688 };
2689
2690 let common = waitable.common(state)?;
2691 let handle = match common.handle {
2692 Some(h) => h,
2693 None => bail_bug!("handle not set when delivering event"),
2694 };
2695 let event = match common.event.take() {
2696 Some(e) => e,
2697 None => bail_bug!("event not set when delivering event"),
2698 };
2699
2700 log::trace!(
2701 "deliver event {event:?} to {guest_task:?} for {waitable:?} (handle {handle}); set {set:?}"
2702 );
2703
2704 waitable.on_delivery(store, self, event)?;
2705
2706 Ok(Some((event, Some((waitable, handle)))))
2707 }
2708
2709 fn handle_callback_code(
2715 self,
2716 store: &mut StoreOpaque,
2717 guest_thread: QualifiedThreadId,
2718 runtime_instance: RuntimeComponentInstanceIndex,
2719 code: u32,
2720 ) -> Result<()> {
2721 let (code, set) = unpack_callback_code(code);
2722
2723 log::trace!("received callback code from {guest_thread:?}: {code} (set: {set})");
2724
2725 let state = store.concurrent_state_mut()?;
2726
2727 state.take_next_switch_item()?;
2728
2729 let get_set = |store: &mut StoreOpaque, handle| -> Result<_> {
2730 let set = store
2731 .instance_state(self.runtime_instance(runtime_instance))
2732 .handle_table()
2733 .waitable_set_rep(handle)?;
2734
2735 Ok(TableId::<WaitableSet>::new(set))
2736 };
2737
2738 match code {
2739 callback_code::EXIT => {
2740 log::trace!("implicit thread {guest_thread:?} completed");
2741 let task = store.concurrent_state_mut()?.get_mut(guest_thread.task)?;
2742 task.exited = true;
2743 task.callback = None;
2744
2745 let runtime_instance = self.runtime_instance(runtime_instance);
2746
2747 store.switch_or_trap_if_may_not_suspend(runtime_instance)?;
2752
2753 store.cleanup_thread(guest_thread, runtime_instance, CleanupTask::Yes)?;
2754 }
2755 callback_code::YIELD => {
2756 let old = state
2759 .get_mut(guest_thread.thread)?
2760 .wake_on_cancel
2761 .replace(WakeOnCancel::Yielding);
2762 if !old.is_none() {
2763 bail_bug!("thread unexpectedly had wake_on_cancel set");
2764 }
2765
2766 let call = GuestCall {
2773 thread: guest_thread,
2774 kind: GuestCallKind::DeliverEvent {
2775 instance: self,
2776 set: None,
2777 },
2778 };
2779 state.push_low_priority(WorkItem::GuestCall {
2782 instance: self.runtime_instance(runtime_instance),
2783 call,
2784 });
2785 }
2786 callback_code::WAIT => {
2787 let set = get_set(store, set)?;
2788 let state = store.concurrent_state_mut()?;
2789 self.wait_with_callback(state, guest_thread, set)?;
2790 }
2791 _ => bail!(Trap::UnsupportedCallbackCode),
2792 }
2793
2794 Ok(())
2795 }
2796
2797 fn wait_with_callback(
2803 self,
2804 state: &mut ConcurrentState,
2805 guest_thread: QualifiedThreadId,
2806 set: TableId<WaitableSet>,
2807 ) -> Result<()> {
2808 state.get_mut(set)?.num_waiting += 1;
2811
2812 if state.get_mut(guest_thread.task)?.event.is_some()
2813 || !state.get_mut(set)?.ready.is_empty()
2814 {
2815 let instance = state.get_mut(guest_thread.task)?.instance;
2817 state.push_high_priority(WorkItem::GuestCall {
2818 instance,
2819 call: GuestCall {
2820 thread: guest_thread,
2821 kind: GuestCallKind::DeliverEvent {
2822 instance: self,
2823 set: Some(set),
2824 },
2825 },
2826 });
2827 return Ok(());
2828 }
2829
2830 let old = state
2836 .get_mut(guest_thread.thread)?
2837 .wake_on_cancel
2838 .replace(WakeOnCancel::Waiting(set));
2839 if !old.is_none() {
2840 bail_bug!("thread unexpectedly had wake_on_cancel set");
2841 }
2842 let old = state
2843 .get_mut(set)?
2844 .waiting
2845 .insert(guest_thread, WaitMode::Callback(self));
2846 if !old.is_none() {
2847 bail_bug!("set's waiting set already had this thread registered");
2848 }
2849 Ok(())
2850 }
2851
2852 unsafe fn stage_call<T: 'static>(
2859 self,
2860 mut store: StoreContextMut<T>,
2861 guest_thread: QualifiedThreadId,
2862 callee: SendSyncPtr<VMFuncRef>,
2863 param_count: usize,
2864 result_count: usize,
2865 async_: bool,
2866 callback: Option<SendSyncPtr<VMFuncRef>>,
2867 post_return: Option<SendSyncPtr<VMFuncRef>>,
2868 host_caller: bool,
2869 ) -> Result<()> {
2870 unsafe fn make_call<T: 'static>(
2885 store: StoreContextMut<T>,
2886 guest_thread: QualifiedThreadId,
2887 callee: SendSyncPtr<VMFuncRef>,
2888 param_count: usize,
2889 result_count: usize,
2890 ) -> impl FnOnce(&mut dyn VMStore) -> Result<[MaybeUninit<ValRaw>; MAX_FLAT_PARAMS]>
2891 + Send
2892 + Sync
2893 + 'static
2894 + use<T> {
2895 let token = StoreToken::new(store);
2896 move |store: &mut dyn VMStore| {
2897 let mut storage = [MaybeUninit::uninit(); MAX_FLAT_PARAMS];
2898
2899 store
2900 .concurrent_state_mut()?
2901 .get_mut(guest_thread.thread)?
2902 .state = GuestThreadState::Running;
2903 let task = store.concurrent_state_mut()?.get_mut(guest_thread.task)?;
2904 let lower = match task.lower_params.take() {
2905 Some(l) => l,
2906 None => bail_bug!("lower_params missing"),
2907 };
2908
2909 lower(store, &mut storage[..param_count])?;
2910
2911 let mut store = token.as_context_mut(store);
2912
2913 unsafe {
2916 crate::Func::call_unchecked_raw(
2917 &mut store,
2918 callee.as_non_null(),
2919 NonNull::new(
2920 &mut storage[..param_count.max(result_count)]
2921 as *mut [MaybeUninit<ValRaw>] as _,
2922 )
2923 .unwrap(),
2924 UncaughtException::Trap,
2925 )?;
2926 }
2927
2928 Ok(storage)
2929 }
2930 }
2931
2932 let call = unsafe {
2936 make_call(
2937 store.as_context_mut(),
2938 guest_thread,
2939 callee,
2940 param_count,
2941 result_count,
2942 )
2943 };
2944
2945 let callee_instance = store
2946 .0
2947 .concurrent_state_mut()?
2948 .get_mut(guest_thread.task)?
2949 .instance;
2950
2951 let fun = if callback.is_some() {
2952 assert!(async_);
2953
2954 Box::new(move |store: &mut dyn VMStore| {
2955 self.add_guest_thread_to_instance_table(
2956 guest_thread.thread,
2957 store,
2958 callee_instance.index,
2959 )?;
2960 let old_thread = store.set_thread(guest_thread)?;
2961 log::trace!(
2962 "stackless call: replaced {old_thread:?} with {guest_thread:?} as current thread"
2963 );
2964
2965 store.enter_instance(callee_instance);
2966
2967 let storage = call(store)?;
2974
2975 store.exit_instance(callee_instance)?;
2976
2977 store.set_thread(old_thread)?;
2978 let state = store.concurrent_state_mut()?;
2979 if let Some(t) = old_thread.guest() {
2980 state.get_mut(t.thread)?.state = GuestThreadState::Running;
2981 }
2982 log::trace!("stackless call: restored {old_thread:?} as current thread");
2983
2984 let code = unsafe { storage[0].assume_init() }.get_i32() as u32;
2987
2988 self.handle_callback_code(store, guest_thread, callee_instance.index, code)
2989 }) as Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send + Sync>
2990 } else {
2991 let token = StoreToken::new(store.as_context_mut());
2992 Box::new(move |store: &mut dyn VMStore| {
2993 self.add_guest_thread_to_instance_table(
2994 guest_thread.thread,
2995 store,
2996 callee_instance.index,
2997 )?;
2998 let old_thread = store.set_thread(guest_thread)?;
2999 log::trace!(
3000 "sync/async-stackful call: replaced {old_thread:?} with {guest_thread:?} as current thread",
3001 );
3002 let flags = self.id().get(store).instance_flags(callee_instance.index);
3003
3004 let callee_async_typed = store
3005 .concurrent_state_mut()?
3006 .get_mut(guest_thread.task)?
3007 .async_typed;
3008
3009 if !async_ && callee_async_typed {
3013 store.enter_instance(callee_instance);
3014 }
3015
3016 if !callee_async_typed {
3017 store.enter_sync_call(callee_instance)?;
3018 }
3019
3020 let storage = call(store)?;
3027
3028 if !callee_async_typed {
3029 store.exit_sync_call(callee_instance)?;
3030 }
3031
3032 if !async_ {
3033 if callee_async_typed {
3039 store.exit_instance(callee_instance)?;
3040 }
3041
3042 let lift = {
3043 let state = store.concurrent_state_mut()?;
3044 if !state.get_mut(guest_thread.task)?.result.is_none() {
3045 bail_bug!("task has already produced a result");
3046 }
3047
3048 match state.get_mut(guest_thread.task)?.lift_result.take() {
3049 Some(lift) => lift,
3050 None => bail_bug!("lift_result field is missing"),
3051 }
3052 };
3053
3054 let result = (lift.lift)(store, unsafe {
3057 mem::transmute::<&[MaybeUninit<ValRaw>], &[ValRaw]>(
3058 &storage[..result_count],
3059 )
3060 })?;
3061
3062 let post_return_arg = match result_count {
3063 0 => ValRaw::i32(0),
3064 1 => unsafe { storage[0].assume_init() },
3067 _ => unreachable!(),
3068 };
3069
3070 unsafe {
3071 call_post_return(
3072 token.as_context_mut(store),
3073 post_return.map(|v| v.as_non_null()),
3074 post_return_arg,
3075 flags,
3076 )?;
3077 }
3078
3079 self.task_complete(store, guest_thread.task, result, Status::Returned)?;
3080 }
3081
3082 store.set_thread(old_thread)?;
3083
3084 store
3085 .concurrent_state_mut()?
3086 .get_mut(guest_thread.task)?
3087 .exited = true;
3088
3089 log::trace!(
3090 "clean up thread; async lifted? {async_} async typed? {callee_async_typed}"
3091 );
3092
3093 if callee_async_typed {
3094 store.switch_or_trap_if_may_not_suspend(callee_instance)?;
3099 }
3100
3101 store.cleanup_thread(guest_thread, callee_instance, CleanupTask::Yes)?;
3103 Ok(())
3104 })
3105 };
3106
3107 store.0.concurrent_state_mut()?.push_work_item(
3108 WorkItem::GuestCall {
3109 instance: callee_instance,
3110 call: GuestCall {
3111 thread: guest_thread,
3112 kind: GuestCallKind::StartImplicit(fun),
3113 },
3114 },
3115 if host_caller {
3116 Priority::High
3117 } else {
3118 Priority::Switch
3119 },
3120 )?;
3121
3122 Ok(())
3123 }
3124
3125 unsafe fn prepare_call<T: 'static>(
3138 self,
3139 mut store: StoreContextMut<T>,
3140 start: NonNull<VMFuncRef>,
3141 return_: NonNull<VMFuncRef>,
3142 caller_instance: RuntimeComponentInstanceIndex,
3143 callee_instance: RuntimeComponentInstanceIndex,
3144 task_return_type: TypeTupleIndex,
3145 callee_async_typed: bool,
3146 memory: *mut VMMemoryDefinition,
3147 string_encoding: StringEncoding,
3148 caller_info: CallerInfo,
3149 ) -> Result<()> {
3150 enum ResultInfo {
3151 Heap { results: u32 },
3152 Stack { result_count: u32 },
3153 }
3154
3155 let result_info = match &caller_info {
3156 CallerInfo::Async {
3157 has_result: true,
3158 params,
3159 } => ResultInfo::Heap {
3160 results: match params.last() {
3161 Some(r) => r.get_u32(),
3162 None => bail_bug!("retptr missing"),
3163 },
3164 },
3165 CallerInfo::Async {
3166 has_result: false, ..
3167 } => ResultInfo::Stack { result_count: 0 },
3168 CallerInfo::Sync {
3169 result_count,
3170 params,
3171 } if *result_count > u32::try_from(MAX_FLAT_RESULTS)? => ResultInfo::Heap {
3172 results: match params.last() {
3173 Some(r) => r.get_u32(),
3174 None => bail_bug!("arg ptr missing"),
3175 },
3176 },
3177 CallerInfo::Sync { result_count, .. } => ResultInfo::Stack {
3178 result_count: *result_count,
3179 },
3180 };
3181
3182 let sync_caller = matches!(caller_info, CallerInfo::Sync { .. });
3183
3184 let start = SendSyncPtr::new(start);
3188 let return_ = SendSyncPtr::new(return_);
3189 let token = StoreToken::new(store.as_context_mut());
3190 let old_thread = store.0.current_guest_thread()?;
3191
3192 let state = store.0.concurrent_state_mut()?;
3193
3194 debug_assert_eq!(
3195 state.get_mut(old_thread.task)?.instance,
3196 self.runtime_instance(caller_instance)
3197 );
3198
3199 let guest_thread = GuestTask::new(
3200 state,
3201 Box::new(move |store, dst| {
3202 let mut store = token.as_context_mut(store);
3203 assert!(dst.len() <= MAX_FLAT_PARAMS);
3204 let mut src = [MaybeUninit::uninit(); MAX_FLAT_PARAMS + 1];
3206 let count = match caller_info {
3207 CallerInfo::Async { params, has_result } => {
3211 let params = ¶ms[..params.len() - usize::from(has_result)];
3212 for (param, src) in params.iter().zip(&mut src) {
3213 src.write(*param);
3214 }
3215 params.len()
3216 }
3217
3218 CallerInfo::Sync { params, .. } => {
3220 for (param, src) in params.iter().zip(&mut src) {
3221 src.write(*param);
3222 }
3223 params.len()
3224 }
3225 };
3226 unsafe {
3233 crate::Func::call_unchecked_raw(
3234 &mut store,
3235 start.as_non_null(),
3236 NonNull::new(
3237 &mut src[..count.max(dst.len())] as *mut [MaybeUninit<ValRaw>] as _,
3238 )
3239 .unwrap(),
3240 UncaughtException::Trap,
3241 )?;
3242 }
3243 dst.copy_from_slice(&src[..dst.len()]);
3244 let task = store.0.current_guest_thread()?.task;
3245 let state = store.0.concurrent_state_mut()?;
3246 Waitable::Guest(task).set_event(
3247 state,
3248 Some(Event::Subtask {
3249 status: Status::Started,
3250 }),
3251 )?;
3252 Ok(())
3253 }),
3254 LiftResult {
3255 lift: Box::new(move |store, src| {
3256 let mut store = token.as_context_mut(store);
3259 let mut my_src = src.to_owned(); if let ResultInfo::Heap { results } = &result_info {
3261 my_src.push(ValRaw::u32(*results));
3262 }
3263
3264 unsafe {
3271 crate::Func::call_unchecked_raw(
3272 &mut store,
3273 return_.as_non_null(),
3274 my_src.as_mut_slice().into(),
3275 UncaughtException::Trap,
3276 )?;
3277 }
3278
3279 let thread = store.0.current_guest_thread()?;
3280 let state = store.0.concurrent_state_mut()?;
3281 if sync_caller {
3282 state.get_mut(thread.task)?.sync_result = SyncResult::Produced(
3283 if let ResultInfo::Stack { result_count } = &result_info {
3284 match result_count {
3285 0 => None,
3286 1 => Some(my_src[0]),
3287 _ => unreachable!(),
3288 }
3289 } else {
3290 None
3291 },
3292 );
3293 }
3294 Ok(Box::new(DummyResult) as Box<dyn Any + Send + Sync>)
3295 }),
3296 ty: task_return_type,
3297 memory: NonNull::new(memory).map(SendSyncPtr::new),
3298 string_encoding,
3299 },
3300 Caller::Guest { thread: old_thread },
3301 None,
3302 self.runtime_instance(callee_instance),
3303 callee_async_typed,
3304 false,
3307 )?;
3308
3309 store.0.set_thread(guest_thread)?;
3312 log::trace!("pushed {guest_thread:?} as current thread; old thread was {old_thread:?}");
3313
3314 Ok(())
3315 }
3316
3317 unsafe fn call_callback<T>(
3322 self,
3323 mut store: StoreContextMut<T>,
3324 function: SendSyncPtr<VMFuncRef>,
3325 event: Event,
3326 handle: u32,
3327 ) -> Result<u32> {
3328 let (ordinal, result) = event.parts();
3329 let params = &mut [
3330 ValRaw::u32(ordinal),
3331 ValRaw::u32(handle),
3332 ValRaw::u32(result),
3333 ];
3334 unsafe {
3339 crate::Func::call_unchecked_raw(
3340 &mut store,
3341 function.as_non_null(),
3342 params.as_mut_slice().into(),
3343 UncaughtException::Trap,
3344 )?;
3345 }
3346 Ok(params[0].get_u32())
3347 }
3348
3349 unsafe fn start_call<T: 'static>(
3362 self,
3363 mut store: StoreContextMut<T>,
3364 callback: *mut VMFuncRef,
3365 post_return: *mut VMFuncRef,
3366 callee: NonNull<VMFuncRef>,
3367 param_count: u32,
3368 result_count: u32,
3369 flags: u32,
3370 storage: Option<&mut [MaybeUninit<ValRaw>]>,
3371 ) -> Result<u32> {
3372 let token = StoreToken::new(store.as_context_mut());
3373 let async_caller = storage.is_none();
3374 let guest_thread = store.0.current_guest_thread()?;
3375 let state = store.0.concurrent_state_mut()?;
3376
3377 if !state.event_loop_running {
3378 bail_bug!("Instance::start_call called without a running event loop");
3379 }
3380
3381 let callee = SendSyncPtr::new(callee);
3382 let param_count = usize::try_from(param_count)?;
3383 assert!(param_count <= MAX_FLAT_PARAMS);
3384 let result_count = usize::try_from(result_count)?;
3385 assert!(result_count <= MAX_FLAT_RESULTS);
3386
3387 let task = state.get_mut(guest_thread.task)?;
3388 let callee_async_typed = task.async_typed;
3389 let callee_instance = task.instance;
3390
3391 task.async_lifted = (flags & START_FLAG_ASYNC_CALLEE) != 0;
3392
3393 if let Some(callback) = NonNull::new(callback) {
3394 let callback = SendSyncPtr::new(callback);
3398 task.callback = Some(Box::new(move |store, event, handle| {
3399 let store = token.as_context_mut(store);
3400 unsafe { self.call_callback::<T>(store, callback, event, handle) }
3401 }));
3402 }
3403
3404 let Caller::Guest { thread: caller } = &task.caller else {
3405 bail_bug!("start_call unexpectedly invoked for host->guest call");
3408 };
3409 let caller = *caller;
3410 let caller_instance = state.get_mut(caller.task)?.instance;
3411
3412 unsafe {
3414 self.stage_call(
3415 store.as_context_mut(),
3416 guest_thread,
3417 callee,
3418 param_count,
3419 result_count,
3420 (flags & START_FLAG_ASYNC_CALLEE) != 0,
3421 NonNull::new(callback).map(SendSyncPtr::new),
3422 NonNull::new(post_return).map(SendSyncPtr::new),
3423 false,
3424 )?;
3425 }
3426
3427 let old_do_not_suspend = if callee_async_typed {
3428 let state = store.0.instance_state(callee_instance).concurrent_state();
3435 let old_do_not_suspend = state.do_not_suspend;
3436 state.do_not_suspend = false;
3437 Some(old_do_not_suspend)
3438 } else {
3439 None
3440 };
3441
3442 let state = store.0.concurrent_state_mut()?;
3443
3444 let guest_waitable = Waitable::Guest(guest_thread.task);
3447 let old_set = guest_waitable.common(state)?.set;
3448 let set = state.get_mut(caller.thread)?.sync_call_set;
3449 guest_waitable.join(state, Some(set))?;
3450
3451 store.0.set_thread(CurrentThread::None)?;
3452
3453 let mut yielded = false;
3469 let (status, waitable) = loop {
3470 store.0.suspend(if yielded {
3471 SuspendReason::Waiting {
3472 set,
3473 thread: caller,
3474 }
3475 } else {
3476 yielded = true;
3477 SuspendReason::YieldingToSubtask { thread: caller }
3478 })?;
3479
3480 if let Some(old_do_not_suspend) = old_do_not_suspend {
3481 store
3482 .0
3483 .instance_state(callee_instance)
3484 .concurrent_state()
3485 .do_not_suspend = old_do_not_suspend;
3486 }
3487
3488 let state = store.0.concurrent_state_mut()?;
3489
3490 log::trace!("taking event for {:?}", guest_thread.task);
3491 let event = guest_waitable.take_event(state)?;
3492 let Some(Event::Subtask { status }) = event else {
3493 bail_bug!("subtasks should only get subtask events, got {event:?}")
3494 };
3495
3496 log::trace!("status {status:?} for {:?}", guest_thread.task);
3497
3498 if status == Status::Returned {
3499 break (status, None);
3501 } else if async_caller {
3502 let handle = store
3506 .0
3507 .instance_state(caller_instance)
3508 .handle_table()
3509 .subtask_insert_guest(guest_thread.task.rep())?;
3510 store
3511 .0
3512 .concurrent_state_mut()?
3513 .get_mut(guest_thread.task)?
3514 .common
3515 .handle = Some(handle);
3516 break (status, Some(handle));
3517 } else {
3518 store.0.switch_or_trap_if_may_not_suspend(caller_instance)?;
3522 }
3523 };
3524
3525 guest_waitable.join(store.0.concurrent_state_mut()?, old_set)?;
3526
3527 store.0.set_thread(caller)?;
3529 store
3530 .0
3531 .concurrent_state_mut()?
3532 .get_mut(caller.thread)?
3533 .state = GuestThreadState::Running;
3534 log::trace!("popped current thread {guest_thread:?}; new thread is {caller:?}");
3535
3536 if let Some(storage) = storage {
3537 let state = store.0.concurrent_state_mut()?;
3541 let task = state.get_mut(guest_thread.task)?;
3542 if let Some(result) = task.sync_result.take()? {
3543 if let Some(result) = result {
3544 storage[0] = MaybeUninit::new(result);
3545 }
3546
3547 if task.exited && task.ready_to_delete() {
3548 Waitable::Guest(guest_thread.task).delete_from(store.0)?;
3549 }
3550 }
3551 }
3552
3553 Ok(status.pack(waitable))
3554 }
3555
3556 pub(crate) fn first_poll<T: 'static, R: Send + 'static>(
3572 self,
3573 mut store: StoreContextMut<'_, T>,
3574 host_task: EnteredHostTask,
3575 result_may_require_realloc: bool,
3576 future: impl Future<Output = Result<R>> + Send + 'static,
3577 lower: impl FnOnce(StoreContextMut<T>, Option<R>, bool, Option<TableId<HostTask>>) -> Result<()>
3578 + Send
3579 + 'static,
3580 ) -> Result<u32> {
3581 let token = StoreToken::new(store.as_context_mut());
3582
3583 let (join_handle, future) = JoinHandle::run(future);
3586 let mut future = Box::pin(future);
3587
3588 let poll = tls::set(store.0, || {
3593 future
3594 .as_mut()
3595 .poll(&mut Context::from_waker(&Waker::noop()))
3596 });
3597
3598 match poll {
3599 Poll::Ready(result) => {
3601 let result = result.transpose()?;
3602 let task = store.0.current_materialized_host_task()?;
3605 lower(store.as_context_mut(), result, true, task)?;
3606 return Ok(Status::Returned.pack(None));
3607 }
3608
3609 Poll::Pending => {}
3611 }
3612
3613 let Some(task) = store.0.materialize_host_task_id()? else {
3617 bail_bug!("current thread is not a host thread")
3618 };
3619 {
3620 let state = &mut store.0.concurrent_state_mut()?.get_mut(task)?.state;
3621 assert!(matches!(state, HostTaskState::CalleeStarted));
3622 *state = HostTaskState::CalleeRunning(join_handle);
3623 }
3624
3625 let future = Box::pin(async move {
3633 let result = match run_with_host_task_set(task, future).await? {
3634 Some(result) => Some(result?),
3635 None => None,
3636 };
3637 let on_complete = move |store: &mut dyn VMStore| {
3638 let mut store = token.as_context_mut(store);
3642 let old = store.0.set_thread(task)?;
3643
3644 let status = if result.is_some() {
3645 Status::Returned
3646 } else {
3647 Status::ReturnCancelled
3648 };
3649
3650 lower(store.as_context_mut(), result, false, Some(task))?;
3651 let state = store.0.concurrent_state_mut()?;
3652 match &mut state.get_mut(task)?.state {
3653 pending @ HostTaskState::CalleeCancelling => {
3656 *pending = HostTaskState::CalleeDone { cancelled: true };
3657 }
3658
3659 other => *other = HostTaskState::CalleeDone { cancelled: false },
3661 }
3662 Waitable::Host(task).set_event(state, Some(Event::Subtask { status }))?;
3663
3664 store.0.set_thread(old)?;
3665 Ok(())
3666 };
3667
3668 tls::get(move |store| {
3669 if result_may_require_realloc {
3670 store
3675 .concurrent_state_mut()?
3676 .push_high_priority(WorkItem::WorkerFunction(AlwaysMut::new(Box::new(
3677 on_complete,
3678 ))));
3679 Ok(())
3680 } else {
3681 on_complete(store)
3684 }
3685 })
3686 });
3687
3688 let caller = match host_task {
3691 Some(caller) => caller,
3692 None => bail_bug!("host task wasn't created but should have been"),
3693 };
3694 let state = store.0.concurrent_state_mut()?;
3695 state.push_future(future);
3696 let instance = state.get_mut(caller.task)?.instance;
3697 let handle = store
3698 .0
3699 .instance_state(instance)
3700 .handle_table()
3701 .subtask_insert_host(task.rep())?;
3702 store.0.concurrent_state_mut()?.get_mut(task)?.common.handle = Some(handle);
3703 log::trace!("assign {task:?} handle {handle} for {caller:?} instance {instance:?}");
3704
3705 store.0.set_thread(caller)?;
3709 Ok(Status::Started.pack(Some(handle)))
3710 }
3711
3712 pub(crate) fn task_return(
3715 self,
3716 store: &mut dyn VMStore,
3717 ty: TypeTupleIndex,
3718 options: OptionsIndex,
3719 storage: &[ValRaw],
3720 ) -> Result<()> {
3721 let guest_thread = store.current_guest_thread()?;
3722 let state = store.concurrent_state_mut()?;
3723 if !state.get_mut(guest_thread.task)?.async_lifted {
3724 bail!(Trap::TaskReturnOrCancelSyncLifted);
3725 }
3726 let lift = state
3727 .get_mut(guest_thread.task)?
3728 .lift_result
3729 .take()
3730 .ok_or_else(|| Trap::TaskCancelOrReturnTwice)?;
3731 if !state.get_mut(guest_thread.task)?.result.is_none() {
3732 bail_bug!("task result unexpectedly already set");
3733 }
3734
3735 let CanonicalOptions {
3736 string_encoding,
3737 data_model,
3738 ..
3739 } = &self.id().get(store).component().env_component().options[options];
3740
3741 let invalid = ty != lift.ty
3742 || string_encoding != &lift.string_encoding
3743 || match data_model {
3744 CanonicalOptionsDataModel::LinearMemory(opts) => match opts.memory {
3745 Some(memory) => {
3746 let expected = lift.memory.map(|v| v.as_ptr()).unwrap_or(ptr::null_mut());
3747 let actual = self.id().get(store).runtime_memory(memory);
3748 expected != actual.as_ptr()
3749 }
3750 None => false,
3753 },
3754 CanonicalOptionsDataModel::Gc { .. } => true,
3756 };
3757
3758 if invalid {
3759 bail!(Trap::TaskReturnInvalid);
3760 }
3761
3762 log::trace!("task.return for {guest_thread:?}");
3763
3764 let result = (lift.lift)(store, storage)?;
3765 self.task_complete(store, guest_thread.task, result, Status::Returned)
3766 }
3767
3768 pub(crate) fn task_cancel(self, store: &mut StoreOpaque) -> Result<()> {
3770 let guest_thread = store.current_guest_thread()?;
3771 let state = store.concurrent_state_mut()?;
3772 let task = state.get_mut(guest_thread.task)?;
3773 if !task.async_lifted {
3774 bail!(Trap::TaskReturnOrCancelSyncLifted);
3775 }
3776 if !task.cancel_request_delivered {
3777 bail!(Trap::TaskCancelNotCancelled);
3778 }
3779 _ = task
3780 .lift_result
3781 .take()
3782 .ok_or_else(|| Trap::TaskCancelOrReturnTwice)?;
3783
3784 if !task.result.is_none() {
3785 bail_bug!("task result should not bet set yet");
3786 }
3787
3788 log::trace!("task.cancel for {guest_thread:?}");
3789
3790 self.task_complete(
3791 store,
3792 guest_thread.task,
3793 Box::new(DummyResult),
3794 Status::ReturnCancelled,
3795 )
3796 }
3797
3798 fn task_complete(
3804 self,
3805 store: &mut StoreOpaque,
3806 guest_task: TableId<GuestTask>,
3807 result: Box<dyn Any + Send + Sync>,
3808 status: Status,
3809 ) -> Result<()> {
3810 store
3811 .component_resource_tables(Some(self))?
3812 .validate_scope_exit()?;
3813
3814 let state = store.concurrent_state_mut()?;
3815 let task = state.get_mut(guest_task)?;
3816
3817 task.event = None;
3821
3822 if let Caller::Host { tx, .. } = &mut task.caller {
3823 if let Some(tx) = tx.take() {
3824 _ = tx.send(result);
3825 }
3826 } else {
3827 task.result = Some(result);
3828 Waitable::Guest(guest_task).set_event(state, Some(Event::Subtask { status }))?;
3829 }
3830
3831 Ok(())
3832 }
3833
3834 pub(crate) fn waitable_set_new(
3836 self,
3837 store: &mut StoreOpaque,
3838 caller_instance: RuntimeComponentInstanceIndex,
3839 ) -> Result<u32> {
3840 let set = store.concurrent_state_mut()?.push(WaitableSet::default())?;
3841 let handle = store
3842 .instance_state(self.runtime_instance(caller_instance))
3843 .handle_table()
3844 .waitable_set_insert(set.rep())?;
3845 log::trace!("new waitable set {set:?} (handle {handle})");
3846 Ok(handle)
3847 }
3848
3849 pub(crate) fn waitable_set_drop(
3851 self,
3852 store: &mut StoreOpaque,
3853 caller_instance: RuntimeComponentInstanceIndex,
3854 set: u32,
3855 ) -> Result<()> {
3856 let rep = store
3857 .instance_state(self.runtime_instance(caller_instance))
3858 .handle_table()
3859 .waitable_set_remove(set)?;
3860
3861 log::trace!("drop waitable set {rep} (handle {set})");
3862
3863 let set = store
3867 .concurrent_state_mut()?
3868 .get_mut(TableId::<WaitableSet>::new(rep))?;
3869 if set.num_waiting > 0 {
3870 bail!(Trap::WaitableSetDropHasWaiters);
3871 }
3872
3873 store
3874 .concurrent_state_mut()?
3875 .delete(TableId::<WaitableSet>::new(rep))?;
3876
3877 Ok(())
3878 }
3879
3880 pub(crate) fn waitable_join(
3882 self,
3883 store: &mut StoreOpaque,
3884 caller_instance: RuntimeComponentInstanceIndex,
3885 waitable_handle: u32,
3886 set_handle: u32,
3887 ) -> Result<()> {
3888 let mut instance = self.id().get_mut(store);
3889 let waitable =
3890 Waitable::from_instance(instance.as_mut(), caller_instance, waitable_handle)?;
3891
3892 let set = if set_handle == 0 {
3893 None
3894 } else {
3895 let set = instance.instance_states().0[caller_instance]
3896 .handle_table()
3897 .waitable_set_rep(set_handle)?;
3898
3899 let state = store.concurrent_state_mut()?;
3900 if let Some(old) = waitable.common(state)?.set
3901 && state.get_mut(old)?.is_sync_call_set
3902 {
3903 bail!(Trap::WaitableSyncAndAsync);
3904 }
3905
3906 Some(TableId::<WaitableSet>::new(set))
3907 };
3908
3909 log::trace!(
3910 "waitable {waitable:?} (handle {waitable_handle}) join set {set:?} (handle {set_handle})",
3911 );
3912
3913 waitable.join(store.concurrent_state_mut()?, set)
3914 }
3915
3916 pub(crate) fn subtask_drop(
3918 self,
3919 store: &mut StoreOpaque,
3920 caller_instance: RuntimeComponentInstanceIndex,
3921 task_id: u32,
3922 ) -> Result<()> {
3923 self.waitable_join(store, caller_instance, task_id, 0)?;
3924
3925 let (rep, is_host) = store
3926 .instance_state(self.runtime_instance(caller_instance))
3927 .handle_table()
3928 .subtask_remove(task_id)?;
3929
3930 let concurrent_state = store.concurrent_state_mut()?;
3931 let (waitable, delete) = if is_host {
3932 let id = TableId::<HostTask>::new(rep);
3933 let task = concurrent_state.get_mut(id)?;
3934 match &task.state {
3935 HostTaskState::CalleeRunning(_) | HostTaskState::CalleeCancelling => {
3936 bail!(Trap::SubtaskDropNotResolved)
3937 }
3938 HostTaskState::CalleeDone { .. } => {}
3939 HostTaskState::CalleeStarted | HostTaskState::CalleeFinished(_) => {
3940 bail_bug!("invalid state for callee in `subtask.drop`")
3941 }
3942 }
3943
3944 (Waitable::Host(id), true)
3945 } else {
3946 let id = TableId::<GuestTask>::new(rep);
3947 let task = concurrent_state.get_mut(id)?;
3948 if task.lift_result.is_some() {
3949 bail!(Trap::SubtaskDropNotResolved);
3950 }
3951 (
3952 Waitable::Guest(id),
3953 concurrent_state.get_mut(id)?.ready_to_delete(),
3954 )
3955 };
3956
3957 waitable.common(concurrent_state)?.handle = None;
3958
3959 if waitable.take_event(concurrent_state)?.is_some() {
3962 bail!(Trap::SubtaskDropNotResolved);
3963 }
3964
3965 if delete {
3966 waitable.delete_from(store)?;
3967 }
3968
3969 log::trace!("subtask_drop {waitable:?} (handle {task_id})");
3970 Ok(())
3971 }
3972
3973 pub(crate) fn waitable_set_wait(
3975 self,
3976 store: &mut StoreOpaque,
3977 options: OptionsIndex,
3978 set: u32,
3979 payload: u32,
3980 ) -> Result<u32> {
3981 let &CanonicalOptions {
3982 instance: caller_instance,
3983 ..
3984 } = &self.id().get(store).component().env_component().options[options];
3985 let caller = self.runtime_instance(caller_instance);
3986 let rep = store
3987 .instance_state(self.runtime_instance(caller_instance))
3988 .handle_table()
3989 .waitable_set_rep(set)?;
3990
3991 self.waitable_check(
3992 store,
3993 caller,
3994 WaitableCheck::Wait,
3995 WaitableCheckParams {
3996 set: TableId::new(rep),
3997 options,
3998 payload,
3999 },
4000 )
4001 }
4002
4003 pub(crate) fn waitable_set_poll(
4005 self,
4006 store: &mut StoreOpaque,
4007 options: OptionsIndex,
4008 set: u32,
4009 payload: u32,
4010 ) -> Result<u32> {
4011 let &CanonicalOptions {
4012 instance: caller_instance,
4013 ..
4014 } = &self.id().get(store).component().env_component().options[options];
4015 let caller = self.runtime_instance(caller_instance);
4016 let rep = store
4017 .instance_state(caller)
4018 .handle_table()
4019 .waitable_set_rep(set)?;
4020
4021 self.waitable_check(
4022 store,
4023 caller,
4024 WaitableCheck::Poll,
4025 WaitableCheckParams {
4026 set: TableId::new(rep),
4027 options,
4028 payload,
4029 },
4030 )
4031 }
4032
4033 pub(crate) fn thread_index(&self, store: &mut dyn VMStore) -> Result<u32> {
4035 let thread_id = store.current_guest_thread()?.thread;
4036 match store
4037 .concurrent_state_mut()?
4038 .get_mut(thread_id)?
4039 .instance_rep
4040 {
4041 Some(r) => Ok(r),
4042 None => bail_bug!("thread should have instance_rep by now"),
4043 }
4044 }
4045
4046 pub(crate) fn thread_new_indirect<T: 'static>(
4048 self,
4049 mut store: StoreContextMut<T>,
4050 runtime_instance: RuntimeComponentInstanceIndex,
4051 _func_ty_idx: TypeFuncIndex, start_func_table_idx: RuntimeTableIndex,
4053 start_func_idx: u32,
4054 context: i32,
4055 ) -> Result<u32> {
4056 log::trace!("creating new thread");
4057
4058 let start_func_ty = FuncType::new(store.engine(), [ValType::I32], []);
4059 let (instance, registry) = self.id().get_mut_and_registry(store.0);
4060 let callee = instance
4061 .index_runtime_func_table(registry, start_func_table_idx, start_func_idx as u64)?
4062 .ok_or_else(|| Trap::ThreadNewIndirectUninitialized)?;
4063 if callee.type_index(store.0) != start_func_ty.type_index() {
4064 bail!(Trap::ThreadNewIndirectInvalidType);
4065 }
4066
4067 let token = StoreToken::new(store.as_context_mut());
4068 let start_func = Box::new(
4069 move |store: &mut dyn VMStore, guest_thread: QualifiedThreadId| -> Result<()> {
4070 let old_thread = store.set_thread(guest_thread)?;
4071 log::trace!(
4072 "thread start: replaced {old_thread:?} with {guest_thread:?} as current thread"
4073 );
4074
4075 let mut store = token.as_context_mut(store);
4076 let mut params = [ValRaw::i32(context)];
4077 unsafe { callee.call_unchecked(store.as_context_mut(), &mut params)? };
4080
4081 store.0.set_thread(old_thread)?;
4082
4083 let runtime_instance = self.runtime_instance(runtime_instance);
4084
4085 store
4088 .0
4089 .switch_or_trap_if_may_not_suspend(runtime_instance)?;
4090
4091 store
4092 .0
4093 .cleanup_thread(guest_thread, runtime_instance, CleanupTask::Yes)?;
4094
4095 log::trace!("explicit thread {guest_thread:?} completed");
4096 let state = store.0.concurrent_state_mut()?;
4097 if let Some(t) = old_thread.guest() {
4098 state.get_mut(t.thread)?.state = GuestThreadState::Running;
4099 }
4100 log::trace!("thread start: restored {old_thread:?} as current thread");
4101
4102 Ok(())
4103 },
4104 );
4105
4106 let current_thread = store.0.current_guest_thread()?;
4107 let state = store.0.concurrent_state_mut()?;
4108 let parent_task = current_thread.task;
4109
4110 let new_thread = GuestThread::new_explicit(state, parent_task, start_func)?;
4111 let thread_id = state.push(new_thread)?;
4112 state.get_mut(parent_task)?.threads.insert(thread_id);
4113
4114 log::trace!("new thread with id {thread_id:?} created");
4115
4116 self.add_guest_thread_to_instance_table(thread_id, store.0, runtime_instance)
4117 }
4118
4119 pub(crate) fn resume_thread(
4120 self,
4121 store: &mut StoreOpaque,
4122 runtime_instance: RuntimeComponentInstanceIndex,
4123 thread_idx: u32,
4124 how: ResumeThread,
4125 ) -> Result<bool> {
4126 let thread_id =
4127 GuestThread::from_instance(self.id().get_mut(store), runtime_instance, thread_idx)?;
4128 let state = store.concurrent_state_mut()?;
4129 let guest_thread = QualifiedThreadId::qualify(state, thread_id)?;
4130
4131 if store.current_guest_thread()? == guest_thread {
4132 bail!(Trap::CannotResumeThread);
4133 }
4134
4135 let state = store.concurrent_state_mut()?;
4136 let thread = state.get_mut(guest_thread.thread)?;
4137 let priority = match how {
4138 ResumeThread::Promote | ResumeThread::Resume => Priority::Switch,
4139 ResumeThread::ResumeLater => Priority::Low,
4140 };
4141
4142 match (&how, &thread.state) {
4143 (ResumeThread::Promote, GuestThreadState::Ready { .. }) => {}
4145 (ResumeThread::Promote, _) => return Ok(false),
4146
4147 (
4150 ResumeThread::Resume | ResumeThread::ResumeLater,
4151 GuestThreadState::NotStartedExplicit(_) | GuestThreadState::Suspended(_),
4152 ) => {}
4153 (ResumeThread::Resume | ResumeThread::ResumeLater, _) => {
4154 bail!(Trap::CannotResumeThread)
4155 }
4156 }
4157
4158 match mem::replace(&mut thread.state, GuestThreadState::Running) {
4159 GuestThreadState::NotStartedExplicit(start_func) => {
4160 log::trace!("starting thread {guest_thread:?}");
4161 let guest_call = WorkItem::GuestCall {
4162 instance: self.runtime_instance(runtime_instance),
4163 call: GuestCall {
4164 thread: guest_thread,
4165 kind: GuestCallKind::StartExplicit(Box::new(move |store| {
4166 start_func(store, guest_thread)
4167 })),
4168 },
4169 };
4170 store
4171 .concurrent_state_mut()?
4172 .push_work_item(guest_call, priority)?;
4173 }
4174 GuestThreadState::Suspended(fiber) => {
4175 log::trace!("resuming thread {thread_id:?} that was suspended");
4176 store.concurrent_state_mut()?.push_work_item(
4177 WorkItem::ResumeFiber {
4178 instance: self.runtime_instance(runtime_instance),
4179 thread: guest_thread,
4180 fiber,
4181 },
4182 priority,
4183 )?;
4184 }
4185 GuestThreadState::Ready { fiber } => {
4186 log::trace!("resuming thread {thread_id:?} that was ready");
4187 thread.state = GuestThreadState::Ready { fiber };
4188 store
4189 .concurrent_state_mut()?
4190 .promote_thread_work_item(guest_thread)?;
4191 }
4192 other @ (GuestThreadState::NotStartedImplicit
4193 | GuestThreadState::Running
4194 | GuestThreadState::Completed) => {
4195 thread.state = other;
4196 }
4197 }
4198 Ok(true)
4199 }
4200
4201 fn add_guest_thread_to_instance_table(
4202 self,
4203 thread_id: TableId<GuestThread>,
4204 store: &mut StoreOpaque,
4205 runtime_instance: RuntimeComponentInstanceIndex,
4206 ) -> Result<u32> {
4207 let guest_id = store
4208 .instance_state(self.runtime_instance(runtime_instance))
4209 .thread_handle_table()
4210 .guest_thread_insert(thread_id.rep())?;
4211 store
4212 .concurrent_state_mut()?
4213 .get_mut(thread_id)?
4214 .instance_rep = Some(guest_id);
4215 Ok(guest_id)
4216 }
4217
4218 pub(crate) fn suspension_intrinsic(
4222 self,
4223 store: &mut StoreOpaque,
4224 caller: RuntimeComponentInstanceIndex,
4225 yielding: bool,
4226 to_thread: SuspensionTarget,
4227 ) -> Result<WaitResult> {
4228 let check_suspend = match to_thread {
4229 SuspensionTarget::Promote(thread) => {
4230 !self.resume_thread(store, caller, thread, ResumeThread::Promote)?
4231 }
4232 SuspensionTarget::Resume(thread) => {
4233 if !self.resume_thread(store, caller, thread, ResumeThread::Resume)? {
4234 bail_bug!(
4235 "`resume_thread` should only ever return false \
4236 when `ResumeThread::Promote` is passed to it"
4237 );
4238 }
4239 false
4240 }
4241 SuspensionTarget::None => true,
4242 };
4243
4244 if check_suspend && !store.switch_if_may_not_suspend(self.runtime_instance(caller))? {
4245 return if yielding {
4246 Ok(WaitResult::Completed)
4247 } else {
4248 Err(Trap::CannotBlockSyncTask.into())
4249 };
4250 }
4251
4252 let guest_thread = store.current_guest_thread()?;
4253
4254 let reason = if yielding {
4255 SuspendReason::Yielding {
4256 thread: guest_thread,
4257 }
4258 } else {
4259 SuspendReason::ExplicitlySuspending {
4260 thread: guest_thread,
4261 }
4262 };
4263
4264 store.suspend(reason)?;
4265
4266 Ok(WaitResult::Completed)
4267 }
4268
4269 fn waitable_check(
4271 self,
4272 store: &mut StoreOpaque,
4273 caller: RuntimeInstance,
4274 check: WaitableCheck,
4275 params: WaitableCheckParams,
4276 ) -> Result<u32> {
4277 let guest_thread = store.current_guest_thread()?;
4278
4279 log::trace!("waitable check for {guest_thread:?}; set {:?}", params.set);
4280
4281 match &check {
4284 WaitableCheck::Wait => {
4285 let set = params.set;
4286
4287 loop {
4292 let state = store.concurrent_state_mut()?;
4293 let task = state.get_mut(guest_thread.task)?;
4294 if !(task.event.is_none() || matches!(task.event, Some(Event::Cancelled)))
4295 || !state.get_mut(set)?.ready.is_empty()
4296 {
4297 break;
4298 }
4299
4300 store.switch_or_trap_if_may_not_suspend(caller)?;
4301
4302 store.suspend(SuspendReason::Waiting {
4303 set,
4304 thread: guest_thread,
4305 })?;
4306 }
4307 }
4308 WaitableCheck::Poll => {}
4309 }
4310
4311 log::trace!(
4312 "waitable check for {guest_thread:?}; set {:?}, part two",
4313 params.set
4314 );
4315
4316 let event = self.get_event(store, guest_thread.task, Some(params.set), false)?;
4318
4319 let (ordinal, handle, result) = match &check {
4320 WaitableCheck::Wait => {
4321 let (event, waitable) = match event {
4322 Some(p) => p,
4323 None => bail_bug!("event expected to be present"),
4324 };
4325 let handle = waitable.map(|(_, v)| v).unwrap_or(0);
4326 let (ordinal, result) = event.parts();
4327 (ordinal, handle, result)
4328 }
4329 WaitableCheck::Poll => {
4330 if let Some((event, waitable)) = event {
4331 let handle = waitable.map(|(_, v)| v).unwrap_or(0);
4332 let (ordinal, result) = event.parts();
4333 (ordinal, handle, result)
4334 } else {
4335 log::trace!(
4336 "no events ready to deliver via waitable-set.poll to {:?}; set {:?}",
4337 guest_thread.task,
4338 params.set
4339 );
4340 let (ordinal, result) = Event::None.parts();
4341 (ordinal, 0, result)
4342 }
4343 }
4344 };
4345 let memory = self.options_memory_mut(store, params.options);
4346 let ptr = crate::component::func::validate_inbounds_dynamic(
4347 &CanonicalAbiInfo::POINTER_PAIR,
4348 memory,
4349 &ValRaw::u32(params.payload),
4350 )?;
4351 memory[ptr + 0..][..4].copy_from_slice(&handle.to_le_bytes());
4352 memory[ptr + 4..][..4].copy_from_slice(&result.to_le_bytes());
4353 Ok(ordinal)
4354 }
4355
4356 pub(crate) fn subtask_cancel(
4358 self,
4359 store: &mut StoreOpaque,
4360 caller_instance: RuntimeComponentInstanceIndex,
4361 async_: bool,
4362 task_id: u32,
4363 ) -> Result<u32> {
4364 let (rep, is_host) = store
4365 .instance_state(self.runtime_instance(caller_instance))
4366 .handle_table()
4367 .subtask_rep(task_id)?;
4368 let waitable = if is_host {
4369 Waitable::Host(TableId::<HostTask>::new(rep))
4370 } else {
4371 Waitable::Guest(TableId::<GuestTask>::new(rep))
4372 };
4373 let concurrent_state = store.concurrent_state_mut()?;
4374
4375 log::trace!("subtask_cancel {waitable:?} (handle {task_id}; async {async_})");
4376
4377 waitable.trap_if_in_waitable_set(concurrent_state)?;
4378
4379 let needs_block;
4380 if let Waitable::Host(host_task) = waitable {
4381 let state = &mut concurrent_state.get_mut(host_task)?.state;
4382 match state {
4383 HostTaskState::CalleeRunning(handle) => {
4390 handle.abort();
4391 *state = HostTaskState::CalleeCancelling;
4392 needs_block = true;
4393 }
4394
4395 HostTaskState::CalleeCancelling | HostTaskState::CalleeDone { cancelled: true } => {
4398 bail!(Trap::SubtaskCancelAfterTerminal);
4399 }
4400 HostTaskState::CalleeDone { cancelled: false } => {
4401 *state = HostTaskState::CalleeDone { cancelled: true };
4404 needs_block = false;
4405 }
4406
4407 HostTaskState::CalleeStarted | HostTaskState::CalleeFinished(_) => {
4410 bail_bug!("invalid states for host callee")
4411 }
4412 }
4413 } else {
4414 let guest_task = TableId::<GuestTask>::new(rep);
4415 let task = concurrent_state.get_mut(guest_task)?;
4416 if !task.already_lowered_parameters() {
4417 store.cancel_guest_subtask_without_lowered_parameters(
4418 self.runtime_instance(caller_instance),
4419 guest_task,
4420 )?;
4421 return Ok(Status::StartCancelled as u32);
4422 } else if !task.returned_or_cancelled() {
4423 task.event = Some(Event::Cancelled);
4431 let runtime_instance = task.instance;
4432 for thread in task.threads.clone() {
4433 let thread = QualifiedThreadId {
4434 task: guest_task,
4435 thread,
4436 };
4437 let concurrent_state = store.concurrent_state_mut()?;
4438 let thread_mut = concurrent_state.get_mut(thread.thread)?;
4439
4440 let yield_ = |store: &mut StoreOpaque| {
4441 let state = store.instance_state(runtime_instance).concurrent_state();
4446 let old_do_not_suspend = state.do_not_suspend;
4447 state.do_not_suspend = false;
4448
4449 let caller = store.current_guest_thread()?;
4450
4451 let state = store.concurrent_state_mut()?;
4456 let set = state.get_mut(caller.thread)?.sync_call_set;
4457 waitable.join(state, Some(set))?;
4458
4459 store.suspend(SuspendReason::YieldingToSubtask { thread: caller })?;
4460
4461 let state = store.concurrent_state_mut()?;
4462 waitable.join(state, None)?;
4463
4464 store
4465 .instance_state(runtime_instance)
4466 .concurrent_state()
4467 .do_not_suspend = old_do_not_suspend;
4468
4469 Ok::<(), crate::Error>(())
4470 };
4471
4472 match thread_mut.wake_on_cancel.take() {
4473 WakeOnCancel::Waiting(set) => {
4474 let item = match concurrent_state.get_mut(set)?.waiting.remove(&thread)
4476 {
4477 Some(WaitMode::Callback(instance)) => WorkItem::GuestCall {
4478 instance: runtime_instance,
4479 call: GuestCall {
4480 thread,
4481 kind: GuestCallKind::DeliverEvent {
4482 instance,
4483 set: Some(set),
4484 },
4485 },
4486 },
4487 other => bail_bug!(
4488 "expected `Some(WaitMode::Callback(_))`; got `{other:?}`"
4489 ),
4490 };
4491 concurrent_state.set_switch_item(item)?;
4492
4493 yield_(store)?;
4494
4495 break;
4496 }
4497 WakeOnCancel::Yielding => {
4498 if concurrent_state.promote_thread_work_item(thread)? {
4499 yield_(store)?;
4500 break;
4501 } else if store
4502 .instance_state(runtime_instance)
4503 .concurrent_state()
4504 .pending
4505 .contains_key(&thread)
4506 {
4507 store
4517 .concurrent_state_mut()?
4518 .get_mut(thread.thread)?
4519 .wake_on_cancel = WakeOnCancel::Yielding;
4520 } else {
4521 bail_bug!("thread with `WakeOnCancel::Yielding` not promotable");
4522 }
4523 }
4524 WakeOnCancel::None => {}
4525 }
4526 }
4527
4528 needs_block = !store
4531 .concurrent_state_mut()?
4532 .get_mut(guest_task)?
4533 .returned_or_cancelled()
4534 } else {
4535 needs_block = false;
4536 }
4537 };
4538
4539 if needs_block {
4543 if async_ {
4544 return Ok(BLOCKED);
4545 }
4546
4547 let old_next_switch_item = {
4550 let state = store.concurrent_state_mut()?;
4551 let item = state.next_switch_item.take();
4552 state.push(item)?
4556 };
4557
4558 store.wait_for_event(self.runtime_instance(caller_instance), waitable)?;
4561
4562 let state = store.concurrent_state_mut()?;
4563 state.next_switch_item = state.delete(old_next_switch_item)?;
4564
4565 }
4567
4568 let event = waitable.take_event(store.concurrent_state_mut()?)?;
4569 if let Some(Event::Subtask {
4570 status: status @ (Status::Returned | Status::ReturnCancelled),
4571 }) = event
4572 {
4573 Ok(status as u32)
4574 } else {
4575 bail!(Trap::SubtaskCancelAfterTerminal);
4576 }
4577 }
4578}
4579
4580pub trait VMComponentAsyncStore {
4588 unsafe fn prepare_call(
4594 &mut self,
4595 instance: Instance,
4596 memory: *mut VMMemoryDefinition,
4597 start: NonNull<VMFuncRef>,
4598 return_: NonNull<VMFuncRef>,
4599 caller_instance: RuntimeComponentInstanceIndex,
4600 callee_instance: RuntimeComponentInstanceIndex,
4601 task_return_type: TypeTupleIndex,
4602 callee_async: bool,
4603 string_encoding: StringEncoding,
4604 result_count: u32,
4605 storage: *mut ValRaw,
4606 storage_len: usize,
4607 ) -> Result<()>;
4608
4609 unsafe fn sync_start(
4612 &mut self,
4613 instance: Instance,
4614 callback: *mut VMFuncRef,
4615 callee: NonNull<VMFuncRef>,
4616 param_count: u32,
4617 storage: *mut MaybeUninit<ValRaw>,
4618 storage_len: usize,
4619 ) -> Result<()>;
4620
4621 unsafe fn async_start(
4624 &mut self,
4625 instance: Instance,
4626 callback: *mut VMFuncRef,
4627 post_return: *mut VMFuncRef,
4628 callee: NonNull<VMFuncRef>,
4629 param_count: u32,
4630 result_count: u32,
4631 flags: u32,
4632 ) -> Result<u32>;
4633
4634 fn future_write(
4636 &mut self,
4637 instance: Instance,
4638 caller: RuntimeComponentInstanceIndex,
4639 ty: TypeFutureTableIndex,
4640 options: OptionsIndex,
4641 future: u32,
4642 address: u32,
4643 ) -> Result<u32>;
4644
4645 fn future_read(
4647 &mut self,
4648 instance: Instance,
4649 caller: RuntimeComponentInstanceIndex,
4650 ty: TypeFutureTableIndex,
4651 options: OptionsIndex,
4652 future: u32,
4653 address: u32,
4654 ) -> Result<u32>;
4655
4656 fn future_drop_writable(
4658 &mut self,
4659 instance: Instance,
4660 ty: TypeFutureTableIndex,
4661 writer: u32,
4662 ) -> Result<()>;
4663
4664 fn stream_write(
4666 &mut self,
4667 instance: Instance,
4668 caller: RuntimeComponentInstanceIndex,
4669 ty: TypeStreamTableIndex,
4670 options: OptionsIndex,
4671 stream: u32,
4672 address: u32,
4673 count: u32,
4674 ) -> Result<u32>;
4675
4676 fn stream_read(
4678 &mut self,
4679 instance: Instance,
4680 caller: RuntimeComponentInstanceIndex,
4681 ty: TypeStreamTableIndex,
4682 options: OptionsIndex,
4683 stream: u32,
4684 address: u32,
4685 count: u32,
4686 ) -> Result<u32>;
4687
4688 fn flat_stream_write(
4691 &mut self,
4692 instance: Instance,
4693 caller: RuntimeComponentInstanceIndex,
4694 ty: TypeStreamTableIndex,
4695 options: OptionsIndex,
4696 payload_size: u32,
4697 payload_align: u32,
4698 stream: u32,
4699 address: u32,
4700 count: u32,
4701 ) -> Result<u32>;
4702
4703 fn flat_stream_read(
4706 &mut self,
4707 instance: Instance,
4708 caller: RuntimeComponentInstanceIndex,
4709 ty: TypeStreamTableIndex,
4710 options: OptionsIndex,
4711 payload_size: u32,
4712 payload_align: u32,
4713 stream: u32,
4714 address: u32,
4715 count: u32,
4716 ) -> Result<u32>;
4717
4718 fn stream_drop_writable(
4720 &mut self,
4721 instance: Instance,
4722 ty: TypeStreamTableIndex,
4723 writer: u32,
4724 ) -> Result<()>;
4725
4726 fn error_context_debug_message(
4728 &mut self,
4729 instance: Instance,
4730 ty: TypeComponentLocalErrorContextTableIndex,
4731 options: OptionsIndex,
4732 err_ctx_handle: u32,
4733 debug_msg_address: u32,
4734 ) -> Result<()>;
4735
4736 fn thread_new_indirect(
4738 &mut self,
4739 instance: Instance,
4740 caller: RuntimeComponentInstanceIndex,
4741 func_ty_idx: TypeFuncIndex,
4742 start_func_table_idx: RuntimeTableIndex,
4743 start_func_idx: u32,
4744 context: i32,
4745 ) -> Result<u32>;
4746}
4747
4748impl<T: 'static> VMComponentAsyncStore for StoreInner<T> {
4750 unsafe fn prepare_call(
4751 &mut self,
4752 instance: Instance,
4753 memory: *mut VMMemoryDefinition,
4754 start: NonNull<VMFuncRef>,
4755 return_: NonNull<VMFuncRef>,
4756 caller_instance: RuntimeComponentInstanceIndex,
4757 callee_instance: RuntimeComponentInstanceIndex,
4758 task_return_type: TypeTupleIndex,
4759 callee_async: bool,
4760 string_encoding: StringEncoding,
4761 result_count_or_max_if_async: u32,
4762 storage: *mut ValRaw,
4763 storage_len: usize,
4764 ) -> Result<()> {
4765 let params = unsafe { core::slice::from_raw_parts(storage, storage_len) }.to_vec();
4769
4770 unsafe {
4771 instance.prepare_call(
4772 StoreContextMut(self),
4773 start,
4774 return_,
4775 caller_instance,
4776 callee_instance,
4777 task_return_type,
4778 callee_async,
4779 memory,
4780 string_encoding,
4781 match result_count_or_max_if_async {
4782 PREPARE_ASYNC_NO_RESULT => CallerInfo::Async {
4783 params,
4784 has_result: false,
4785 },
4786 PREPARE_ASYNC_WITH_RESULT => CallerInfo::Async {
4787 params,
4788 has_result: true,
4789 },
4790 result_count => CallerInfo::Sync {
4791 params,
4792 result_count,
4793 },
4794 },
4795 )
4796 }
4797 }
4798
4799 unsafe fn sync_start(
4800 &mut self,
4801 instance: Instance,
4802 callback: *mut VMFuncRef,
4803 callee: NonNull<VMFuncRef>,
4804 param_count: u32,
4805 storage: *mut MaybeUninit<ValRaw>,
4806 storage_len: usize,
4807 ) -> Result<()> {
4808 unsafe {
4809 instance
4810 .start_call(
4811 StoreContextMut(self),
4812 callback,
4813 ptr::null_mut(),
4814 callee,
4815 param_count,
4816 1,
4817 START_FLAG_ASYNC_CALLEE,
4818 Some(core::slice::from_raw_parts_mut(storage, storage_len)),
4822 )
4823 .map(drop)
4824 }
4825 }
4826
4827 unsafe fn async_start(
4828 &mut self,
4829 instance: Instance,
4830 callback: *mut VMFuncRef,
4831 post_return: *mut VMFuncRef,
4832 callee: NonNull<VMFuncRef>,
4833 param_count: u32,
4834 result_count: u32,
4835 flags: u32,
4836 ) -> Result<u32> {
4837 unsafe {
4838 instance.start_call(
4839 StoreContextMut(self),
4840 callback,
4841 post_return,
4842 callee,
4843 param_count,
4844 result_count,
4845 flags,
4846 None,
4847 )
4848 }
4849 }
4850
4851 fn future_write(
4852 &mut self,
4853 instance: Instance,
4854 caller: RuntimeComponentInstanceIndex,
4855 ty: TypeFutureTableIndex,
4856 options: OptionsIndex,
4857 future: u32,
4858 address: u32,
4859 ) -> Result<u32> {
4860 instance
4861 .guest_write(
4862 StoreContextMut(self),
4863 caller,
4864 TransmitIndex::Future(ty),
4865 options,
4866 None,
4867 future,
4868 address,
4869 1,
4870 )
4871 .map(|result| result.encode())
4872 }
4873
4874 fn future_read(
4875 &mut self,
4876 instance: Instance,
4877 caller: RuntimeComponentInstanceIndex,
4878 ty: TypeFutureTableIndex,
4879 options: OptionsIndex,
4880 future: u32,
4881 address: u32,
4882 ) -> Result<u32> {
4883 instance
4884 .guest_read(
4885 StoreContextMut(self),
4886 caller,
4887 TransmitIndex::Future(ty),
4888 options,
4889 None,
4890 future,
4891 address,
4892 1,
4893 )
4894 .map(|result| result.encode())
4895 }
4896
4897 fn stream_write(
4898 &mut self,
4899 instance: Instance,
4900 caller: RuntimeComponentInstanceIndex,
4901 ty: TypeStreamTableIndex,
4902 options: OptionsIndex,
4903 stream: u32,
4904 address: u32,
4905 count: u32,
4906 ) -> Result<u32> {
4907 instance
4908 .guest_write(
4909 StoreContextMut(self),
4910 caller,
4911 TransmitIndex::Stream(ty),
4912 options,
4913 None,
4914 stream,
4915 address,
4916 count,
4917 )
4918 .map(|result| result.encode())
4919 }
4920
4921 fn stream_read(
4922 &mut self,
4923 instance: Instance,
4924 caller: RuntimeComponentInstanceIndex,
4925 ty: TypeStreamTableIndex,
4926 options: OptionsIndex,
4927 stream: u32,
4928 address: u32,
4929 count: u32,
4930 ) -> Result<u32> {
4931 instance
4932 .guest_read(
4933 StoreContextMut(self),
4934 caller,
4935 TransmitIndex::Stream(ty),
4936 options,
4937 None,
4938 stream,
4939 address,
4940 count,
4941 )
4942 .map(|result| result.encode())
4943 }
4944
4945 fn future_drop_writable(
4946 &mut self,
4947 instance: Instance,
4948 ty: TypeFutureTableIndex,
4949 writer: u32,
4950 ) -> Result<()> {
4951 instance.guest_drop_writable(self, TransmitIndex::Future(ty), writer)
4952 }
4953
4954 fn flat_stream_write(
4955 &mut self,
4956 instance: Instance,
4957 caller: RuntimeComponentInstanceIndex,
4958 ty: TypeStreamTableIndex,
4959 options: OptionsIndex,
4960 payload_size: u32,
4961 payload_align: u32,
4962 stream: u32,
4963 address: u32,
4964 count: u32,
4965 ) -> Result<u32> {
4966 instance
4967 .guest_write(
4968 StoreContextMut(self),
4969 caller,
4970 TransmitIndex::Stream(ty),
4971 options,
4972 Some(FlatAbi {
4973 size: payload_size,
4974 align: payload_align,
4975 }),
4976 stream,
4977 address,
4978 count,
4979 )
4980 .map(|result| result.encode())
4981 }
4982
4983 fn flat_stream_read(
4984 &mut self,
4985 instance: Instance,
4986 caller: RuntimeComponentInstanceIndex,
4987 ty: TypeStreamTableIndex,
4988 options: OptionsIndex,
4989 payload_size: u32,
4990 payload_align: u32,
4991 stream: u32,
4992 address: u32,
4993 count: u32,
4994 ) -> Result<u32> {
4995 instance
4996 .guest_read(
4997 StoreContextMut(self),
4998 caller,
4999 TransmitIndex::Stream(ty),
5000 options,
5001 Some(FlatAbi {
5002 size: payload_size,
5003 align: payload_align,
5004 }),
5005 stream,
5006 address,
5007 count,
5008 )
5009 .map(|result| result.encode())
5010 }
5011
5012 fn stream_drop_writable(
5013 &mut self,
5014 instance: Instance,
5015 ty: TypeStreamTableIndex,
5016 writer: u32,
5017 ) -> Result<()> {
5018 instance.guest_drop_writable(self, TransmitIndex::Stream(ty), writer)
5019 }
5020
5021 fn error_context_debug_message(
5022 &mut self,
5023 instance: Instance,
5024 ty: TypeComponentLocalErrorContextTableIndex,
5025 options: OptionsIndex,
5026 err_ctx_handle: u32,
5027 debug_msg_address: u32,
5028 ) -> Result<()> {
5029 instance.error_context_debug_message(
5030 StoreContextMut(self),
5031 ty,
5032 options,
5033 err_ctx_handle,
5034 debug_msg_address,
5035 )
5036 }
5037
5038 fn thread_new_indirect(
5039 &mut self,
5040 instance: Instance,
5041 caller: RuntimeComponentInstanceIndex,
5042 func_ty_idx: TypeFuncIndex,
5043 start_func_table_idx: RuntimeTableIndex,
5044 start_func_idx: u32,
5045 context: i32,
5046 ) -> Result<u32> {
5047 instance.thread_new_indirect(
5048 StoreContextMut(self),
5049 caller,
5050 func_ty_idx,
5051 start_func_table_idx,
5052 start_func_idx,
5053 context,
5054 )
5055 }
5056}
5057
5058type HostTaskFuture = Pin<Box<dyn Future<Output = Result<()>> + Send + 'static>>;
5059
5060async fn run_with_host_task_set<F>(task: TableId<HostTask>, future: F) -> Result<F::Output>
5063where
5064 F: Future,
5065{
5066 let mut future = pin!(future);
5067 future::poll_fn(|cx| {
5068 let old_thread = match tls::get(|store| store.set_thread(task)) {
5069 Ok(thread) => thread,
5070 Err(error) => return Poll::Ready(Err(error)),
5071 };
5072 let result = future.as_mut().poll(cx);
5073 match tls::get(|store| store.set_thread(old_thread)) {
5074 Ok(_) => result.map(Ok),
5075 Err(error) => Poll::Ready(Err(error)),
5076 }
5077 })
5078 .await
5079}
5080
5081pub(crate) struct HostTask {
5085 common: WaitableCommon,
5086
5087 call_context: CallContext,
5090
5091 state: HostTaskState,
5092
5093 group: TaskGroupId,
5094}
5095
5096enum HostTaskState {
5097 CalleeStarted,
5102
5103 CalleeRunning(JoinHandle),
5108
5109 CalleeCancelling,
5113
5114 CalleeFinished(LiftedResult),
5118
5119 CalleeDone { cancelled: bool },
5122}
5123
5124impl HostTask {
5125 fn new(
5126 concurrent_state: &mut ConcurrentState,
5127 state: HostTaskState,
5128 caller: QualifiedThreadId,
5129 ) -> Result<Self> {
5130 let group = concurrent_state.get_mut(caller.task)?.group;
5131 concurrent_state.increment_group_ref_count(group)?;
5132
5133 Ok(Self {
5134 common: WaitableCommon::default(),
5135 call_context: CallContext::default(),
5136 state,
5137 group,
5138 })
5139 }
5140}
5141
5142impl TableDebug for HostTask {
5143 fn type_name() -> &'static str {
5144 "HostTask"
5145 }
5146}
5147
5148type CallbackFn = Box<dyn Fn(&mut dyn VMStore, Event, u32) -> Result<u32> + Send + Sync + 'static>;
5149
5150enum Caller {
5152 Host {
5154 tx: Option<oneshot::Sender<LiftedResult>>,
5156 host_future_present: bool,
5159 caller: Option<TableId<HostTask>>,
5163 },
5164 Guest {
5166 thread: QualifiedThreadId,
5168 },
5169}
5170
5171struct LiftResult {
5174 lift: RawLift,
5175 ty: TypeTupleIndex,
5176 memory: Option<SendSyncPtr<VMMemoryDefinition>>,
5177 string_encoding: StringEncoding,
5178}
5179
5180#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
5185pub(crate) struct QualifiedThreadId {
5186 task: TableId<GuestTask>,
5187 thread: TableId<GuestThread>,
5188}
5189
5190impl QualifiedThreadId {
5191 fn qualify(
5192 state: &mut ConcurrentState,
5193 thread: TableId<GuestThread>,
5194 ) -> Result<QualifiedThreadId> {
5195 Ok(QualifiedThreadId {
5196 task: state.get_mut(thread)?.parent_task,
5197 thread,
5198 })
5199 }
5200}
5201
5202impl fmt::Debug for QualifiedThreadId {
5203 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
5204 f.debug_tuple("QualifiedThreadId")
5205 .field(&self.task.rep())
5206 .field(&self.thread.rep())
5207 .finish()
5208 }
5209}
5210
5211enum GuestThreadState {
5212 NotStartedImplicit,
5213 NotStartedExplicit(
5214 Box<dyn FnOnce(&mut dyn VMStore, QualifiedThreadId) -> Result<()> + Send + Sync>,
5215 ),
5216 Running,
5217 Suspended(StoreFiber<'static>),
5218 Ready {
5219 fiber: StoreFiber<'static>,
5220 },
5221 Completed,
5222}
5223
5224impl fmt::Debug for GuestThreadState {
5225 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
5226 match self {
5227 Self::NotStartedImplicit => f.debug_tuple("NotStartedImplicit").finish(),
5228 Self::NotStartedExplicit(_) => f.debug_tuple("NotStartedExplicit").finish(),
5229 Self::Running => f.debug_tuple("Running").finish(),
5230 Self::Suspended(_) => f.debug_tuple("Suspended").finish(),
5231 Self::Ready { .. } => f.debug_struct("Ready").finish(),
5232 Self::Completed => f.debug_tuple("Completed").finish(),
5233 }
5234 }
5235}
5236
5237#[derive(Copy, Clone, PartialEq, Eq, Debug)]
5238enum WakeOnCancel {
5239 None,
5240 Waiting(TableId<WaitableSet>),
5241 Yielding,
5242}
5243
5244impl WakeOnCancel {
5245 fn is_none(self) -> bool {
5246 matches!(self, WakeOnCancel::None)
5247 }
5248
5249 fn replace(&mut self, other: WakeOnCancel) -> Self {
5250 let old = *self;
5251 *self = other;
5252 old
5253 }
5254
5255 fn take(&mut self) -> Self {
5256 self.replace(WakeOnCancel::None)
5257 }
5258}
5259
5260pub struct GuestThread {
5261 context: [u32; NUM_COMPONENT_CONTEXT_SLOTS],
5264 parent_task: TableId<GuestTask>,
5266 wake_on_cancel: WakeOnCancel,
5269 state: GuestThreadState,
5271 instance_rep: Option<u32>,
5274 sync_call_set: TableId<WaitableSet>,
5276 old_do_not_suspend: Option<bool>,
5279}
5280
5281impl GuestThread {
5282 fn from_instance(
5285 state: Pin<&mut ComponentInstance>,
5286 caller_instance: RuntimeComponentInstanceIndex,
5287 guest_thread: u32,
5288 ) -> Result<TableId<Self>> {
5289 let rep = state.instance_states().0[caller_instance]
5290 .thread_handle_table()
5291 .guest_thread_rep(guest_thread)?;
5292 Ok(TableId::new(rep))
5293 }
5294
5295 fn new_implicit(state: &mut ConcurrentState, parent_task: TableId<GuestTask>) -> Result<Self> {
5296 let sync_call_set = state.push(WaitableSet {
5297 is_sync_call_set: true,
5298 ..WaitableSet::default()
5299 })?;
5300 Ok(Self {
5301 context: [0; NUM_COMPONENT_CONTEXT_SLOTS],
5302 parent_task,
5303 wake_on_cancel: WakeOnCancel::None,
5304 state: GuestThreadState::NotStartedImplicit,
5305 instance_rep: None,
5306 sync_call_set,
5307 old_do_not_suspend: None,
5308 })
5309 }
5310
5311 fn new_explicit(
5312 state: &mut ConcurrentState,
5313 parent_task: TableId<GuestTask>,
5314 start_func: Box<
5315 dyn FnOnce(&mut dyn VMStore, QualifiedThreadId) -> Result<()> + Send + Sync,
5316 >,
5317 ) -> Result<Self> {
5318 let sync_call_set = state.push(WaitableSet {
5319 is_sync_call_set: true,
5320 ..WaitableSet::default()
5321 })?;
5322 Ok(Self {
5323 context: [0; NUM_COMPONENT_CONTEXT_SLOTS],
5324 parent_task,
5325 wake_on_cancel: WakeOnCancel::None,
5326 state: GuestThreadState::NotStartedExplicit(start_func),
5327 instance_rep: None,
5328 sync_call_set,
5329 old_do_not_suspend: None,
5330 })
5331 }
5332}
5333
5334impl TableDebug for GuestThread {
5335 fn type_name() -> &'static str {
5336 "GuestThread"
5337 }
5338}
5339
5340enum SyncResult {
5341 NotProduced,
5342 Produced(Option<ValRaw>),
5343 Taken,
5344}
5345
5346impl SyncResult {
5347 fn take(&mut self) -> Result<Option<Option<ValRaw>>> {
5348 Ok(match mem::replace(self, SyncResult::Taken) {
5349 SyncResult::NotProduced => None,
5350 SyncResult::Produced(val) => Some(val),
5351 SyncResult::Taken => {
5352 bail_bug!("attempted to take a synchronous result that was already taken")
5353 }
5354 })
5355 }
5356}
5357
5358#[derive(Debug)]
5359enum HostFutureState {
5360 NotApplicable,
5361 Live,
5362 Dropped,
5363}
5364
5365pub(crate) struct GuestTask {
5367 common: WaitableCommon,
5369 lower_params: Option<RawLower>,
5371 lift_result: Option<LiftResult>,
5373 result: Option<LiftedResult>,
5376 callback: Option<CallbackFn>,
5379 caller: Caller,
5381 call_context: CallContext,
5386 sync_result: SyncResult,
5389 cancel_request_delivered: bool,
5393 starting_sent: bool,
5396 instance: RuntimeInstance,
5403 event: Option<Event>,
5405 exited: bool,
5407 threads: HashSet<TableId<GuestThread>>,
5409 host_future_state: HostFutureState,
5412 async_typed: bool,
5415 async_lifted: bool,
5418
5419 decremented_interesting_task_count: bool,
5420
5421 group: TaskGroupId,
5422}
5423
5424impl GuestTask {
5425 fn already_lowered_parameters(&self) -> bool {
5426 self.lower_params.is_none()
5428 }
5429
5430 fn returned_or_cancelled(&self) -> bool {
5431 self.lift_result.is_none()
5433 }
5434
5435 fn ready_to_delete(&self) -> bool {
5436 let threads_completed = self.threads.is_empty();
5437 let has_sync_result = matches!(self.sync_result, SyncResult::Produced(_));
5438 let pending_completion_event = matches!(
5439 self.common.event,
5440 Some(Event::Subtask {
5441 status: Status::Returned | Status::ReturnCancelled
5442 })
5443 );
5444 let ready = threads_completed
5445 && !has_sync_result
5446 && !pending_completion_event
5447 && !matches!(self.host_future_state, HostFutureState::Live);
5448 log::trace!(
5449 "ready to delete? {ready} (threads_completed: {}, has_sync_result: {}, pending_completion_event: {}, host_future_state: {:?})",
5450 threads_completed,
5451 has_sync_result,
5452 pending_completion_event,
5453 self.host_future_state
5454 );
5455 ready
5456 }
5457
5458 fn new(
5459 state: &mut ConcurrentState,
5460 lower_params: RawLower,
5461 lift_result: LiftResult,
5462 caller: Caller,
5463 callback: Option<CallbackFn>,
5464 instance: RuntimeInstance,
5465 async_typed: bool,
5466 async_lifted: bool,
5467 ) -> Result<QualifiedThreadId> {
5468 let host_future_state = match &caller {
5469 Caller::Guest { .. } => HostFutureState::NotApplicable,
5470 Caller::Host {
5471 host_future_present,
5472 ..
5473 } => {
5474 if *host_future_present {
5475 HostFutureState::Live
5476 } else {
5477 HostFutureState::NotApplicable
5478 }
5479 }
5480 };
5481
5482 let group = match caller {
5483 Caller::Guest { thread } => {
5484 let group = state.get_mut(thread.task)?.group;
5485 state.increment_group_ref_count(group)?;
5486 group
5487 }
5488 Caller::Host { .. } => state.make_task_group()?,
5489 };
5490
5491 let task = state.push(Self {
5492 common: WaitableCommon::default(),
5493 lower_params: Some(lower_params),
5494 lift_result: Some(lift_result),
5495 result: None,
5496 callback,
5497 caller,
5498 call_context: CallContext::default(),
5499 sync_result: SyncResult::NotProduced,
5500 cancel_request_delivered: false,
5501 starting_sent: false,
5502 instance,
5503 event: None,
5504 exited: false,
5505 threads: HashSet::new(),
5506 host_future_state,
5507 async_typed,
5508 async_lifted,
5509 decremented_interesting_task_count: false,
5510 group,
5511 })?;
5512 let new_thread = GuestThread::new_implicit(state, task)?;
5513 let thread = state.push(new_thread)?;
5514 state.get_mut(task)?.threads.insert(thread);
5515 state.interesting_tasks += 1;
5516 let thread = QualifiedThreadId { task, thread };
5517 log::trace!("new implicit thread {thread:?} for instance {instance:?}");
5518 Ok(thread)
5519 }
5520}
5521
5522impl TableDebug for GuestTask {
5523 fn type_name() -> &'static str {
5524 "GuestTask"
5525 }
5526}
5527
5528#[derive(Default)]
5530struct WaitableCommon {
5531 event: Option<Event>,
5533 set: Option<TableId<WaitableSet>>,
5535 handle: Option<u32>,
5537}
5538
5539#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
5541enum Waitable {
5542 Host(TableId<HostTask>),
5544 Guest(TableId<GuestTask>),
5546 Transmit(TableId<TransmitHandle>),
5548}
5549
5550impl Waitable {
5551 fn from_instance(
5554 state: Pin<&mut ComponentInstance>,
5555 caller_instance: RuntimeComponentInstanceIndex,
5556 waitable: u32,
5557 ) -> Result<Self> {
5558 use crate::runtime::vm::component::Waitable;
5559
5560 let (waitable, kind) = state.instance_states().0[caller_instance]
5561 .handle_table()
5562 .waitable_rep(waitable)?;
5563
5564 Ok(match kind {
5565 Waitable::Subtask { is_host: true } => Self::Host(TableId::new(waitable)),
5566 Waitable::Subtask { is_host: false } => Self::Guest(TableId::new(waitable)),
5567 Waitable::Stream | Waitable::Future => Self::Transmit(TableId::new(waitable)),
5568 })
5569 }
5570
5571 fn rep(&self) -> u32 {
5573 match self {
5574 Self::Host(id) => id.rep(),
5575 Self::Guest(id) => id.rep(),
5576 Self::Transmit(id) => id.rep(),
5577 }
5578 }
5579
5580 fn join(&self, state: &mut ConcurrentState, set: Option<TableId<WaitableSet>>) -> Result<()> {
5584 log::trace!("waitable {self:?} join set {set:?}");
5585
5586 let old = mem::replace(&mut self.common(state)?.set, set);
5587
5588 if let Some(old) = old {
5589 match *self {
5590 Waitable::Host(id) => state.remove_child(id, old),
5591 Waitable::Guest(id) => state.remove_child(id, old),
5592 Waitable::Transmit(id) => state.remove_child(id, old),
5593 }?;
5594
5595 state.get_mut(old)?.ready.remove(self);
5596 }
5597
5598 if let Some(set) = set {
5599 match *self {
5600 Waitable::Host(id) => state.add_child(id, set),
5601 Waitable::Guest(id) => state.add_child(id, set),
5602 Waitable::Transmit(id) => state.add_child(id, set),
5603 }?;
5604
5605 if self.common(state)?.event.is_some() {
5606 self.mark_ready(state)?;
5607 }
5608 }
5609
5610 Ok(())
5611 }
5612
5613 fn common<'a>(&self, state: &'a mut ConcurrentState) -> Result<&'a mut WaitableCommon> {
5615 Ok(match self {
5616 Self::Host(id) => &mut state.get_mut(*id)?.common,
5617 Self::Guest(id) => &mut state.get_mut(*id)?.common,
5618 Self::Transmit(id) => &mut state.get_mut(*id)?.common,
5619 })
5620 }
5621
5622 fn trap_if_in_waitable_set(&self, state: &mut ConcurrentState) -> Result<()> {
5628 if self.common(state)?.set.is_some() {
5629 bail!(Trap::WaitableSyncAndAsync);
5630 }
5631 Ok(())
5632 }
5633
5634 fn set_event(&self, state: &mut ConcurrentState, event: Option<Event>) -> Result<()> {
5638 log::trace!("set event for {self:?}: {event:?}");
5639 self.common(state)?.event = event;
5640 self.mark_ready(state)
5641 }
5642
5643 fn take_event(&self, state: &mut ConcurrentState) -> Result<Option<Event>> {
5645 let common = self.common(state)?;
5646 let event = common.event.take();
5647 if let Some(set) = self.common(state)?.set {
5648 state.get_mut(set)?.ready.remove(self);
5649 }
5650
5651 Ok(event)
5652 }
5653
5654 fn mark_ready(&self, state: &mut ConcurrentState) -> Result<()> {
5658 if let Some(set) = self.common(state)?.set {
5659 let set_state = state.get_mut(set)?;
5660 set_state.ready.insert(*self);
5661
5662 if let Some((thread, mode)) = set_state.waiting.pop_first() {
5663 let wake_on_cancel = state.get_mut(thread.thread)?.wake_on_cancel.take();
5664 assert!(wake_on_cancel.is_none() || wake_on_cancel == WakeOnCancel::Waiting(set));
5665
5666 let item = match mode {
5667 WaitMode::Fiber(fiber) => Some(WorkItem::ResumeFiber {
5668 instance: state.get_mut(thread.task)?.instance,
5669 thread,
5670 fiber,
5671 }),
5672 WaitMode::Callback(instance) => Some(WorkItem::GuestCall {
5673 instance: state.get_mut(thread.task)?.instance,
5674 call: GuestCall {
5675 thread,
5676 kind: GuestCallKind::DeliverEvent {
5677 instance,
5678 set: Some(set),
5679 },
5680 },
5681 }),
5682 };
5683
5684 if let Some(item) = item {
5685 state.push_high_priority(item);
5686 }
5687 }
5688 }
5689 Ok(())
5690 }
5691
5692 fn delete_from(&self, store: &mut StoreOpaque) -> Result<()> {
5694 match self {
5695 Self::Host(task) => {
5696 log::trace!("delete host task {task:?}");
5697 let state = store.concurrent_state_mut()?;
5698 let task = state.delete(*task)?;
5699
5700 state.decrement_group_ref_count(task.group)?;
5701 }
5702 Self::Guest(task) => {
5703 log::trace!("delete guest task {task:?}");
5704 let state = store.concurrent_state_mut()?;
5705 let task = state.delete(*task)?;
5706
5707 state.decrement_group_ref_count(task.group)?;
5708
5709 debug_assert!(task.decremented_interesting_task_count);
5716 }
5717 Self::Transmit(task) => {
5718 store.concurrent_state_mut()?.delete(*task)?;
5719 }
5720 }
5721
5722 Ok(())
5723 }
5724}
5725
5726impl fmt::Debug for Waitable {
5727 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
5728 match self {
5729 Self::Host(id) => write!(f, "{id:?}"),
5730 Self::Guest(id) => write!(f, "{id:?}"),
5731 Self::Transmit(id) => write!(f, "{id:?}"),
5732 }
5733 }
5734}
5735
5736#[derive(Default)]
5738struct WaitableSet {
5739 ready: BTreeSet<Waitable>,
5741 waiting: BTreeMap<QualifiedThreadId, WaitMode>,
5743 num_waiting: usize,
5746 is_sync_call_set: bool,
5749}
5750
5751impl WaitableSet {
5752 fn stop_waiting(&mut self) -> Result<()> {
5754 self.num_waiting = match self.num_waiting.checked_sub(1) {
5755 Some(n) => n,
5756 None => bail_bug!("waiter not accounted for in waitable set"),
5757 };
5758 Ok(())
5759 }
5760}
5761
5762impl TableDebug for WaitableSet {
5763 fn type_name() -> &'static str {
5764 "WaitableSet"
5765 }
5766}
5767
5768type RawLower =
5770 Box<dyn FnOnce(&mut dyn VMStore, &mut [MaybeUninit<ValRaw>]) -> Result<()> + Send + Sync>;
5771
5772type RawLift = Box<
5774 dyn FnOnce(&mut dyn VMStore, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>> + Send + Sync,
5775>;
5776
5777type LiftedResult = Box<dyn Any + Send + Sync>;
5781
5782struct DummyResult;
5785
5786#[derive(Default)]
5788pub struct ConcurrentInstanceState {
5789 backpressure: u16,
5791 do_not_enter: bool,
5793 do_not_suspend: bool,
5796 pending: BTreeMap<QualifiedThreadId, GuestCallKind>,
5799}
5800
5801impl ConcurrentInstanceState {
5802 pub fn pending_is_empty(&self) -> bool {
5803 self.pending.is_empty()
5804 }
5805}
5806
5807#[derive(Debug, Copy, Clone)]
5808pub(crate) enum CurrentThread {
5809 Guest(QualifiedThreadId),
5812 Host(TableId<HostTask>),
5814 DeferredHost(QualifiedThreadId),
5817 None,
5820}
5821
5822impl CurrentThread {
5823 fn guest(&self) -> Option<&QualifiedThreadId> {
5824 match self {
5825 Self::Guest(id) => Some(id),
5826 _ => None,
5827 }
5828 }
5829
5830 fn guest_task(&self) -> Option<TableId<GuestTask>> {
5831 match self {
5832 Self::Guest(id) => Some(id.task),
5833 _ => None,
5834 }
5835 }
5836
5837 fn is_none(&self) -> bool {
5838 matches!(self, Self::None)
5839 }
5840}
5841
5842impl From<QualifiedThreadId> for CurrentThread {
5843 fn from(id: QualifiedThreadId) -> Self {
5844 Self::Guest(id)
5845 }
5846}
5847
5848impl From<TableId<HostTask>> for CurrentThread {
5849 fn from(id: TableId<HostTask>) -> Self {
5850 Self::Host(id)
5851 }
5852}
5853
5854enum Priority {
5855 Switch,
5856 High,
5857 Low,
5858}
5859
5860pub struct ConcurrentState {
5862 unforced_current_thread: CurrentThread,
5868
5869 deferred_host_call_context: Option<CallContext>,
5875
5876 futures: AlwaysMut<Option<FuturesUnordered<HostTaskFuture>>>,
5881 table: AlwaysMut<ResourceTable>,
5883 switch_item: Option<WorkItem>,
5891 next_switch_item: Option<WorkItem>,
5897 high_priority: VecDeque<WorkItem>,
5899 low_priority: VecDeque<WorkItem>,
5901 suspend_reason: Option<SuspendReason>,
5905 worker: Option<StoreFiber<'static>>,
5909 worker_item: Option<WorkerItem>,
5911
5912 global_error_context_ref_counts:
5925 BTreeMap<TypeComponentGlobalErrorContextTableIndex, GlobalErrorContextRefCount>,
5926
5927 interesting_tasks: usize,
5940
5941 interesting_tasks_empty_waker: Option<Waker>,
5945
5946 ready_for_concurrent_call_waker: Option<Waker>,
5951
5952 event_loop_running: bool,
5954
5955 #[cfg(feature = "task-group-hook")]
5957 task_group_hook: Option<Box<dyn TaskGroupHook>>,
5958}
5959
5960impl Default for ConcurrentState {
5961 fn default() -> Self {
5962 Self {
5963 unforced_current_thread: CurrentThread::None,
5964 deferred_host_call_context: None,
5965 table: AlwaysMut::new(ResourceTable::new()),
5966 futures: AlwaysMut::new(Some(FuturesUnordered::new())),
5967 switch_item: None,
5968 next_switch_item: None,
5969 high_priority: VecDeque::new(),
5970 low_priority: VecDeque::new(),
5971 suspend_reason: None,
5972 worker: None,
5973 worker_item: None,
5974 global_error_context_ref_counts: BTreeMap::new(),
5975 interesting_tasks: 0,
5976 interesting_tasks_empty_waker: None,
5977 ready_for_concurrent_call_waker: None,
5978 event_loop_running: false,
5979 #[cfg(feature = "task-group-hook")]
5980 task_group_hook: None,
5981 }
5982 }
5983}
5984
5985impl ConcurrentState {
5986 pub(crate) fn take_fibers_and_futures(
6003 &mut self,
6004 fibers: &mut Vec<StoreFiber<'static>>,
6005 futures: &mut Vec<FuturesUnordered<HostTaskFuture>>,
6006 ) {
6007 let mut items = Vec::new();
6008 for (_, entry) in self.table.get_mut().iter_mut() {
6009 if let Some(set) = entry.downcast_mut::<WaitableSet>() {
6010 for mode in mem::take(&mut set.waiting).into_values() {
6011 match mode {
6012 WaitMode::Fiber(fiber) => {
6013 fibers.push(fiber);
6014 }
6015 WaitMode::Callback(_) => {}
6016 }
6017 }
6018 } else if let Some(thread) = entry.downcast_mut::<GuestThread>() {
6019 if let GuestThreadState::Suspended(fiber) | GuestThreadState::Ready { fiber, .. } =
6020 mem::replace(&mut thread.state, GuestThreadState::Completed)
6021 {
6022 fibers.push(fiber);
6023 }
6024 } else if let Some(item) = entry.downcast_mut::<Option<WorkItem>>() {
6025 if let Some(item) = item.take() {
6026 items.push(item);
6027 }
6028 }
6029 }
6030
6031 if let Some(fiber) = self.worker.take() {
6032 fibers.push(fiber);
6033 }
6034
6035 let mut handle_item = |item| match item {
6036 WorkItem::ResumeFiber { fiber, .. } => {
6037 fibers.push(fiber);
6038 }
6039 WorkItem::PushFuture(future) => {
6040 self.futures
6041 .get_mut()
6042 .as_mut()
6043 .unwrap()
6044 .push(future.into_inner());
6045 }
6046 WorkItem::ResumeThread { .. }
6047 | WorkItem::GuestCall { .. }
6048 | WorkItem::WorkerFunction(_) => {}
6049 };
6050
6051 for item in items {
6052 handle_item(item);
6053 }
6054 if let Some(item) = self.switch_item.take() {
6055 handle_item(item);
6056 }
6057 if let Some(item) = self.next_switch_item.take() {
6058 handle_item(item);
6059 }
6060 for item in mem::take(&mut self.high_priority) {
6061 handle_item(item);
6062 }
6063 for item in mem::take(&mut self.low_priority) {
6064 handle_item(item);
6065 }
6066
6067 if let Some(them) = self.futures.get_mut().take() {
6068 futures.push(them);
6069 }
6070 }
6071
6072 #[cfg(feature = "gc")]
6073 pub(crate) fn trace_fiber_roots(
6074 &mut self,
6075 modules: &ModuleRegistry,
6076 unwind: &dyn Unwind,
6077 gc_roots_list: &mut GcRootsList,
6078 ) {
6079 let ConcurrentState {
6080 table,
6081 worker,
6082 switch_item,
6083 next_switch_item,
6084 high_priority,
6085 low_priority,
6086
6087 futures: _,
6091
6092 worker_item: _,
6094 unforced_current_thread: _,
6095 deferred_host_call_context: _,
6096 suspend_reason: _,
6097 global_error_context_ref_counts: _,
6098 interesting_tasks: _,
6099 interesting_tasks_empty_waker: _,
6100 ready_for_concurrent_call_waker: _,
6101 event_loop_running: _,
6102 #[cfg(feature = "task-group-hook")]
6103 task_group_hook: _,
6104 } = self;
6105
6106 for (_, entry) in table.get_mut().iter_mut() {
6107 if let Some(set) = entry.downcast_mut::<WaitableSet>() {
6108 for mode in set.waiting.values_mut() {
6109 match mode {
6110 WaitMode::Fiber(fiber) => {
6111 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6112 }
6113 WaitMode::Callback(_) => {}
6114 }
6115 }
6116 } else if let Some(thread) = entry.downcast_mut::<GuestThread>() {
6117 if let GuestThreadState::Suspended(fiber) | GuestThreadState::Ready { fiber, .. } =
6118 &mut thread.state
6119 {
6120 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6121 }
6122 } else if let Some(Some(WorkItem::ResumeFiber { fiber, .. })) =
6123 entry.downcast_mut::<Option<WorkItem>>()
6124 {
6125 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6126 }
6127 }
6128
6129 if let Some(fiber) = worker {
6130 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6131 }
6132
6133 let mut handle_item = |item: &mut WorkItem| match item {
6134 WorkItem::ResumeFiber { fiber, .. } => {
6135 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6136 }
6137 WorkItem::PushFuture(_future) => {
6138 }
6141 WorkItem::ResumeThread { .. }
6142 | WorkItem::GuestCall { .. }
6143 | WorkItem::WorkerFunction(_) => {}
6144 };
6145
6146 if let Some(item) = switch_item {
6147 handle_item(item);
6148 }
6149 if let Some(item) = next_switch_item {
6150 handle_item(item);
6151 }
6152 for item in high_priority {
6153 handle_item(item);
6154 }
6155 for item in low_priority {
6156 handle_item(item);
6157 }
6158 }
6159
6160 fn push<V: Send + Sync + 'static>(
6161 &mut self,
6162 value: V,
6163 ) -> Result<TableId<V>, ResourceTableError> {
6164 self.table.get_mut().push(value).map(TableId::from)
6165 }
6166
6167 fn get_mut<V: 'static>(&mut self, id: TableId<V>) -> Result<&mut V, ResourceTableError> {
6168 self.table.get_mut().get_mut(&Resource::from(id))
6169 }
6170
6171 pub fn add_child<T: 'static, U: 'static>(
6172 &mut self,
6173 child: TableId<T>,
6174 parent: TableId<U>,
6175 ) -> Result<(), ResourceTableError> {
6176 self.table
6177 .get_mut()
6178 .add_child(Resource::from(child), Resource::from(parent))
6179 }
6180
6181 pub fn remove_child<T: 'static, U: 'static>(
6182 &mut self,
6183 child: TableId<T>,
6184 parent: TableId<U>,
6185 ) -> Result<(), ResourceTableError> {
6186 self.table
6187 .get_mut()
6188 .remove_child(Resource::from(child), Resource::from(parent))
6189 }
6190
6191 fn delete<V: 'static>(&mut self, id: TableId<V>) -> Result<V, ResourceTableError> {
6192 self.table.get_mut().delete(Resource::from(id))
6193 }
6194
6195 fn push_future(&mut self, future: HostTaskFuture) {
6196 self.push_high_priority(WorkItem::PushFuture(AlwaysMut::new(future)));
6203 }
6204
6205 fn set_switch_item(&mut self, item: WorkItem) -> Result<()> {
6206 log::trace!("set switch item: {item:?}");
6207
6208 if self.switch_item.is_some() {
6209 bail_bug!("switch item already set");
6210 }
6211
6212 self.switch_item = Some(item);
6213
6214 Ok(())
6215 }
6216
6217 fn take_next_switch_item(&mut self) -> Result<()> {
6218 if let Some(item) = self.next_switch_item.take() {
6219 self.set_switch_item(item)?;
6220 }
6221 Ok(())
6222 }
6223
6224 fn push_high_priority(&mut self, item: WorkItem) {
6225 log::trace!("push high priority: {item:?}");
6226 self.high_priority.push_front(item);
6227 }
6228
6229 fn push_low_priority(&mut self, item: WorkItem) {
6230 log::trace!("push low priority: {item:?}");
6231 self.low_priority.push_front(item);
6232 }
6233
6234 fn push_work_item(&mut self, item: WorkItem, priority: Priority) -> Result<()> {
6235 match priority {
6236 Priority::Switch => self.set_switch_item(item)?,
6237 Priority::High => self.push_high_priority(item),
6238 Priority::Low => self.push_low_priority(item),
6239 }
6240
6241 Ok(())
6242 }
6243
6244 fn promote_instance_local_thread_work_item(
6245 &mut self,
6246 current_instance: RuntimeInstance,
6247 ) -> Result<bool> {
6248 log::trace!("promote thread work items for {current_instance:?}");
6249
6250 self.promote_work_item_matching(|item: &WorkItem| {
6251 let result = match item {
6252 WorkItem::ResumeThread { instance, .. }
6253 | WorkItem::ResumeFiber { instance, .. }
6254 | WorkItem::GuestCall { instance, .. } => *instance == current_instance,
6255 _ => false,
6256 };
6257
6258 log::trace!("candidate {item:?}: {result}");
6259 result
6260 })
6261 }
6262
6263 fn promote_thread_work_item(&mut self, thread: QualifiedThreadId) -> Result<bool> {
6264 self.promote_work_item_matching(|item: &WorkItem| match item {
6265 WorkItem::ResumeThread {
6266 thread: item_thread,
6267 ..
6268 }
6269 | WorkItem::GuestCall {
6270 call:
6271 GuestCall {
6272 thread: item_thread,
6273 ..
6274 },
6275 ..
6276 } => *item_thread == thread,
6277 _ => false,
6278 })
6279 }
6280
6281 fn promote_work_item_matching<F>(&mut self, mut predicate: F) -> Result<bool>
6282 where
6283 F: FnMut(&WorkItem) -> bool,
6284 {
6285 for item in mem::take(&mut self.high_priority).into_iter().rev() {
6290 if self.switch_item.is_none() && predicate(&item) {
6291 self.set_switch_item(item)?;
6292 } else {
6293 self.push_high_priority(item);
6294 }
6295 }
6296
6297 if self.switch_item.is_none() {
6298 for item in mem::take(&mut self.low_priority).into_iter().rev() {
6299 if self.switch_item.is_none() && predicate(&item) {
6300 self.set_switch_item(item)?;
6301 } else {
6302 self.push_low_priority(item);
6303 }
6304 }
6305 }
6306
6307 Ok(self.switch_item.is_some())
6308 }
6309
6310 pub fn call_context(&mut self, task: Scope) -> Result<&mut CallContext> {
6313 match task {
6314 Scope::HostId(task) => {
6315 let task: TableId<HostTask> = TableId::new(task);
6316 Ok(&mut self.get_mut(task)?.call_context)
6317 }
6318 Scope::Id(task) => {
6319 let task: TableId<GuestTask> = TableId::new(task);
6320 Ok(&mut self.get_mut(task)?.call_context)
6321 }
6322 }
6323 }
6324
6325 pub(crate) fn deferred_host_call_context(&mut self) -> Option<&mut CallContext> {
6326 self.deferred_host_call_context.as_mut()
6327 }
6328
6329 fn futures_mut(&mut self) -> Result<&mut FuturesUnordered<HostTaskFuture>> {
6330 match self.futures.get_mut().as_mut() {
6331 Some(f) => Ok(f),
6332 None => bail_bug!("futures field of concurrent state is currently taken"),
6333 }
6334 }
6335
6336 pub(crate) fn table(&mut self) -> &mut ResourceTable {
6337 self.table.get_mut()
6338 }
6339
6340 fn debug_assert_deferred_host_invariant(&self) {
6341 debug_assert_eq!(
6342 self.deferred_host_call_context.is_some(),
6343 matches!(self.unforced_current_thread, CurrentThread::DeferredHost(_)),
6344 "a deferred host thread and call context must exist together",
6345 );
6346 }
6347
6348 fn materialize_host_task(&mut self) -> Result<CurrentThread> {
6349 self.debug_assert_deferred_host_invariant();
6350 let caller = match self.unforced_current_thread {
6351 CurrentThread::DeferredHost(caller) => caller,
6352 thread => return Ok(thread),
6353 };
6354
6355 let task = HostTask::new(self, HostTaskState::CalleeStarted, caller)?;
6357 let task = self.push(task)?;
6358 let call_context = self
6359 .deferred_host_call_context
6360 .take()
6361 .expect("deferred host call context should be present");
6362 self.get_mut(task)
6363 .expect("newly inserted host task should be present")
6364 .call_context = call_context;
6365 self.unforced_current_thread = CurrentThread::Host(task);
6366 self.debug_assert_deferred_host_invariant();
6367 log::trace!("new host task materialized {task:?}");
6368 Ok(CurrentThread::Host(task))
6369 }
6370
6371 fn materialize_current_host_task_id(&mut self) -> Result<Option<TableId<HostTask>>> {
6372 match self.materialize_host_task()? {
6373 CurrentThread::Host(id) => Ok(Some(id)),
6374 CurrentThread::None => Ok(None),
6375 CurrentThread::Guest(_) => {
6376 bail_bug!("tried to materialize a host task id from a guest thread")
6377 }
6378 CurrentThread::DeferredHost(_) => {
6379 bail_bug!(
6380 "current thread is a deferred host thread which should have been materialized"
6381 )
6382 }
6383 }
6384 }
6385
6386 pub(crate) fn materialize_current_scope(&mut self) -> Result<Scope> {
6387 match self.materialize_host_task()? {
6388 CurrentThread::Host(id) => Ok(Scope::HostId(id.rep())),
6389 _ => bail_bug!("current scope is not a deferred host scope"),
6390 }
6391 }
6392}
6393
6394fn for_any_lower<
6397 F: FnOnce(&mut dyn VMStore, &mut [MaybeUninit<ValRaw>]) -> Result<()> + Send + Sync,
6398>(
6399 fun: F,
6400) -> F {
6401 fun
6402}
6403
6404fn for_any_lift<
6406 F: FnOnce(&mut dyn VMStore, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>> + Send + Sync,
6407>(
6408 fun: F,
6409) -> F {
6410 fun
6411}
6412
6413fn check_ambient_store(id: StoreId) {
6414 let message = "\
6415 `Future`s which depend on asynchronous component tasks, streams, or \
6416 futures to complete may only be polled from the event loop of the \
6417 store to which they belong. Please use \
6418 `StoreContextMut::{run_concurrent,spawn}` to poll or await them.\
6419 ";
6420 tls::try_get(|store| {
6421 let matched = match store {
6422 tls::TryGet::Some(store) => store.id() == id,
6423 tls::TryGet::Taken | tls::TryGet::None => false,
6424 };
6425
6426 if !matched {
6427 panic!("{message}")
6428 }
6429 });
6430}
6431
6432fn unpack_callback_code(code: u32) -> (u32, u32) {
6433 (code & 0xF, code >> 4)
6434}
6435
6436struct WaitableCheckParams {
6440 set: TableId<WaitableSet>,
6441 options: OptionsIndex,
6442 payload: u32,
6443}
6444
6445enum WaitableCheck {
6448 Wait,
6449 Poll,
6450}
6451
6452pub(crate) struct PreparedCall<R> {
6454 handle: Func,
6456 thread: QualifiedThreadId,
6458 param_count: usize,
6460 rx: oneshot::Receiver<LiftedResult>,
6463 runtime_instance: RuntimeInstance,
6465 _phantom: PhantomData<R>,
6466}
6467
6468impl<R> PreparedCall<R> {
6469 pub(crate) fn task_id(&self) -> TaskId {
6471 TaskId {
6472 task: self.thread.task,
6473 runtime_instance: self.runtime_instance,
6474 }
6475 }
6476}
6477
6478pub(crate) struct TaskId {
6480 task: TableId<GuestTask>,
6481 runtime_instance: RuntimeInstance,
6482}
6483
6484impl TaskId {
6485 pub(crate) fn host_future_dropped(&self, store: &mut StoreOpaque) -> Result<()> {
6491 let task = store.concurrent_state_mut()?.get_mut(self.task)?;
6492 let delete = if !task.already_lowered_parameters() {
6493 store.cancel_guest_subtask_without_lowered_parameters(
6494 self.runtime_instance,
6495 self.task,
6496 )?;
6497 true
6498 } else {
6499 task.host_future_state = HostFutureState::Dropped;
6500 task.ready_to_delete()
6501 };
6502 if delete {
6503 Waitable::Guest(self.task).delete_from(store)?
6504 }
6505 Ok(())
6506 }
6507}
6508
6509pub(crate) fn prepare_call<T, R>(
6515 mut store: StoreContextMut<T>,
6516 handle: Func,
6517 param_count: usize,
6518 host_future_present: bool,
6519 lower_params: impl FnOnce(StoreContextMut<T>, &mut [MaybeUninit<ValRaw>]) -> Result<()>
6520 + Send
6521 + Sync
6522 + 'static,
6523 lift_result: impl FnOnce(&mut StoreOpaque, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>>
6524 + Send
6525 + Sync
6526 + 'static,
6527) -> Result<PreparedCall<R>> {
6528 if !store.0.may_enter() {
6529 bail!(Trap::CannotEnterComponent);
6530 }
6531
6532 let (options, _flags, ty, raw_options) = handle.abi_info(store.0);
6533
6534 let instance = handle.instance().id().get(store.0);
6535 let options = &instance.component().env_component().options[options];
6536 let ty = &instance.component().types()[ty];
6537 let async_typed = ty.async_;
6538 let async_lifted = raw_options.async_;
6539 let task_return_type = ty.results;
6540 let component_instance = raw_options.instance;
6541 let callback = options.callback.map(|i| instance.runtime_callback(i));
6542 let memory = options
6543 .memory()
6544 .map(|i| instance.runtime_memory(i))
6545 .map(SendSyncPtr::new);
6546 let string_encoding = options.string_encoding;
6547 let token = StoreToken::new(store.as_context_mut());
6548 let caller = store.0.materialize_host_task_id()?;
6549 let state = store.0.concurrent_state_mut()?;
6550
6551 let (tx, rx) = oneshot::channel();
6552
6553 let instance = handle.instance().runtime_instance(component_instance);
6554 let thread = GuestTask::new(
6555 state,
6556 Box::new(for_any_lower(move |store, params| {
6557 lower_params(token.as_context_mut(store), params)
6558 })),
6559 LiftResult {
6560 lift: Box::new(for_any_lift(move |store, result| {
6561 lift_result(store, result)
6562 })),
6563 ty: task_return_type,
6564 memory,
6565 string_encoding,
6566 },
6567 Caller::Host {
6568 tx: Some(tx),
6569 host_future_present,
6570 caller,
6571 },
6572 callback.map(|callback| {
6573 let callback = SendSyncPtr::new(callback);
6574 let instance = handle.instance();
6575 Box::new(move |store: &mut dyn VMStore, event, handle| {
6576 let store = token.as_context_mut(store);
6577 unsafe { instance.call_callback(store, callback, event, handle) }
6580 }) as CallbackFn
6581 }),
6582 instance,
6583 async_typed,
6584 async_lifted,
6585 )?;
6586
6587 Ok(PreparedCall {
6588 handle,
6589 thread,
6590 param_count,
6591 runtime_instance: instance,
6592 rx,
6593 _phantom: PhantomData,
6594 })
6595}
6596
6597pub(crate) struct StagedCall<R> {
6598 store: StoreId,
6599 rx: oneshot::Receiver<LiftedResult>,
6600 _marker: PhantomData<fn() -> R>,
6601 group: TaskGroupId,
6602}
6603
6604impl<R> StagedCall<R> {
6605 pub(crate) fn new<T: 'static>(
6612 mut store: StoreContextMut<T>,
6613 prepared: PreparedCall<R>,
6614 ) -> Result<StagedCall<R>> {
6615 let PreparedCall {
6616 handle,
6617 thread,
6618 param_count,
6619 rx,
6620 ..
6621 } = prepared;
6622
6623 stage_call0(store.as_context_mut(), handle, thread, param_count)?;
6624
6625 Ok(StagedCall {
6626 store: store.0.id(),
6627 rx,
6628 _marker: PhantomData,
6629 group: store.0.concurrent_state_mut()?.get_mut(thread.task)?.group,
6630 })
6631 }
6632}
6633
6634impl<R> Future for StagedCall<R>
6635where
6636 R: 'static,
6637{
6638 type Output = Result<R>;
6639
6640 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
6641 check_ambient_store(self.store);
6642 Pin::new(&mut self.rx).poll(cx).map(|result| match result {
6643 Ok(r) => match r.downcast() {
6644 Ok(r) => Ok(*r),
6645 Err(_) => bail_bug!("wrong type of value produced"),
6646 },
6647 Err(oneshot::Canceled) => bail_bug!("channel erroneously dropped"),
6648 })
6649 }
6650}
6651
6652fn stage_call0<T: 'static>(
6655 store: StoreContextMut<T>,
6656 handle: Func,
6657 guest_thread: QualifiedThreadId,
6658 param_count: usize,
6659) -> Result<()> {
6660 let (_options, _, _ty, raw_options) = handle.abi_info(store.0);
6661 let is_concurrent = raw_options.async_;
6662 let callback = raw_options.callback;
6663 let instance = handle.instance();
6664 let callee = handle.lifted_core_func(store.0);
6665 let post_return = raw_options
6666 .post_return
6667 .map(|i| instance.id().get(store.0).runtime_post_return(i));
6668 let callback = callback.map(|i| {
6669 let instance = instance.id().get(store.0);
6670 SendSyncPtr::new(instance.runtime_callback(i))
6671 });
6672
6673 log::trace!("queueing call {guest_thread:?}");
6674
6675 unsafe {
6679 instance.stage_call(
6680 store,
6681 guest_thread,
6682 SendSyncPtr::new(callee),
6683 param_count,
6684 1,
6685 is_concurrent,
6686 callback,
6687 post_return.map(SendSyncPtr::new),
6688 true,
6689 )
6690 }
6691}
6692
6693#[cfg(all(test, feature = "cranelift", feature = "wat"))]
6694mod tests {
6695 use super::*;
6696 use crate::component::{Component, Linker};
6697 use crate::store::AsStoreOpaque;
6698 use crate::{Config, Engine};
6699
6700 fn host_subtask(
6701 state: HostTaskState,
6702 event: Option<Event>,
6703 ) -> Result<(Store<()>, Instance, TableId<HostTask>, u32)> {
6704 let mut config = Config::new();
6705 config.wasm_component_model_async(true);
6706 let engine = Engine::new(&config)?;
6707 let component = Component::new(&engine, "(component)")?;
6708 let mut store = Store::new(&engine, ());
6709 let instance = Linker::new(&engine).instantiate(&mut store, &component)?;
6710 let store_opaque = store.as_store_opaque();
6711 let concurrent_state = store_opaque.concurrent_state_mut()?;
6712 let group = concurrent_state.make_task_group()?;
6715 let task = concurrent_state.push(HostTask {
6716 common: WaitableCommon::default(),
6717 call_context: CallContext::default(),
6718 state,
6719 group,
6720 })?;
6721 let handle = store_opaque
6722 .instance_state(instance.runtime_instance(RuntimeComponentInstanceIndex::from_u32(0)))
6723 .handle_table()
6724 .subtask_insert_host(task.rep())?;
6725 let common = &mut store_opaque.concurrent_state_mut()?.get_mut(task)?.common;
6726 common.handle = Some(handle);
6727 common.event = event;
6728 Ok((store, instance, task, handle))
6729 }
6730
6731 #[test]
6732 fn host_subtask_drop_during_cancellation() -> Result<()> {
6733 for abort_completed in [false, true] {
6734 let (handle, future) = JoinHandle::run(future::pending::<()>());
6735 let mut future = pin!(future);
6736 let (mut store, instance, task, handle) =
6737 host_subtask(HostTaskState::CalleeRunning(handle), None)?;
6738 let store = store.as_store_opaque();
6739 let caller = RuntimeComponentInstanceIndex::from_u32(0);
6740 assert_eq!(
6741 instance.subtask_cancel(store, caller, true, handle)?,
6742 BLOCKED
6743 );
6744 if abort_completed {
6745 assert!(matches!(
6748 future
6749 .as_mut()
6750 .poll(&mut Context::from_waker(Waker::noop())),
6751 Poll::Ready(None),
6752 ));
6753 }
6754 for async_ in [false, true] {
6755 let err = instance
6756 .subtask_cancel(store, caller, async_, handle)
6757 .unwrap_err();
6758 assert_eq!(err.downcast::<Trap>()?, Trap::SubtaskCancelAfterTerminal);
6759 }
6760 let err = instance.subtask_drop(store, caller, handle).unwrap_err();
6761 assert_eq!(err.downcast::<Trap>()?, Trap::SubtaskDropNotResolved);
6762 assert!(store.concurrent_state_mut()?.get_mut(task).is_ok());
6763 }
6764 Ok(())
6765 }
6766
6767 #[test]
6768 fn host_subtask_cancel_after_completion() -> Result<()> {
6769 for async_ in [false, true] {
6770 let (mut store, instance, task, handle) = host_subtask(
6771 HostTaskState::CalleeDone { cancelled: false },
6772 Some(Event::Subtask {
6773 status: Status::Returned,
6774 }),
6775 )?;
6776 let store = store.as_store_opaque();
6777 let caller = RuntimeComponentInstanceIndex::from_u32(0);
6778 assert_eq!(
6779 instance.subtask_cancel(store, caller, async_, handle)?,
6780 Status::Returned as u32,
6781 );
6782 let err = instance
6783 .subtask_cancel(store, caller, async_, handle)
6784 .unwrap_err();
6785 assert_eq!(err.downcast::<Trap>()?, Trap::SubtaskCancelAfterTerminal);
6786 instance.subtask_drop(store, caller, handle)?;
6787 assert!(store.concurrent_state_mut()?.get_mut(task).is_err());
6788 }
6789 Ok(())
6790 }
6791
6792 #[test]
6793 fn host_subtask_drop_requires_terminal_event_delivery() -> Result<()> {
6794 for (cancelled, status) in [
6795 (false, Status::Returned),
6796 (true, Status::Returned),
6797 (true, Status::ReturnCancelled),
6798 ] {
6799 for delivered in [false, true] {
6800 let event = if delivered {
6801 None
6802 } else {
6803 Some(Event::Subtask { status })
6804 };
6805 let (mut store, instance, task, handle) =
6806 host_subtask(HostTaskState::CalleeDone { cancelled }, event)?;
6807 let store = store.as_store_opaque();
6808 let result = instance.subtask_drop(
6809 store,
6810 RuntimeComponentInstanceIndex::from_u32(0),
6811 handle,
6812 );
6813 if delivered {
6814 result?;
6815 assert!(store.concurrent_state_mut()?.get_mut(task).is_err());
6816 } else {
6817 let err = result.unwrap_err();
6818 assert_eq!(err.downcast::<Trap>()?, Trap::SubtaskDropNotResolved);
6819 assert!(store.concurrent_state_mut()?.get_mut(task).is_ok());
6820 }
6821 }
6822 }
6823 Ok(())
6824 }
6825}