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::{AsContext, AsContextMut, Result, StoreContext, StoreContextMut, ValRaw, bail};
71use crate::{Instance as ModuleInstance, bail_bug};
72use alloc::borrow::ToOwned;
73use alloc::collections::{BTreeMap, BTreeSet, VecDeque};
74use core::any::Any;
75use core::cell::UnsafeCell;
76use core::fmt;
77use core::future;
78use core::future::Future;
79use core::marker::PhantomData;
80use core::mem::{self, ManuallyDrop, MaybeUninit};
81use core::ops::DerefMut;
82use core::pin::{Pin, pin};
83use core::ptr::{self, NonNull};
84use core::task::{Context, Poll, Waker};
85use futures::channel::oneshot;
86use futures::stream::{FuturesUnordered, StreamExt};
87use futures_and_streams::{FlatAbi, ReturnCode, TransmitHandle, TransmitIndex};
88use table::{TableDebug, TableId};
89use wasmtime_environ::component::{
90 CanonicalAbiInfo, CanonicalOptions, CanonicalOptionsDataModel, MAX_FLAT_PARAMS,
91 MAX_FLAT_RESULTS, OptionsIndex, PREPARE_ASYNC_NO_RESULT, PREPARE_ASYNC_WITH_RESULT,
92 RuntimeComponentInstanceIndex, RuntimeTableIndex, StringEncoding,
93 TypeComponentGlobalErrorContextTableIndex, TypeComponentLocalErrorContextTableIndex,
94 TypeFutureTableIndex, TypeStreamTableIndex, TypeTupleIndex,
95};
96use wasmtime_environ::packed_option::ReservedValue;
97use wasmtime_environ::{ModuleInternedTypeIndex, NUM_COMPONENT_CONTEXT_SLOTS, Trap};
98#[cfg(feature = "gc")]
99use wasmtime_unwinder::Unwind;
100
101pub use abort::JoinHandle;
102pub use func::{FuncCallConcurrent, TypedFuncCallConcurrent};
103pub use future_stream_any::{FutureAny, StreamAny};
104pub use futures_and_streams::{
105 Destination, DirectDestination, DirectSource, ErrorContext, FutureConsumer, FutureProducer,
106 FutureReader, GuardedFutureReader, GuardedStreamReader, ReadBuffer, Source, StreamConsumer,
107 StreamProducer, StreamReader, StreamResult, VecBuffer, WriteBuffer,
108};
109pub(crate) use futures_and_streams::{ResourcePair, lower_error_context_to_index};
110#[cfg(feature = "task-group-hook")]
111pub use task_group_hook::TaskGroupHook;
112pub use task_group_hook::TaskGroupId;
113
114mod abort;
115mod error_contexts;
116mod func;
117mod future_stream_any;
118mod futures_and_streams;
119pub(crate) mod table;
120#[cfg(feature = "task-group-hook")]
121mod task_group_hook;
122#[cfg(not(feature = "task-group-hook"))]
123mod task_group_hook_disabled;
124#[cfg(not(feature = "task-group-hook"))]
125use task_group_hook_disabled as task_group_hook;
126pub(crate) mod tls;
127
128const BLOCKED: u32 = 0xffff_ffff;
131
132#[derive(Clone, Copy, Eq, PartialEq, Debug)]
134pub enum Status {
135 Starting = 0,
136 Started = 1,
137 Returned = 2,
138 StartCancelled = 3,
139 ReturnCancelled = 4,
140}
141
142impl Status {
143 pub fn pack(self, waitable: Option<u32>) -> u32 {
149 assert!(matches!(self, Status::Returned) == waitable.is_none());
150 let waitable = waitable.unwrap_or(0);
151 assert!(waitable < (1 << 28));
152 (waitable << 4) | (self as u32)
153 }
154}
155
156#[derive(Clone, Copy, Debug)]
159enum Event {
160 None,
161 Subtask {
162 status: Status,
163 },
164 StreamRead {
165 code: ReturnCode,
166 pending: Option<(TypeStreamTableIndex, u32)>,
167 },
168 StreamWrite {
169 code: ReturnCode,
170 pending: Option<(TypeStreamTableIndex, u32)>,
171 },
172 FutureRead {
173 code: ReturnCode,
174 pending: Option<(TypeFutureTableIndex, u32)>,
175 },
176 FutureWrite {
177 code: ReturnCode,
178 pending: Option<(TypeFutureTableIndex, u32)>,
179 },
180 Cancelled,
181}
182
183impl Event {
184 fn parts(self) -> (u32, u32) {
189 const EVENT_NONE: u32 = 0;
190 const EVENT_SUBTASK: u32 = 1;
191 const EVENT_STREAM_READ: u32 = 2;
192 const EVENT_STREAM_WRITE: u32 = 3;
193 const EVENT_FUTURE_READ: u32 = 4;
194 const EVENT_FUTURE_WRITE: u32 = 5;
195 const EVENT_CANCELLED: u32 = 6;
196 match self {
197 Event::None => (EVENT_NONE, 0),
198 Event::Cancelled => (EVENT_CANCELLED, 0),
199 Event::Subtask { status } => (EVENT_SUBTASK, status as u32),
200 Event::StreamRead { code, .. } => (EVENT_STREAM_READ, code.encode()),
201 Event::StreamWrite { code, .. } => (EVENT_STREAM_WRITE, code.encode()),
202 Event::FutureRead { code, .. } => (EVENT_FUTURE_READ, code.encode()),
203 Event::FutureWrite { code, .. } => (EVENT_FUTURE_WRITE, code.encode()),
204 }
205 }
206}
207
208mod callback_code {
210 pub const EXIT: u32 = 0;
211 pub const YIELD: u32 = 1;
212 pub const WAIT: u32 = 2;
213}
214
215const START_FLAG_ASYNC_CALLEE: u32 = wasmtime_environ::component::START_FLAG_ASYNC_CALLEE as u32;
217const START_FLAG_ASYNC_CALLER: u32 = wasmtime_environ::component::START_FLAG_ASYNC_CALLER as u32;
218
219pub struct Access<'a, T: 'static, D: HasData + ?Sized = HasSelf<T>> {
225 store: StoreContextMut<'a, T>,
226 get_data: fn(&mut T) -> D::Data<'_>,
227}
228
229impl<'a, T, D> Access<'a, T, D>
230where
231 D: HasData + ?Sized,
232 T: 'static,
233{
234 pub fn new(store: StoreContextMut<'a, T>, get_data: fn(&mut T) -> D::Data<'_>) -> Self {
236 Self { store, get_data }
237 }
238
239 pub fn data_mut(&mut self) -> &mut T {
241 self.store.data_mut()
242 }
243
244 pub fn get(&mut self) -> D::Data<'_> {
246 (self.get_data)(self.data_mut())
247 }
248
249 pub fn spawn(&mut self, task: impl for<'fut> AccessorTask<'fut, T, D>) -> Result<JoinHandle>
253 where
254 T: 'static,
255 {
256 let accessor = Accessor {
257 get_data: self.get_data,
258 token: StoreToken::new(self.store.as_context_mut()),
259 };
260 self.store
261 .as_context_mut()
262 .spawn_with_accessor(accessor, task)
263 }
264
265 pub fn getter(&self) -> fn(&mut T) -> D::Data<'_> {
268 self.get_data
269 }
270}
271
272impl<'a, T, D> AsContext for Access<'a, T, D>
273where
274 D: HasData + ?Sized,
275 T: 'static,
276{
277 type Data = T;
278
279 fn as_context(&self) -> StoreContext<'_, T> {
280 self.store.as_context()
281 }
282}
283
284impl<'a, T, D> AsContextMut for Access<'a, T, D>
285where
286 D: HasData + ?Sized,
287 T: 'static,
288{
289 fn as_context_mut(&mut self) -> StoreContextMut<'_, T> {
290 self.store.as_context_mut()
291 }
292}
293
294pub struct Accessor<T: 'static, D = HasSelf<T>>
354where
355 D: HasData + ?Sized,
356{
357 token: StoreToken<T>,
358 get_data: fn(&mut T) -> D::Data<'_>,
359}
360
361pub trait AsAccessor {
378 type Data: 'static;
380
381 type AccessorData: HasData + ?Sized;
384
385 fn as_accessor(&self) -> &Accessor<Self::Data, Self::AccessorData>;
387}
388
389impl<T: AsAccessor + ?Sized> AsAccessor for &T {
390 type Data = T::Data;
391 type AccessorData = T::AccessorData;
392
393 fn as_accessor(&self) -> &Accessor<Self::Data, Self::AccessorData> {
394 T::as_accessor(self)
395 }
396}
397
398impl<T, D: HasData + ?Sized> AsAccessor for Accessor<T, D> {
399 type Data = T;
400 type AccessorData = D;
401
402 fn as_accessor(&self) -> &Accessor<T, D> {
403 self
404 }
405}
406
407const _: () = {
430 const fn assert<T: Send + Sync>() {}
431 assert::<Accessor<UnsafeCell<u32>>>();
432};
433
434impl<T> Accessor<T> {
435 pub(crate) fn new(token: StoreToken<T>) -> Self {
444 Self {
445 token,
446 get_data: |x| x,
447 }
448 }
449}
450
451impl<T, D> Accessor<T, D>
452where
453 D: HasData + ?Sized,
454{
455 pub fn with<R>(&self, fun: impl FnOnce(Access<'_, T, D>) -> R) -> R {
473 tls::get(|vmstore| {
474 fun(Access {
475 store: self.token.as_context_mut(vmstore),
476 get_data: self.get_data,
477 })
478 })
479 }
480
481 pub fn getter(&self) -> fn(&mut T) -> D::Data<'_> {
484 self.get_data
485 }
486
487 pub fn with_getter<D2: HasData>(
504 &self,
505 get_data: fn(&mut T) -> D2::Data<'_>,
506 ) -> Accessor<T, D2> {
507 Accessor {
508 token: self.token,
509 get_data,
510 }
511 }
512
513 pub fn spawn(&self, task: impl for<'fut> AccessorTask<'fut, T, D>) -> Result<JoinHandle>
529 where
530 T: 'static,
531 {
532 let accessor = self.clone_for_spawn();
533 self.with(|mut access| access.as_context_mut().spawn_with_accessor(accessor, task))
534 }
535
536 fn clone_for_spawn(&self) -> Self {
537 Self {
538 token: self.token,
539 get_data: self.get_data,
540 }
541 }
542
543 pub fn poll_no_interesting_tasks(&self, cx: &mut Context<'_>) -> Poll<()> {
579 self.with(|mut access| {
580 let store = access.as_context_mut().0;
581 let state = store.concurrent_state_mut_without_forcing_current_thread();
582 if state.interesting_tasks == 0 {
583 Poll::Ready(())
584 } else {
585 state.interesting_tasks_empty_waker = Some(cx.waker().clone());
586 Poll::Pending
587 }
588 })
589 }
590
591 pub fn poll_ready_for_concurrent_call(&self, func: Func, cx: &mut Context<'_>) -> Poll<()> {
608 self.with(|mut access| {
609 let store = access.as_context_mut().0;
610 let (_, _, _, raw_options) = func.abi_info(store);
611 let instance = func.instance().runtime_instance(raw_options.instance);
612 let state = store.instance_state(instance).concurrent_state();
613 if state.backpressure == 0 {
614 Poll::Ready(())
615 } else {
616 store
617 .concurrent_state_mut_without_forcing_current_thread()
618 .ready_for_concurrent_call_waker = Some(cx.waker().clone());
619 Poll::Pending
620 }
621 })
622 }
623}
624
625pub trait AccessorTask<'fut, T, D = HasSelf<T>>:
647 AsyncFnOnce(&Accessor<T, D>) -> Result<()> + Send + 'static
648where
649 D: HasData + ?Sized,
650{
651 fn run(self, accessor: &'fut Accessor<T, D>) -> impl Future<Output = Result<()>> + Send + 'fut;
653}
654
655impl<'fut, F, Fut, T, D> AccessorTask<'fut, T, D> for F
656where
657 T: 'static,
658 F: AsyncFnOnce(&Accessor<T, D>) -> Result<()>,
659 F: FnOnce(&'fut Accessor<T, D>) -> Fut + Send + 'static,
660 Fut: Future<Output = Result<()>> + Send + 'fut,
661 D: HasData,
662{
663 fn run(self, accessor: &'fut Accessor<T, D>) -> impl Future<Output = Result<()>> + Send + 'fut {
664 (self)(accessor)
665 }
666}
667
668enum CallerInfo {
671 Async {
673 params: Vec<ValRaw>,
674 has_result: bool,
675 },
676 Sync {
678 params: Vec<ValRaw>,
679 result_count: u32,
680 },
681}
682
683enum WaitMode {
685 Fiber(StoreFiber<'static>),
687 Callback(Instance),
690}
691
692impl fmt::Debug for WaitMode {
693 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
694 match self {
695 Self::Fiber(_) => f.debug_tuple("Fiber").finish(),
696 Self::Callback(instance) => f.debug_tuple("Callback").field(instance).finish(),
697 }
698 }
699}
700
701#[derive(Debug)]
703enum SuspendReason {
704 Waiting {
707 set: TableId<WaitableSet>,
708 thread: QualifiedThreadId,
709 },
710 YieldingToSubtask { thread: QualifiedThreadId },
713 NeedWork,
716 Yielding { thread: QualifiedThreadId },
719 ExplicitlySuspending { thread: QualifiedThreadId },
722}
723
724enum GuestCallKind {
726 DeliverEvent {
729 instance: Instance,
731 set: Option<TableId<WaitableSet>>,
736 },
737 StartImplicit(Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send + Sync>),
743 StartExplicit(Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send + Sync>),
744}
745
746impl fmt::Debug for GuestCallKind {
747 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
748 match self {
749 Self::DeliverEvent { instance, set } => f
750 .debug_struct("DeliverEvent")
751 .field("instance", instance)
752 .field("set", set)
753 .finish(),
754 Self::StartImplicit(_) => f.debug_tuple("StartImplicit").finish(),
755 Self::StartExplicit(_) => f.debug_tuple("StartExplicit").finish(),
756 }
757 }
758}
759
760#[derive(Copy, Clone, Debug)]
762pub enum SuspensionTarget {
763 Resume(u32),
764 Promote(u32),
765 None,
766}
767
768#[derive(Copy, Clone, Debug)]
770pub enum ResumeThread {
771 Promote,
772 Resume,
773 ResumeLater,
774}
775
776#[derive(Debug)]
778struct GuestCall {
779 thread: QualifiedThreadId,
780 kind: GuestCallKind,
781}
782
783impl GuestCall {
784 fn is_ready(&self, store: &mut StoreOpaque) -> Result<bool> {
794 let task = store.concurrent_state_mut()?.get_mut(self.thread.task)?;
795 let async_typed = task.async_typed;
796 let instance = task.instance;
797 let state = store.instance_state(instance).concurrent_state();
798
799 let ready = match &self.kind {
800 GuestCallKind::DeliverEvent { .. } => !state.do_not_enter,
801 GuestCallKind::StartImplicit(_) => {
802 !async_typed || !(state.do_not_enter || state.backpressure > 0)
803 }
804 GuestCallKind::StartExplicit(_) => true,
805 };
806 log::trace!(
807 "call {self:?} ready? {ready} (do_not_enter: {}; backpressure: {})",
808 state.do_not_enter,
809 state.backpressure
810 );
811 Ok(ready)
812 }
813}
814
815enum WorkerItem {
817 GuestCall(GuestCall),
818 Function(AlwaysMut<Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send>>),
819}
820
821enum WorkItem {
824 PushFuture(AlwaysMut<HostTaskFuture>),
826 ResumeFiber {
828 instance: RuntimeInstance,
829 thread: QualifiedThreadId,
830 fiber: StoreFiber<'static>,
831 },
832 ResumeThread {
834 instance: RuntimeInstance,
835 thread: QualifiedThreadId,
836 },
837 GuestCall {
839 instance: RuntimeInstance,
840 call: GuestCall,
841 },
842 WorkerFunction(AlwaysMut<Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send>>),
844}
845
846impl fmt::Debug for WorkItem {
847 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
848 match self {
849 Self::PushFuture(_) => f.debug_tuple("PushFuture").finish(),
850 Self::ResumeFiber {
851 instance, thread, ..
852 } => f
853 .debug_struct("ResumeFiber")
854 .field("instance", instance)
855 .field("thread", thread)
856 .finish(),
857 Self::ResumeThread { instance, thread } => f
858 .debug_struct("ResumeThread")
859 .field("instance", instance)
860 .field("thread", thread)
861 .finish(),
862 Self::GuestCall { instance, call } => f
863 .debug_struct("GuestCall")
864 .field("instance", instance)
865 .field("call", call)
866 .finish(),
867 Self::WorkerFunction(_) => f.debug_tuple("WorkerFunction").finish(),
868 }
869 }
870}
871
872#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
874pub(crate) enum WaitResult {
875 Cancelled,
876 Completed,
877}
878
879pub(crate) fn poll_and_block<R: Send + Sync + 'static>(
887 store: &mut dyn VMStore,
888 host_task: EnteredHostTask,
889 future: impl Future<Output = Result<R>> + Send + 'static,
890) -> Result<R> {
891 let mut future = Box::pin(future);
898 let poll = tls::set(store, || {
899 future
900 .as_mut()
901 .poll(&mut Context::from_waker(&Waker::noop()))
902 });
903
904 let caller = match host_task {
905 Some(caller) => caller,
906 None => bail_bug!("host task wasn't created but should have been"),
907 };
908
909 let task = match poll {
910 Poll::Ready(result) => return result,
912
913 Poll::Pending => {
918 let Some(task) = store.materialize_host_task_id()? else {
919 bail_bug!("current thread is not a host thread")
920 };
921
922 let future = Box::pin(async move {
925 let result = run_with_host_task_set(task, future).await??;
926 tls::get(move |store| {
927 let state = store.concurrent_state_mut()?;
928 let host_state = &mut state.get_mut(task)?.state;
929 assert!(matches!(host_state, HostTaskState::CalleeStarted));
930 *host_state = HostTaskState::CalleeFinished(Box::new(result));
931
932 Waitable::Host(task).set_event(
933 state,
934 Some(Event::Subtask {
935 status: Status::Returned,
936 }),
937 )?;
938
939 Ok(())
940 })
941 }) as HostTaskFuture;
942
943 let caller_instance = store.concurrent_state_mut()?.get_mut(caller.task)?.instance;
944 store.switch_or_trap_if_may_not_suspend(caller_instance)?;
945
946 let state = store.concurrent_state_mut()?;
947 state.push_future(future);
948
949 let set = state.get_mut(caller.thread)?.sync_call_set;
950 Waitable::Host(task).join(state, Some(set))?;
951
952 store.suspend(SuspendReason::Waiting {
953 set,
954 thread: caller,
955 })?;
956
957 Waitable::Host(task).join(store.concurrent_state_mut()?, None)?;
961 task
962 }
963 };
964
965 let host_state = &mut store.concurrent_state_mut()?.get_mut(task)?.state;
967 match mem::replace(host_state, HostTaskState::CalleeDone { cancelled: false }) {
968 HostTaskState::CalleeFinished(result) => Ok(match result.downcast() {
969 Ok(result) => *result,
970 Err(_) => bail_bug!("host task finished with wrong type of result"),
971 }),
972 _ => bail_bug!("unexpected host task state after completion"),
973 }
974}
975
976fn handle_guest_call(store: &mut dyn VMStore, call: GuestCall) -> Result<()> {
978 match call.kind {
979 GuestCallKind::DeliverEvent { instance, set } => {
980 if let Some(set) = set {
983 store.concurrent_state_mut()?.get_mut(set)?.stop_waiting()?;
984 }
985 let (event, waitable) = match instance.get_event(store, call.thread.task, set, true)? {
986 Some(pair) => pair,
987 None => match set {
988 Some(set) => {
993 log::trace!(
994 "event for {:?} on {set:?} no longer present; waiting again",
995 call.thread
996 );
997 return instance.wait_with_callback(
998 store.store_opaque_mut(),
999 call.thread,
1000 set,
1001 );
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.saved_next_switch_items.is_empty());
1127 assert!(state.unforced_current_thread.is_none());
1128 assert!(state.deferred_host_call_context.is_none());
1129 assert!(state.futures_mut().unwrap().is_empty());
1130 assert!(state.global_error_context_ref_counts.is_empty());
1131 }
1132
1133 #[doc(hidden)]
1138 pub fn concurrent_state_table_size(&mut self) -> usize {
1139 self.0
1140 .concurrent_state_mut()
1141 .unwrap()
1142 .table
1143 .get_mut()
1144 .iter_mut()
1145 .count()
1146 }
1147
1148 pub fn spawn(mut self, task: impl for<'fut> AccessorTask<'fut, T>) -> Result<JoinHandle>
1158 where
1159 T: 'static,
1160 {
1161 let accessor = Accessor::new(StoreToken::new(self.as_context_mut()));
1162 self.spawn_with_accessor(accessor, task)
1163 }
1164
1165 fn spawn_with_accessor<D>(
1168 self,
1169 accessor: Accessor<T, D>,
1170 task: impl for<'fut> AccessorTask<'fut, T, D>,
1171 ) -> Result<JoinHandle>
1172 where
1173 T: 'static,
1174 D: HasData + ?Sized,
1175 {
1176 let (handle, future) = JoinHandle::run(async move { task.run(&accessor).await });
1180 self.0
1181 .concurrent_state_mut()?
1182 .push_future(Box::pin(async move { future.await.unwrap_or(Ok(())) }));
1183 Ok(handle)
1184 }
1185
1186 pub async fn run_concurrent<R>(self, fun: impl AsyncFnOnce(&Accessor<T>) -> R) -> Result<R>
1270 where
1271 T: Send + 'static,
1272 {
1273 ensure!(
1274 self.0.concurrency_support(),
1275 "cannot use `run_concurrent` when Config::concurrency_support disabled",
1276 );
1277 self.do_run_concurrent(fun, false).await
1278 }
1279
1280 pub(super) async fn run_concurrent_trap_on_idle<R>(
1281 self,
1282 fun: impl AsyncFnOnce(&Accessor<T>) -> R,
1283 ) -> Result<R> {
1284 self.do_run_concurrent(fun, true).await
1285 }
1286
1287 async fn do_run_concurrent<R>(
1288 mut self,
1289 fun: impl AsyncFnOnce(&Accessor<T>) -> R,
1290 trap_on_idle: bool,
1291 ) -> Result<R> {
1292 debug_assert!(self.0.concurrency_support());
1293 let already_running = self
1294 .0
1295 .concurrent_state_mut_already_forced_current_thread()
1296 .event_loop_running;
1297 if already_running {
1298 bail!("Recursive `StoreContextMut::run_concurrent` calls not supported")
1299 }
1300 let token = StoreToken::new(self.as_context_mut());
1301
1302 struct Dropper<'a, T: 'static, V> {
1303 store: StoreContextMut<'a, T>,
1304 value: ManuallyDrop<V>,
1305 }
1306
1307 impl<'a, T, V> Drop for Dropper<'a, T, V> {
1308 fn drop(&mut self) {
1309 self.store
1310 .0
1311 .concurrent_state_mut_already_forced_current_thread()
1312 .event_loop_running = false;
1313
1314 tls::set(self.store.0, || {
1315 unsafe { ManuallyDrop::drop(&mut self.value) }
1320 });
1321 }
1322 }
1323
1324 let accessor = &Accessor::new(token);
1325 self.0
1326 .concurrent_state_mut_already_forced_current_thread()
1327 .event_loop_running = true;
1328 let dropper = &mut Dropper {
1329 store: self,
1330 value: ManuallyDrop::new(fun(accessor)),
1331 };
1332 let future = unsafe { Pin::new_unchecked(dropper.value.deref_mut()) };
1334
1335 let result = dropper
1336 .store
1337 .as_context_mut()
1338 .poll_until(future, trap_on_idle)
1339 .await;
1340
1341 if result.is_err() {
1342 dropper.store.0.set_trapped();
1343 }
1344
1345 result
1346 }
1347
1348 async fn poll_until<R>(
1354 mut self,
1355 mut future: Pin<&mut impl Future<Output = R>>,
1356 trap_on_idle: bool,
1357 ) -> Result<R> {
1358 struct Reset<'a, T: 'static> {
1359 store: StoreContextMut<'a, T>,
1360 futures: Option<FuturesUnordered<HostTaskFuture>>,
1361 }
1362
1363 impl<'a, T> Drop for Reset<'a, T> {
1364 fn drop(&mut self) {
1365 if let Some(futures) = self.futures.take() {
1366 *self
1367 .store
1368 .0
1369 .concurrent_state_mut_already_forced_current_thread()
1370 .futures
1371 .get_mut() = Some(futures);
1372 }
1373 }
1374 }
1375
1376 const MAX_TURNS_WITHOUT_YIELD: usize = 128;
1380 let mut turns_without_yield = 0;
1381
1382 loop {
1383 let futures = self.0.concurrent_state_mut()?.futures.get_mut().take();
1387 let mut reset = Reset {
1388 store: self.as_context_mut(),
1389 futures,
1390 };
1391 let mut next = match reset.futures.as_mut() {
1392 Some(f) => pin!(f.next()),
1393 None => bail_bug!("concurrent state missing futures field"),
1394 };
1395
1396 enum PollResult<R> {
1397 Complete(R),
1398 ProcessWork {
1399 ready: Option<WorkItem>,
1400 low_priority: bool,
1401 },
1402 }
1403
1404 let result = future::poll_fn(|cx| {
1405 if let Poll::Ready(value) = tls::set(reset.store.0, || future.as_mut().poll(cx)) {
1408 return Poll::Ready(Ok(PollResult::Complete(value)));
1409 }
1410
1411 if reset.store.0.trapped() {
1420 return Poll::Ready(Err(Trap::CannotEnterComponent.into()));
1421 }
1422
1423 let next = match tls::set(reset.store.0, || next.as_mut().poll(cx)) {
1427 Poll::Ready(Some(output)) => {
1428 match output {
1429 Err(e) => return Poll::Ready(Err(e)),
1430 Ok(()) => {}
1431 }
1432 Poll::Ready(true)
1433 }
1434 Poll::Ready(None) => Poll::Ready(false),
1435 Poll::Pending => Poll::Pending,
1436 };
1437
1438 let state = reset.store.0.concurrent_state_mut()?;
1453 let mut ready = state.switch_item.take();
1454 let mut low_priority = false;
1455 if ready.is_none() {
1456 ready = state.high_priority.pop_back();
1457 if ready.is_none() {
1458 ready = state.low_priority.pop_back();
1459 low_priority = true;
1460 }
1461 }
1462 if ready.is_some() {
1463 return Poll::Ready(Ok(PollResult::ProcessWork {
1464 ready,
1465 low_priority,
1466 }));
1467 }
1468
1469 return match next {
1473 Poll::Ready(true) => {
1474 Poll::Ready(Ok(PollResult::ProcessWork {
1480 ready: None,
1481 low_priority: false,
1482 }))
1483 }
1484 Poll::Ready(false) => {
1485 if let Poll::Ready(value) =
1489 tls::set(reset.store.0, || future.as_mut().poll(cx))
1490 {
1491 Poll::Ready(Ok(PollResult::Complete(value)))
1492 } else {
1493 if trap_on_idle {
1499 Poll::Ready(Err(if reset.store.0.any_may_not_suspend()? {
1506 Trap::CannotBlockSyncTask.into()
1507 } else {
1508 Trap::AsyncDeadlock.into()
1510 }))
1511 } else {
1512 Poll::Pending
1516 }
1517 }
1518 }
1519 Poll::Pending => Poll::Pending,
1524 };
1525 })
1526 .await;
1527
1528 drop(reset);
1532
1533 match result? {
1534 PollResult::Complete(value) => break Ok(value),
1537 PollResult::ProcessWork {
1540 ready,
1541 low_priority,
1542 } => {
1543 struct Dispose<'a, T: 'static> {
1544 store: StoreContextMut<'a, T>,
1545 ready: Option<WorkItem>,
1546 }
1547
1548 impl<'a, T> Drop for Dispose<'a, T> {
1549 fn drop(&mut self) {
1550 if let Some(item) = self.ready.take() {
1551 match item {
1552 WorkItem::ResumeFiber { mut fiber, .. } => {
1553 fiber.dispose(self.store.0);
1554 }
1555 WorkItem::PushFuture(future) => {
1556 tls::set(self.store.0, move || drop(future))
1557 }
1558 _ => {}
1559 }
1560 }
1561 }
1562 }
1563
1564 let mut dispose = Dispose {
1565 store: self.as_context_mut(),
1566 ready,
1567 };
1568
1569 if low_priority {
1591 dispose.store.0.yield_now().await;
1592 turns_without_yield = 0;
1593 }
1594
1595 if let Some(item) = dispose.ready.take() {
1596 dispose
1597 .store
1598 .as_context_mut()
1599 .handle_work_item(item)
1600 .await?;
1601 }
1602
1603 turns_without_yield += 1;
1604 if turns_without_yield == MAX_TURNS_WITHOUT_YIELD {
1605 turns_without_yield = 0;
1606 dispose.store.0.yield_now().await;
1607 }
1608 }
1609 }
1610 }
1611 }
1612
1613 async fn handle_work_item(self, item: WorkItem) -> Result<()> {
1615 log::trace!("handle work item {item:?}");
1616 match item {
1617 WorkItem::PushFuture(future) => {
1618 self.0
1619 .concurrent_state_mut()?
1620 .futures_mut()?
1621 .push(future.into_inner());
1622 }
1623 WorkItem::ResumeFiber { fiber, .. } => {
1624 self.0.resume_fiber(fiber).await?;
1625 }
1626 WorkItem::ResumeThread { thread, .. } => {
1627 if let GuestThreadState::Ready { fiber, .. } = mem::replace(
1628 &mut self.0.concurrent_state_mut()?.get_mut(thread.thread)?.state,
1629 GuestThreadState::Running,
1630 ) {
1631 self.0.resume_fiber(fiber).await?;
1632 } else {
1633 bail_bug!("cannot resume non-pending thread {thread:?}");
1634 }
1635 }
1636 WorkItem::GuestCall { call, .. } => {
1637 if call.is_ready(self.0)? {
1638 self.0
1639 .concurrent_state_mut()?
1640 .get_mut(call.thread.thread)?
1641 .wake_on_cancel = WakeOnCancel::None;
1642 self.run_on_worker(WorkerItem::GuestCall(call)).await?;
1643 } else {
1644 let state = self.0.concurrent_state_mut()?;
1645 let task = state.get_mut(call.thread.task)?;
1646 if !task.starting_sent {
1647 task.starting_sent = true;
1648 if let GuestCallKind::StartImplicit(_) = &call.kind {
1649 Waitable::Guest(call.thread.task).set_event(
1650 state,
1651 Some(Event::Subtask {
1652 status: Status::Starting,
1653 }),
1654 )?;
1655 }
1656 }
1657
1658 let instance = state.get_mut(call.thread.task)?.instance;
1659 if let GuestCallKind::DeliverEvent { set: Some(set), .. } = &call.kind {
1668 let state = self.0.concurrent_state_mut()?;
1669 if !state.get_mut(*set)?.ready.is_empty() {
1670 state.wake_waiter(*set)?;
1671 }
1672 }
1673 self.0
1674 .instance_state(instance)
1675 .concurrent_state()
1676 .pending
1677 .insert(call.thread, call.kind);
1678
1679 self.0.concurrent_state_mut()?.take_next_switch_item()?;
1683 }
1684 }
1685 WorkItem::WorkerFunction(fun) => {
1686 self.run_on_worker(WorkerItem::Function(fun)).await?;
1687 }
1688 }
1689
1690 Ok(())
1691 }
1692
1693 async fn run_on_worker(self, item: WorkerItem) -> Result<()> {
1695 let worker = if let Some(fiber) = self.0.concurrent_state_mut()?.worker.take() {
1696 fiber
1697 } else {
1698 unsafe {
1717 fiber::make_fiber_unchecked(self.0, move |store| {
1718 loop {
1719 let Some(item) = store.concurrent_state_mut()?.worker_item.take() else {
1720 bail_bug!("worker_item not present when resuming fiber")
1721 };
1722 match item {
1723 WorkerItem::GuestCall(call) => handle_guest_call(store, call)?,
1724 WorkerItem::Function(fun) => fun.into_inner()(store)?,
1725 }
1726
1727 store.suspend(SuspendReason::NeedWork)?;
1728 }
1729 })?
1730 }
1731 };
1732
1733 let worker_item = &mut self.0.concurrent_state_mut()?.worker_item;
1734 assert!(worker_item.is_none());
1735 *worker_item = Some(item);
1736
1737 self.0.resume_fiber(worker).await
1738 }
1739
1740 pub(crate) fn wrap_call<F, R>(self, closure: F) -> impl Future<Output = Result<R>> + 'static
1745 where
1746 T: 'static,
1747 F: FnOnce(&Accessor<T>) -> Pin<Box<dyn Future<Output = Result<R>> + Send + '_>>
1748 + Send
1749 + Sync
1750 + 'static,
1751 R: Send + Sync + 'static,
1752 {
1753 let token = StoreToken::new(self);
1754 async move {
1755 let mut accessor = Accessor::new(token);
1756 closure(&mut accessor).await
1757 }
1758 }
1759
1760 pub(crate) async fn start_instance(
1761 &mut self,
1762 instance: ModuleInstance,
1763 callee: Option<RuntimeInstance>,
1764 ) -> Result<ModuleInstance> {
1765 let (tx, rx) = oneshot::channel();
1766 let token = StoreToken::new(self.as_context_mut());
1767 self.0.queue_task(move |store| {
1768 _ = tx.send(
1769 super::instance::start_raw(&mut token.as_context_mut(store), instance, callee)
1770 .map(|()| instance),
1771 );
1772 Ok(())
1773 })?;
1774 self.as_context_mut()
1775 .run_concurrent_trap_on_idle(async |_| {
1776 rx.await
1777 .map_err(|_| format_err!("oneshot channel canceled"))
1778 })
1779 .await??
1780 }
1781}
1782
1783pub type EnteredHostTask = Option<QualifiedThreadId>;
1790
1791impl StoreOpaque {
1792 #[inline]
1796 pub(crate) fn current_thread(&mut self) -> Result<CurrentThread> {
1797 if !self.concurrency_support() {
1799 return Ok(CurrentThread::None);
1800 }
1801
1802 if !self
1805 .vm_store_context_mut()
1806 .current_thread_mut()
1807 .is_deferred()
1808 {
1809 return Ok(self
1810 .concurrent_state_mut_already_forced_current_thread()
1811 .unforced_current_thread);
1812 }
1813
1814 self.force_deferred_current_thread()
1815 }
1816
1817 #[cold]
1820 fn force_deferred_current_thread(&mut self) -> Result<CurrentThread> {
1821 let state = self.concurrent_state_mut_without_forcing_current_thread();
1830 let id = match state.unforced_current_thread.guest_task() {
1831 Some(task) => state.get_mut(task)?.instance.instance,
1832 None => bail_bug!("deferred component-model thread with non-guest base"),
1833 };
1834
1835 let mut frames = Vec::new();
1838 let mut cur = *self.vm_store_context_mut().current_thread_mut();
1839 while let Some(ptr) = cur.as_deferred() {
1840 let deferred = unsafe { ptr.as_non_null().as_ref() };
1845 frames.push((
1846 deferred.callee_async != 0,
1847 deferred.callee_instance,
1848 deferred.saved_context,
1849 ));
1850 cur = deferred.parent;
1851 }
1852
1853 *self.vm_store_context_mut().current_thread_mut() = VMLazyThread::forced();
1857
1858 let current_context = *self.vm_store_context_mut().component_context_mut();
1861
1862 for (callee_async, callee_instance, saved_context) in frames.into_iter().rev() {
1866 *self.vm_store_context_mut().component_context_mut() = saved_context;
1870 let callee = RuntimeInstance {
1871 instance: id,
1872 index: RuntimeComponentInstanceIndex::from_u32(callee_instance),
1873 };
1874 self.enter_guest_sync_call(callee_async, callee)?;
1875 }
1876
1877 *self.vm_store_context_mut().component_context_mut() = current_context;
1879
1880 Ok(self
1881 .concurrent_state_mut_without_forcing_current_thread()
1882 .unforced_current_thread)
1883 }
1884
1885 fn current_guest_thread(&mut self) -> Result<QualifiedThreadId> {
1886 match self.current_thread()?.guest() {
1887 Some(id) => Ok(*id),
1888 None => bail_bug!("current thread is not a guest thread"),
1889 }
1890 }
1891
1892 pub(crate) fn current_materialized_host_task(&mut self) -> Result<Option<TableId<HostTask>>> {
1896 match self.current_thread()? {
1897 CurrentThread::Host(id) => Ok(Some(id)),
1898 CurrentThread::DeferredHost(_) | CurrentThread::None => Ok(None),
1899 _ => bail_bug!("current thread is not a host thread"),
1900 }
1901 }
1902
1903 fn materialize_host_task_id(&mut self) -> Result<Option<TableId<HostTask>>> {
1906 Ok(self
1907 .concurrent_state_mut()?
1908 .materialize_current_host_task_id()?)
1909 }
1910
1911 fn enter_sync_call(&mut self, callee: RuntimeInstance) -> Result<()> {
1912 log::trace!("enter sync-typed call {callee:?}");
1913 let state = self.instance_state(callee).concurrent_state();
1914 let old_do_not_suspend = state.do_not_suspend;
1915 state.do_not_suspend = true;
1916
1917 let thread = self.current_guest_thread()?;
1918 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
1919 if thread.old_do_not_suspend.is_some() {
1920 bail_bug!("current thread already has `old_do_not_suspend` value");
1921 }
1922
1923 thread.old_do_not_suspend = Some(old_do_not_suspend);
1924
1925 Ok(())
1926 }
1927
1928 fn exit_sync_call(&mut self, callee: RuntimeInstance) -> Result<()> {
1929 log::trace!("exit sync-typed call {callee:?}");
1930 let thread = self.current_guest_thread()?;
1931 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
1932 let Some(old_do_not_suspend) = thread.old_do_not_suspend.take() else {
1933 bail_bug!("current thread missing `old_do_not_suspend` value");
1934 };
1935 let state = self.instance_state(callee).concurrent_state();
1936 state.do_not_suspend = old_do_not_suspend;
1937 Ok(())
1938 }
1939
1940 pub(crate) fn enter_guest_sync_call(
1952 &mut self,
1953 callee_async_typed: bool,
1954 callee: RuntimeInstance,
1955 ) -> Result<()> {
1956 log::trace!("enter sync-lifted call {callee:?}");
1957 if !self.concurrency_support() {
1958 return self.enter_call_not_concurrent();
1959 }
1960
1961 let state = self.concurrent_state_mut()?;
1965 let item = state.next_switch_item.take();
1966 state.saved_next_switch_items.push(item);
1967
1968 let thread = self.current_thread()?;
1969 let caller = if let Some(thread) = thread.guest() {
1970 Caller::Guest { thread: *thread }
1971 } else {
1972 Caller::Host {
1973 tx: None,
1974 host_future_present: false,
1975 caller: self.materialize_host_task_id()?,
1976 }
1977 };
1978
1979 let state = self.concurrent_state_mut()?;
1980 let guest_thread = GuestTask::new(
1981 state,
1982 Box::new(move |_, _| bail_bug!("cannot lower params in sync call")),
1983 LiftResult {
1984 lift: Box::new(move |_, _| bail_bug!("cannot lift result in sync call")),
1985 ty: TypeTupleIndex::reserved_value(),
1986 memory: None,
1987 string_encoding: StringEncoding::Utf8,
1988 },
1989 caller,
1990 None,
1991 callee,
1992 callee_async_typed,
1993 false,
1994 )?;
1995
1996 Instance::from_wasmtime(self, callee.instance).add_guest_thread_to_instance_table(
1997 guest_thread.thread,
1998 self,
1999 callee.index,
2000 )?;
2001 self.set_thread(guest_thread)?;
2002
2003 if !callee_async_typed {
2004 self.enter_sync_call(callee)?;
2005 }
2006
2007 Ok(())
2008 }
2009
2010 pub(crate) fn exit_guest_sync_call(&mut self) -> Result<()> {
2018 if !self.concurrency_support() {
2019 return Ok(self.exit_call_not_concurrent());
2020 }
2021
2022 let thread = match self.current_thread()?.guest() {
2023 Some(t) => *t,
2024 None => bail_bug!("expected task when exiting"),
2025 };
2026 let task = self.concurrent_state_mut()?.get_mut(thread.task)?;
2027 let instance = task.instance;
2028
2029 let caller = match &task.caller {
2030 &Caller::Guest { thread } => thread.into(),
2031 &Caller::Host { caller, .. } => caller
2032 .map(CurrentThread::Host)
2033 .unwrap_or(CurrentThread::None),
2034 };
2035 task.lift_result = None;
2036 task.exited = true;
2037 let async_typed = task.async_typed;
2038
2039 if !async_typed {
2040 self.exit_sync_call(instance)?;
2041 }
2042
2043 self.set_thread(caller)?;
2044
2045 log::trace!("exit sync-lifted call {instance:?}");
2046
2047 if async_typed {
2048 self.switch_or_trap_if_may_not_suspend(instance)?;
2053 }
2054
2055 self.cleanup_thread(thread, instance, CleanupTask::Yes)?;
2056
2057 let state = self.concurrent_state_mut()?;
2058 let Some(item) = state.saved_next_switch_items.pop() else {
2059 bail_bug!("unable to pop from `saved_next_switch_items`");
2060 };
2061 if let Some(item) = mem::replace(&mut state.next_switch_item, item) {
2062 state.push_high_priority(item);
2065 bail_bug!("`next_switch_item` unexpectedly already set");
2066 }
2067
2068 Ok(())
2069 }
2070
2071 pub(crate) fn host_task_create(&mut self) -> Result<EnteredHostTask> {
2078 if !self.concurrency_support() {
2079 self.enter_call_not_concurrent()?;
2080 return Ok(None);
2081 }
2082 let caller = self.current_guest_thread()?;
2083 log::trace!("new deferred host task with caller {caller:?}");
2084
2085 self.set_thread(CurrentThread::DeferredHost(caller))?;
2086 let state = self.concurrent_state_mut()?;
2087 debug_assert!(state.deferred_host_call_context.is_none());
2088 state.deferred_host_call_context = Some(CallContext::default());
2089 state.debug_assert_deferred_host_invariant();
2090 Ok(Some(caller))
2091 }
2092
2093 pub(crate) fn host_task_delete(
2100 &mut self,
2101 original_task: EnteredHostTask,
2102 materialized_task: Option<TableId<HostTask>>,
2103 ) -> Result<()> {
2104 match original_task {
2105 Some(caller) => {
2106 self.set_thread(caller)?;
2107 if materialized_task.is_none() {
2108 let state = self.concurrent_state_mut()?;
2109 let context = state
2110 .deferred_host_call_context
2111 .take()
2112 .expect("deferred host call context should be present");
2113 debug_assert!(context.is_empty());
2114 state.debug_assert_deferred_host_invariant();
2115 }
2116 log::trace!(
2117 "delete host task with caller {original_task:?} and materialized as {materialized_task:?}"
2118 );
2119 if let Some(task) = materialized_task {
2120 Waitable::Host(task).delete_from(self)?;
2121 }
2122 }
2123 None => {
2124 debug_assert!(materialized_task.is_none());
2125 self.exit_call_not_concurrent();
2126 }
2127 }
2128 Ok(())
2129 }
2130
2131 fn instance_state(&mut self, instance: RuntimeInstance) -> &mut InstanceState {
2134 self.component_instance_mut(instance.instance)
2135 .instance_state(instance.index)
2136 }
2137
2138 pub(crate) fn set_thread(&mut self, thread: impl Into<CurrentThread>) -> Result<CurrentThread> {
2144 let thread = thread.into();
2145 let state = self.concurrent_state_mut()?;
2146 state.debug_assert_deferred_host_invariant();
2147 let old_thread = mem::replace(&mut state.unforced_current_thread, thread);
2148
2149 state.handle_thread_switch(old_thread, thread)?;
2150
2151 if let Some(old_thread) = old_thread.guest() {
2159 let old_context = *self.vm_store_context_mut().component_context_mut();
2160 self.concurrent_state_mut()?
2161 .get_mut(old_thread.thread)?
2162 .context = old_context;
2163 }
2164 if cfg!(debug_assertions) {
2165 *self.vm_store_context_mut().component_context_mut() =
2166 [u32::MAX; NUM_COMPONENT_CONTEXT_SLOTS];
2167 }
2168 if let Some(thread) = thread.guest() {
2169 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
2170 let context = thread.context;
2171 if cfg!(debug_assertions) {
2172 thread.context = [u32::MAX; NUM_COMPONENT_CONTEXT_SLOTS];
2173 }
2174 *self.vm_store_context_mut().component_context_mut() = context;
2175 }
2176
2177 *self.vm_store_context_mut().current_thread_mut() = if thread.is_none() {
2179 VMLazyThread::none()
2180 } else {
2181 VMLazyThread::forced()
2182 };
2183
2184 Ok(old_thread)
2185 }
2186
2187 fn switch_or_trap_if_may_not_suspend(&mut self, instance: RuntimeInstance) -> Result<()> {
2189 if self.switch_if_may_not_suspend(instance)? {
2190 Ok(())
2191 } else {
2192 Err(Trap::CannotBlockSyncTask.into())
2193 }
2194 }
2195
2196 fn switch_if_may_not_suspend(&mut self, instance: RuntimeInstance) -> Result<bool> {
2200 self.concurrent_state_mut()?;
2204
2205 Ok(!self.concurrency_support()
2206 || !self
2207 .instance_state(instance)
2208 .concurrent_state()
2209 .do_not_suspend
2210 || self
2211 .concurrent_state_mut()?
2212 .promote_instance_local_thread_work_item(instance)?)
2213 }
2214
2215 fn enter_instance(&mut self, instance: RuntimeInstance) {
2219 log::trace!("enter {instance:?}");
2220 self.instance_state(instance)
2221 .concurrent_state()
2222 .do_not_enter = true;
2223 }
2224
2225 fn exit_instance(&mut self, instance: RuntimeInstance) -> Result<()> {
2229 log::trace!("exit {instance:?}");
2230 self.instance_state(instance)
2231 .concurrent_state()
2232 .do_not_enter = false;
2233 self.partition_pending(instance)
2234 }
2235
2236 fn partition_pending(&mut self, instance: RuntimeInstance) -> Result<()> {
2244 for (thread, kind) in
2245 mem::take(&mut self.instance_state(instance).concurrent_state().pending).into_iter()
2246 {
2247 let call = GuestCall { thread, kind };
2248 if call.is_ready(self)? {
2249 self.concurrent_state_mut()?
2250 .push_high_priority(WorkItem::GuestCall { instance, call });
2251 } else {
2252 self.instance_state(instance)
2253 .concurrent_state()
2254 .pending
2255 .insert(call.thread, call.kind);
2256 }
2257 }
2258
2259 if let Some(waker) = self
2260 .concurrent_state_mut()?
2261 .ready_for_concurrent_call_waker
2262 .take()
2263 {
2264 waker.wake();
2265 }
2266
2267 Ok(())
2268 }
2269
2270 pub(crate) fn backpressure_modify(
2272 &mut self,
2273 caller_instance: RuntimeInstance,
2274 modify: impl FnOnce(u16) -> Option<u16>,
2275 ) -> Result<()> {
2276 let state = self.instance_state(caller_instance).concurrent_state();
2277 let old = state.backpressure;
2278 let new = modify(old).ok_or_else(|| Trap::BackpressureOverflow)?;
2279 state.backpressure = new;
2280
2281 if old > 0 && new == 0 {
2282 self.partition_pending(caller_instance)?;
2285 }
2286
2287 Ok(())
2288 }
2289
2290 async fn resume_fiber(&mut self, fiber: StoreFiber<'static>) -> Result<()> {
2293 let old_thread = self.current_thread()?;
2294 log::trace!("resume_fiber: save current thread {old_thread:?}");
2295
2296 let fiber = fiber::resolve_or_release(self, fiber).await?;
2297
2298 self.set_thread(old_thread)?;
2299
2300 let state = self.concurrent_state_mut()?;
2301
2302 if let Some(ot) = old_thread.guest() {
2303 state.get_mut(ot.thread)?.state = GuestThreadState::Running;
2304 }
2305 log::trace!("resume_fiber: restore current thread {old_thread:?}");
2306
2307 if let Some(mut fiber) = fiber {
2308 log::trace!("resume_fiber: suspend reason {:?}", &state.suspend_reason);
2309 let reason = match state.suspend_reason.take() {
2311 Some(r) => r,
2312 None => bail_bug!("suspend reason missing when resuming fiber"),
2313 };
2314 match reason {
2315 SuspendReason::NeedWork => {
2316 if state.worker.is_none() {
2317 state.worker = Some(fiber);
2318 } else {
2319 fiber.dispose(self);
2320 }
2321 }
2322 SuspendReason::Yielding { thread } => {
2323 state.get_mut(thread.thread)?.state = GuestThreadState::Ready { fiber };
2324 let instance = state.get_mut(thread.task)?.instance;
2325 state.push_low_priority(WorkItem::ResumeThread { instance, thread });
2326 }
2327 SuspendReason::ExplicitlySuspending { thread } => {
2328 state.get_mut(thread.thread)?.state = GuestThreadState::Suspended(fiber);
2329 }
2330 SuspendReason::Waiting { set, thread } => {
2331 let old = state
2332 .get_mut(set)?
2333 .waiting
2334 .insert(thread, WaitMode::Fiber(fiber));
2335 assert!(old.is_none());
2336 }
2337 SuspendReason::YieldingToSubtask { thread } => {
2338 let item = WorkItem::ResumeFiber {
2347 instance: state.get_mut(thread.task)?.instance,
2348 thread,
2349 fiber,
2350 };
2351
2352 log::trace!("set next switch item to {item:?}");
2353 if state.next_switch_item.replace(item).is_some() {
2354 bail_bug!(
2357 "`ConcurrentState::next_switch_item` was already `Some(_)` when \
2358 a thread wanted to wait on a subtask"
2359 );
2360 }
2361 }
2362 };
2363 } else {
2364 log::trace!("resume_fiber: fiber has exited");
2365 }
2366
2367 Ok(())
2368 }
2369
2370 fn suspend(&mut self, reason: SuspendReason) -> Result<()> {
2376 log::trace!("suspend fiber: {reason:?}");
2377
2378 let state = self.concurrent_state_mut()?;
2379
2380 let (save_and_restore_thread, save_and_restore_next_switch_item) = match &reason {
2387 SuspendReason::Yielding { .. }
2388 | SuspendReason::Waiting { .. }
2389 | SuspendReason::ExplicitlySuspending { .. } => {
2390 if state.switch_item.is_none() {
2393 state.take_next_switch_item()?;
2394 }
2395
2396 (true, false)
2397 }
2398 SuspendReason::YieldingToSubtask { .. } => (true, true),
2399 SuspendReason::NeedWork => (false, false),
2400 };
2401
2402 let old_next_switch_item = if save_and_restore_next_switch_item {
2403 let item = state.next_switch_item.take();
2404 Some(state.push(item)?)
2408 } else {
2409 None
2410 };
2411
2412 let old_guest_thread = if save_and_restore_thread {
2413 self.current_thread()?
2414 } else {
2415 CurrentThread::None
2416 };
2417
2418 let waiting_set = match &reason {
2419 SuspendReason::Waiting { set, .. } => Some(*set),
2420 _ => None,
2421 };
2422
2423 let suspend_reason = &mut self.concurrent_state_mut()?.suspend_reason;
2424 assert!(suspend_reason.is_none());
2425 *suspend_reason = Some(reason);
2426
2427 if !self.fiber_async_state_mut().can_block() {
2430 return Err(format_err!("future dropped"));
2431 }
2432
2433 if let Some(set) = waiting_set {
2436 self.concurrent_state_mut()?.get_mut(set)?.num_waiting += 1;
2437 }
2438
2439 self.with_blocking(|_, cx| cx.suspend(StoreFiberYield::ReleaseStore))?;
2440
2441 if let Some(set) = waiting_set {
2442 self.concurrent_state_mut()?.get_mut(set)?.stop_waiting()?;
2443 }
2444
2445 if save_and_restore_thread {
2446 self.set_thread(old_guest_thread)?;
2447 }
2448
2449 if let Some(item) = old_next_switch_item {
2450 let state = self.concurrent_state_mut()?;
2451 state.next_switch_item = state.delete(item)?;
2452 }
2453
2454 Ok(())
2455 }
2456
2457 fn wait_for_event(
2458 &mut self,
2459 caller_instance: RuntimeInstance,
2460 waitable: Waitable,
2461 ) -> Result<()> {
2462 let caller = self.current_guest_thread()?;
2463 let state = self.concurrent_state_mut()?;
2464
2465 waitable.trap_if_in_waitable_set(state)?;
2466
2467 let set = state.get_mut(caller.thread)?.sync_call_set;
2468 waitable.join(state, Some(set))?;
2469
2470 self.switch_or_trap_if_may_not_suspend(caller_instance)?;
2471
2472 self.suspend(SuspendReason::Waiting {
2473 set,
2474 thread: caller,
2475 })?;
2476 let state = self.concurrent_state_mut()?;
2477
2478 waitable.join(state, None)
2479 }
2480
2481 fn cleanup_thread(
2503 &mut self,
2504 guest_thread: QualifiedThreadId,
2505 runtime_instance: RuntimeInstance,
2506 cleanup_task: CleanupTask,
2507 ) -> Result<()> {
2508 let state = self.concurrent_state_mut()?;
2509 state.take_next_switch_item()?;
2512 let thread_data = state.get_mut(guest_thread.thread)?;
2513 let sync_call_set = thread_data.sync_call_set;
2514 if let Some(guest_id) = thread_data.instance_rep {
2515 self.instance_state(runtime_instance)
2516 .thread_handle_table()
2517 .guest_thread_remove(guest_id)?;
2518 }
2519 let state = self.concurrent_state_mut()?;
2520
2521 for waitable in mem::take(&mut state.get_mut(sync_call_set)?.ready) {
2523 if let Some(Event::Subtask {
2524 status: Status::Returned | Status::ReturnCancelled,
2525 }) = waitable.common(self.concurrent_state_mut()?)?.event
2526 {
2527 waitable.delete_from(self)?;
2528 }
2529 }
2530
2531 let state = self.concurrent_state_mut()?;
2532 state.delete(guest_thread.thread)?;
2533 state.delete(sync_call_set)?;
2534 let task = state.get_mut(guest_thread.task)?;
2535 task.threads.remove(&guest_thread.thread);
2536
2537 if task.threads.is_empty() && !task.returned_or_cancelled() {
2538 bail!(Trap::NoAsyncResult);
2539 }
2540 let ready_to_delete = task.ready_to_delete();
2541
2542 if !task.decremented_interesting_task_count && task.exited && task.returned_or_cancelled() {
2543 task.decremented_interesting_task_count = true;
2544
2545 debug_assert!(state.interesting_tasks > 0);
2546 state.interesting_tasks -= 1;
2547 if state.interesting_tasks == 0
2548 && let Some(waker) = state.interesting_tasks_empty_waker.take()
2549 {
2550 waker.wake();
2551 }
2552 }
2553
2554 match cleanup_task {
2555 CleanupTask::Yes => {
2556 if ready_to_delete {
2557 Waitable::Guest(guest_thread.task).delete_from(self)?;
2558 }
2559 }
2560 CleanupTask::No => {}
2561 }
2562
2563 Ok(())
2564 }
2565
2566 fn cancel_guest_subtask_without_lowered_parameters(
2579 &mut self,
2580 caller_instance: RuntimeInstance,
2581 guest_task: TableId<GuestTask>,
2582 ) -> Result<()> {
2583 let concurrent_state = self.concurrent_state_mut()?;
2584 let task = concurrent_state.get_mut(guest_task)?;
2585 assert!(!task.already_lowered_parameters());
2586 task.lower_params = None;
2590 task.lift_result = None;
2591 task.exited = true;
2592 let instance = task.instance;
2593
2594 assert_eq!(1, task.threads.len());
2597 let thread = *task.threads.iter().next().unwrap();
2598 self.cleanup_thread(
2599 QualifiedThreadId {
2600 task: guest_task,
2601 thread,
2602 },
2603 caller_instance,
2604 CleanupTask::No,
2605 )?;
2606
2607 let pending = &mut self.instance_state(instance).concurrent_state().pending;
2609 let pending_count = pending.len();
2610 pending.retain(|thread, _| thread.task != guest_task);
2611 if pending.len() == pending_count {
2613 bail!(Trap::SubtaskCancelAfterTerminal);
2614 }
2615 Ok(())
2616 }
2617
2618 pub(crate) fn current_scope(&mut self) -> Result<Option<CurrentScope>> {
2621 if !self.concurrency_support() {
2622 return Ok(self
2623 .current_scope_id_not_concurrent()?
2624 .map(|id| CurrentScope::Id(Scope::Id(id))));
2625 }
2626
2627 Ok(match self.current_thread()? {
2628 CurrentThread::Guest(id) => Some(CurrentScope::Id(Scope::Id(id.task.rep()))),
2629 CurrentThread::Host(id) => Some(CurrentScope::Id(Scope::HostId(id.rep()))),
2630 CurrentThread::DeferredHost(_) => Some(CurrentScope::DeferredHost),
2631 CurrentThread::None => return Ok(None),
2632 })
2633 }
2634
2635 pub(crate) fn queue_task(
2636 &mut self,
2637 task: impl FnOnce(&mut dyn VMStore) -> Result<()> + Send + 'static,
2638 ) -> Result<()> {
2639 self.concurrent_state_mut()?
2640 .push_high_priority(WorkItem::WorkerFunction(AlwaysMut::new(Box::new(task))));
2641 Ok(())
2642 }
2643
2644 fn any_may_not_suspend(&mut self) -> Result<bool> {
2653 Ok(self
2661 .concurrent_state_mut()?
2662 .table
2663 .get_mut()
2664 .iter_mut()
2665 .filter_map(|(_, entry)| {
2666 if let Some(task) = entry.downcast_ref::<GuestTask>() {
2667 Some(task.instance)
2668 } else {
2669 None
2670 }
2671 })
2672 .collect::<Vec<_>>()
2673 .into_iter()
2674 .any(|instance| {
2675 self.instance_state(instance)
2676 .concurrent_state()
2677 .do_not_suspend
2678 }))
2679 }
2680}
2681
2682enum CleanupTask {
2683 Yes,
2684 No,
2685}
2686
2687impl Instance {
2688 fn get_event(
2691 self,
2692 store: &mut StoreOpaque,
2693 guest_task: TableId<GuestTask>,
2694 set: Option<TableId<WaitableSet>>,
2695 cancellable: bool,
2696 ) -> Result<Option<(Event, Option<(Waitable, u32)>)>> {
2697 let state = store.concurrent_state_mut()?;
2698
2699 let task = state.get_mut(guest_task)?;
2700 let event = &mut task.event;
2701 if let Some(ev) = event
2702 && (cancellable || !matches!(ev, Event::Cancelled))
2703 {
2704 log::trace!("deliver event {ev:?} to {guest_task:?}");
2705
2706 if matches!(ev, Event::Cancelled) {
2707 task.cancel_request_delivered = true;
2708 }
2709
2710 let ev = *ev;
2711 *event = None;
2712 return Ok(Some((ev, None)));
2713 }
2714
2715 let set = match set {
2716 Some(set) => set,
2717 None => return Ok(None),
2718 };
2719 let waitable = match state.get_mut(set)?.ready.pop_first() {
2720 Some(v) => v,
2721 None => return Ok(None),
2722 };
2723
2724 let common = waitable.common(state)?;
2725 let handle = match common.handle {
2726 Some(h) => h,
2727 None => bail_bug!("handle not set when delivering event"),
2728 };
2729 let event = match common.event.take() {
2730 Some(e) => e,
2731 None => bail_bug!("event not set when delivering event"),
2732 };
2733
2734 log::trace!(
2735 "deliver event {event:?} to {guest_task:?} for {waitable:?} (handle {handle}); set {set:?}"
2736 );
2737
2738 waitable.on_delivery(store, self, event)?;
2739
2740 Ok(Some((event, Some((waitable, handle)))))
2741 }
2742
2743 fn handle_callback_code(
2749 self,
2750 store: &mut StoreOpaque,
2751 guest_thread: QualifiedThreadId,
2752 runtime_instance: RuntimeComponentInstanceIndex,
2753 code: u32,
2754 ) -> Result<()> {
2755 let (code, set) = unpack_callback_code(code);
2756
2757 log::trace!("received callback code from {guest_thread:?}: {code} (set: {set})");
2758
2759 let state = store.concurrent_state_mut()?;
2760
2761 state.take_next_switch_item()?;
2762
2763 let get_set = |store: &mut StoreOpaque, handle| -> Result<_> {
2764 let set = store
2765 .instance_state(self.runtime_instance(runtime_instance))
2766 .handle_table()
2767 .waitable_set_rep(handle)?;
2768
2769 Ok(TableId::<WaitableSet>::new(set))
2770 };
2771
2772 match code {
2773 callback_code::EXIT => {
2774 log::trace!("implicit thread {guest_thread:?} completed");
2775 let task = store.concurrent_state_mut()?.get_mut(guest_thread.task)?;
2776 task.exited = true;
2777 task.callback = None;
2778
2779 let runtime_instance = self.runtime_instance(runtime_instance);
2780
2781 store.switch_or_trap_if_may_not_suspend(runtime_instance)?;
2786
2787 store.cleanup_thread(guest_thread, runtime_instance, CleanupTask::Yes)?;
2788 }
2789 callback_code::YIELD => {
2790 let old = state
2793 .get_mut(guest_thread.thread)?
2794 .wake_on_cancel
2795 .replace(WakeOnCancel::Yielding);
2796 if !old.is_none() {
2797 bail_bug!("thread unexpectedly had wake_on_cancel set");
2798 }
2799
2800 let call = GuestCall {
2807 thread: guest_thread,
2808 kind: GuestCallKind::DeliverEvent {
2809 instance: self,
2810 set: None,
2811 },
2812 };
2813 state.push_low_priority(WorkItem::GuestCall {
2816 instance: self.runtime_instance(runtime_instance),
2817 call,
2818 });
2819 }
2820 callback_code::WAIT => {
2821 let set = get_set(store, set)?;
2822 self.wait_with_callback(store, guest_thread, set)?;
2823 }
2824 _ => bail!(Trap::UnsupportedCallbackCode),
2825 }
2826
2827 Ok(())
2828 }
2829
2830 fn wait_with_callback(
2836 self,
2837 store: &mut StoreOpaque,
2838 guest_thread: QualifiedThreadId,
2839 set: TableId<WaitableSet>,
2840 ) -> Result<()> {
2841 let state = store.concurrent_state_mut()?;
2842 state.get_mut(set)?.num_waiting += 1;
2845
2846 if state.get_mut(guest_thread.task)?.event.is_some()
2847 || !state.get_mut(set)?.ready.is_empty()
2848 {
2849 let instance = state.get_mut(guest_thread.task)?.instance;
2851 state.push_high_priority(WorkItem::GuestCall {
2852 instance,
2853 call: GuestCall {
2854 thread: guest_thread,
2855 kind: GuestCallKind::DeliverEvent {
2856 instance: self,
2857 set: Some(set),
2858 },
2859 },
2860 });
2861 return Ok(());
2862 }
2863
2864 let instance = store
2870 .concurrent_state_mut()?
2871 .get_mut(guest_thread.task)?
2872 .instance;
2873 store.switch_or_trap_if_may_not_suspend(instance)?;
2874
2875 let state = store.concurrent_state_mut()?;
2881 let old = state
2882 .get_mut(guest_thread.thread)?
2883 .wake_on_cancel
2884 .replace(WakeOnCancel::Waiting(set));
2885 if !old.is_none() {
2886 bail_bug!("thread unexpectedly had wake_on_cancel set");
2887 }
2888 let old = state
2889 .get_mut(set)?
2890 .waiting
2891 .insert(guest_thread, WaitMode::Callback(self));
2892 if !old.is_none() {
2893 bail_bug!("set's waiting set already had this thread registered");
2894 }
2895 Ok(())
2896 }
2897
2898 unsafe fn stage_call<T: 'static>(
2905 self,
2906 mut store: StoreContextMut<T>,
2907 guest_thread: QualifiedThreadId,
2908 callee: SendSyncPtr<VMFuncRef>,
2909 param_count: usize,
2910 result_count: usize,
2911 async_: bool,
2912 callback: Option<SendSyncPtr<VMFuncRef>>,
2913 post_return: Option<SendSyncPtr<VMFuncRef>>,
2914 host_caller: bool,
2915 ) -> Result<()> {
2916 unsafe fn make_call<T: 'static>(
2931 store: StoreContextMut<T>,
2932 guest_thread: QualifiedThreadId,
2933 callee: SendSyncPtr<VMFuncRef>,
2934 param_count: usize,
2935 result_count: usize,
2936 ) -> impl FnOnce(&mut dyn VMStore) -> Result<[MaybeUninit<ValRaw>; MAX_FLAT_PARAMS]>
2937 + Send
2938 + Sync
2939 + 'static
2940 + use<T> {
2941 let token = StoreToken::new(store);
2942 move |store: &mut dyn VMStore| {
2943 let mut storage = [MaybeUninit::uninit(); MAX_FLAT_PARAMS];
2944
2945 store
2946 .concurrent_state_mut()?
2947 .get_mut(guest_thread.thread)?
2948 .state = GuestThreadState::Running;
2949 let task = store.concurrent_state_mut()?.get_mut(guest_thread.task)?;
2950 let lower = match task.lower_params.take() {
2951 Some(l) => l,
2952 None => bail_bug!("lower_params missing"),
2953 };
2954
2955 lower(store, &mut storage[..param_count])?;
2956
2957 let mut store = token.as_context_mut(store);
2958
2959 unsafe {
2962 crate::Func::call_unchecked_raw(
2963 &mut store,
2964 callee.as_non_null(),
2965 NonNull::new(
2966 &mut storage[..param_count.max(result_count)]
2967 as *mut [MaybeUninit<ValRaw>] as _,
2968 )
2969 .unwrap(),
2970 UncaughtException::Trap,
2971 )?;
2972 }
2973
2974 Ok(storage)
2975 }
2976 }
2977
2978 let call = unsafe {
2982 make_call(
2983 store.as_context_mut(),
2984 guest_thread,
2985 callee,
2986 param_count,
2987 result_count,
2988 )
2989 };
2990
2991 let callee_instance = store
2992 .0
2993 .concurrent_state_mut()?
2994 .get_mut(guest_thread.task)?
2995 .instance;
2996
2997 let fun = if callback.is_some() {
2998 assert!(async_);
2999
3000 Box::new(move |store: &mut dyn VMStore| {
3001 self.add_guest_thread_to_instance_table(
3002 guest_thread.thread,
3003 store,
3004 callee_instance.index,
3005 )?;
3006 let old_thread = store.set_thread(guest_thread)?;
3007 log::trace!(
3008 "stackless call: replaced {old_thread:?} with {guest_thread:?} as current thread"
3009 );
3010
3011 store.enter_instance(callee_instance);
3012
3013 let storage = call(store)?;
3020
3021 store.exit_instance(callee_instance)?;
3022
3023 store.set_thread(old_thread)?;
3024 let state = store.concurrent_state_mut()?;
3025 if let Some(t) = old_thread.guest() {
3026 state.get_mut(t.thread)?.state = GuestThreadState::Running;
3027 }
3028 log::trace!("stackless call: restored {old_thread:?} as current thread");
3029
3030 let code = unsafe { storage[0].assume_init() }.get_i32() as u32;
3033
3034 self.handle_callback_code(store, guest_thread, callee_instance.index, code)
3035 }) as Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send + Sync>
3036 } else {
3037 let token = StoreToken::new(store.as_context_mut());
3038 Box::new(move |store: &mut dyn VMStore| {
3039 self.add_guest_thread_to_instance_table(
3040 guest_thread.thread,
3041 store,
3042 callee_instance.index,
3043 )?;
3044 let old_thread = store.set_thread(guest_thread)?;
3045 log::trace!(
3046 "sync/async-stackful call: replaced {old_thread:?} with {guest_thread:?} as current thread",
3047 );
3048 let flags = self.id().get(store).instance_flags(callee_instance.index);
3049
3050 let callee_async_typed = store
3051 .concurrent_state_mut()?
3052 .get_mut(guest_thread.task)?
3053 .async_typed;
3054
3055 if !async_ && callee_async_typed {
3059 store.enter_instance(callee_instance);
3060 }
3061
3062 if !callee_async_typed {
3063 store.enter_sync_call(callee_instance)?;
3064 }
3065
3066 let storage = call(store)?;
3073
3074 if !callee_async_typed {
3075 store.exit_sync_call(callee_instance)?;
3076 }
3077
3078 if !async_ {
3079 if callee_async_typed {
3085 store.exit_instance(callee_instance)?;
3086 }
3087
3088 let lift = {
3089 let state = store.concurrent_state_mut()?;
3090 if !state.get_mut(guest_thread.task)?.result.is_none() {
3091 bail_bug!("task has already produced a result");
3092 }
3093
3094 match state.get_mut(guest_thread.task)?.lift_result.take() {
3095 Some(lift) => lift,
3096 None => bail_bug!("lift_result field is missing"),
3097 }
3098 };
3099
3100 let result = (lift.lift)(store, unsafe {
3103 mem::transmute::<&[MaybeUninit<ValRaw>], &[ValRaw]>(
3104 &storage[..result_count],
3105 )
3106 })?;
3107
3108 let post_return_arg = match result_count {
3109 0 => ValRaw::i32(0),
3110 1 => unsafe { storage[0].assume_init() },
3113 _ => unreachable!(),
3114 };
3115
3116 unsafe {
3117 call_post_return(
3118 token.as_context_mut(store),
3119 post_return.map(|v| v.as_non_null()),
3120 post_return_arg,
3121 flags,
3122 )?;
3123 }
3124
3125 self.task_complete(store, guest_thread.task, result, Status::Returned)?;
3126 }
3127
3128 store.set_thread(old_thread)?;
3129
3130 store
3131 .concurrent_state_mut()?
3132 .get_mut(guest_thread.task)?
3133 .exited = true;
3134
3135 log::trace!(
3136 "clean up thread; async lifted? {async_} async typed? {callee_async_typed}"
3137 );
3138
3139 if callee_async_typed {
3140 store.switch_or_trap_if_may_not_suspend(callee_instance)?;
3145 }
3146
3147 store.cleanup_thread(guest_thread, callee_instance, CleanupTask::Yes)?;
3149 Ok(())
3150 })
3151 };
3152
3153 store.0.concurrent_state_mut()?.push_work_item(
3154 WorkItem::GuestCall {
3155 instance: callee_instance,
3156 call: GuestCall {
3157 thread: guest_thread,
3158 kind: GuestCallKind::StartImplicit(fun),
3159 },
3160 },
3161 if host_caller {
3162 Priority::High
3163 } else {
3164 Priority::Switch
3165 },
3166 )?;
3167
3168 Ok(())
3169 }
3170
3171 unsafe fn prepare_call<T: 'static>(
3184 self,
3185 mut store: StoreContextMut<T>,
3186 start: NonNull<VMFuncRef>,
3187 return_: NonNull<VMFuncRef>,
3188 caller_instance: RuntimeComponentInstanceIndex,
3189 callee_instance: RuntimeComponentInstanceIndex,
3190 task_return_type: TypeTupleIndex,
3191 callee_async_typed: bool,
3192 memory: *mut VMMemoryDefinition,
3193 string_encoding: StringEncoding,
3194 caller_info: CallerInfo,
3195 ) -> Result<()> {
3196 enum ResultInfo {
3197 Heap { results: u32 },
3198 Stack { result_count: u32 },
3199 }
3200
3201 let result_info = match &caller_info {
3202 CallerInfo::Async {
3203 has_result: true,
3204 params,
3205 } => ResultInfo::Heap {
3206 results: match params.last() {
3207 Some(r) => r.get_u32(),
3208 None => bail_bug!("retptr missing"),
3209 },
3210 },
3211 CallerInfo::Async {
3212 has_result: false, ..
3213 } => ResultInfo::Stack { result_count: 0 },
3214 CallerInfo::Sync {
3215 result_count,
3216 params,
3217 } if *result_count > u32::try_from(MAX_FLAT_RESULTS)? => ResultInfo::Heap {
3218 results: match params.last() {
3219 Some(r) => r.get_u32(),
3220 None => bail_bug!("arg ptr missing"),
3221 },
3222 },
3223 CallerInfo::Sync { result_count, .. } => ResultInfo::Stack {
3224 result_count: *result_count,
3225 },
3226 };
3227
3228 let sync_caller = matches!(caller_info, CallerInfo::Sync { .. });
3229
3230 let start = SendSyncPtr::new(start);
3234 let return_ = SendSyncPtr::new(return_);
3235 let token = StoreToken::new(store.as_context_mut());
3236 let old_thread = store.0.current_guest_thread()?;
3237
3238 let state = store.0.concurrent_state_mut()?;
3239
3240 debug_assert_eq!(
3241 state.get_mut(old_thread.task)?.instance,
3242 self.runtime_instance(caller_instance)
3243 );
3244
3245 let guest_thread = GuestTask::new(
3246 state,
3247 Box::new(move |store, dst| {
3248 let mut store = token.as_context_mut(store);
3249 assert!(dst.len() <= MAX_FLAT_PARAMS);
3250 let mut src = [MaybeUninit::uninit(); MAX_FLAT_PARAMS + 1];
3252 let count = match caller_info {
3253 CallerInfo::Async { params, has_result } => {
3257 let params = ¶ms[..params.len() - usize::from(has_result)];
3258 for (param, src) in params.iter().zip(&mut src) {
3259 src.write(*param);
3260 }
3261 params.len()
3262 }
3263
3264 CallerInfo::Sync { params, .. } => {
3266 for (param, src) in params.iter().zip(&mut src) {
3267 src.write(*param);
3268 }
3269 params.len()
3270 }
3271 };
3272 unsafe {
3279 crate::Func::call_unchecked_raw(
3280 &mut store,
3281 start.as_non_null(),
3282 NonNull::new(
3283 &mut src[..count.max(dst.len())] as *mut [MaybeUninit<ValRaw>] as _,
3284 )
3285 .unwrap(),
3286 UncaughtException::Trap,
3287 )?;
3288 }
3289 dst.copy_from_slice(&src[..dst.len()]);
3290 let task = store.0.current_guest_thread()?.task;
3291 let state = store.0.concurrent_state_mut()?;
3292 Waitable::Guest(task).set_event(
3293 state,
3294 Some(Event::Subtask {
3295 status: Status::Started,
3296 }),
3297 )?;
3298 Ok(())
3299 }),
3300 LiftResult {
3301 lift: Box::new(move |store, src| {
3302 let mut store = token.as_context_mut(store);
3305 let mut my_src = src.to_owned(); if let ResultInfo::Heap { results } = &result_info {
3307 my_src.push(ValRaw::u32(*results));
3308 }
3309
3310 unsafe {
3317 crate::Func::call_unchecked_raw(
3318 &mut store,
3319 return_.as_non_null(),
3320 my_src.as_mut_slice().into(),
3321 UncaughtException::Trap,
3322 )?;
3323 }
3324
3325 let thread = store.0.current_guest_thread()?;
3326 let state = store.0.concurrent_state_mut()?;
3327 if sync_caller {
3328 state.get_mut(thread.task)?.sync_result = SyncResult::Produced(
3329 if let ResultInfo::Stack { result_count } = &result_info {
3330 match result_count {
3331 0 => None,
3332 1 => Some(my_src[0]),
3333 _ => unreachable!(),
3334 }
3335 } else {
3336 None
3337 },
3338 );
3339 }
3340 Ok(Box::new(DummyResult) as Box<dyn Any + Send + Sync>)
3341 }),
3342 ty: task_return_type,
3343 memory: NonNull::new(memory).map(SendSyncPtr::new),
3344 string_encoding,
3345 },
3346 Caller::Guest { thread: old_thread },
3347 None,
3348 self.runtime_instance(callee_instance),
3349 callee_async_typed,
3350 false,
3353 )?;
3354
3355 store.0.set_thread(guest_thread)?;
3358 log::trace!("pushed {guest_thread:?} as current thread; old thread was {old_thread:?}");
3359
3360 Ok(())
3361 }
3362
3363 unsafe fn call_callback<T>(
3368 self,
3369 mut store: StoreContextMut<T>,
3370 function: SendSyncPtr<VMFuncRef>,
3371 event: Event,
3372 handle: u32,
3373 ) -> Result<u32> {
3374 let (ordinal, result) = event.parts();
3375 let params = &mut [
3376 ValRaw::u32(ordinal),
3377 ValRaw::u32(handle),
3378 ValRaw::u32(result),
3379 ];
3380 unsafe {
3385 crate::Func::call_unchecked_raw(
3386 &mut store,
3387 function.as_non_null(),
3388 params.as_mut_slice().into(),
3389 UncaughtException::Trap,
3390 )?;
3391 }
3392 Ok(params[0].get_u32())
3393 }
3394
3395 unsafe fn start_call<T: 'static>(
3413 self,
3414 mut store: StoreContextMut<T>,
3415 callback: *mut VMFuncRef,
3416 post_return: *mut VMFuncRef,
3417 callee: NonNull<VMFuncRef>,
3418 param_count: u32,
3419 result_count: u32,
3420 flags: u32,
3421 storage: &mut [MaybeUninit<ValRaw>],
3422 ) -> Result<()> {
3423 let token = StoreToken::new(store.as_context_mut());
3424 let async_caller = (flags & START_FLAG_ASYNC_CALLER) != 0;
3425 let guest_thread = store.0.current_guest_thread()?;
3426 let state = store.0.concurrent_state_mut()?;
3427
3428 if !state.event_loop_running {
3429 bail_bug!("Instance::start_call called without a running event loop");
3430 }
3431
3432 let callee = SendSyncPtr::new(callee);
3433 let param_count = usize::try_from(param_count)?;
3434 assert!(param_count <= MAX_FLAT_PARAMS);
3435 let result_count = usize::try_from(result_count)?;
3436 assert!(result_count <= MAX_FLAT_RESULTS);
3437
3438 let task = state.get_mut(guest_thread.task)?;
3439 let callee_async_typed = task.async_typed;
3440 let callee_instance = task.instance;
3441
3442 task.async_lifted = (flags & START_FLAG_ASYNC_CALLEE) != 0;
3443
3444 if let Some(callback) = NonNull::new(callback) {
3445 let callback = SendSyncPtr::new(callback);
3449 task.callback = Some(Box::new(move |store, event, handle| {
3450 let store = token.as_context_mut(store);
3451 unsafe { self.call_callback::<T>(store, callback, event, handle) }
3452 }));
3453 }
3454
3455 let Caller::Guest { thread: caller } = &task.caller else {
3456 bail_bug!("start_call unexpectedly invoked for host->guest call");
3459 };
3460 let caller = *caller;
3461 let caller_instance = state.get_mut(caller.task)?.instance;
3462
3463 unsafe {
3465 self.stage_call(
3466 store.as_context_mut(),
3467 guest_thread,
3468 callee,
3469 param_count,
3470 result_count,
3471 (flags & START_FLAG_ASYNC_CALLEE) != 0,
3472 NonNull::new(callback).map(SendSyncPtr::new),
3473 NonNull::new(post_return).map(SendSyncPtr::new),
3474 false,
3475 )?;
3476 }
3477
3478 let old_do_not_suspend = if callee_async_typed {
3479 let state = store.0.instance_state(callee_instance).concurrent_state();
3486 let old_do_not_suspend = state.do_not_suspend;
3487 state.do_not_suspend = false;
3488 Some(old_do_not_suspend)
3489 } else {
3490 None
3491 };
3492
3493 let state = store.0.concurrent_state_mut()?;
3494
3495 let guest_waitable = Waitable::Guest(guest_thread.task);
3498 let old_set = guest_waitable.common(state)?.set;
3499 let set = state.get_mut(caller.thread)?.sync_call_set;
3500 guest_waitable.join(state, Some(set))?;
3501
3502 store.0.set_thread(CurrentThread::None)?;
3503
3504 let mut yielded = false;
3520 let (status, waitable) = loop {
3521 store.0.suspend(if yielded {
3522 SuspendReason::Waiting {
3523 set,
3524 thread: caller,
3525 }
3526 } else {
3527 yielded = true;
3528 SuspendReason::YieldingToSubtask { thread: caller }
3529 })?;
3530
3531 if let Some(old_do_not_suspend) = old_do_not_suspend {
3532 store
3533 .0
3534 .instance_state(callee_instance)
3535 .concurrent_state()
3536 .do_not_suspend = old_do_not_suspend;
3537 }
3538
3539 let state = store.0.concurrent_state_mut()?;
3540
3541 log::trace!("taking event for {:?}", guest_thread.task);
3542 let event = guest_waitable.take_event(state)?;
3543 let Some(Event::Subtask { status }) = event else {
3544 bail_bug!("subtasks should only get subtask events, got {event:?}")
3545 };
3546
3547 log::trace!("status {status:?} for {:?}", guest_thread.task);
3548
3549 if status == Status::Returned {
3550 break (status, None);
3552 } else if async_caller {
3553 let handle = store
3557 .0
3558 .instance_state(caller_instance)
3559 .handle_table()
3560 .subtask_insert_guest(guest_thread.task.rep())?;
3561 store
3562 .0
3563 .concurrent_state_mut()?
3564 .get_mut(guest_thread.task)?
3565 .common
3566 .handle = Some(handle);
3567 break (status, Some(handle));
3568 } else {
3569 store.0.switch_or_trap_if_may_not_suspend(caller_instance)?;
3573 }
3574 };
3575
3576 guest_waitable.join(store.0.concurrent_state_mut()?, old_set)?;
3577
3578 store.0.set_thread(caller)?;
3580 store
3581 .0
3582 .concurrent_state_mut()?
3583 .get_mut(caller.thread)?
3584 .state = GuestThreadState::Running;
3585 log::trace!("popped current thread {guest_thread:?}; new thread is {caller:?}");
3586
3587 if async_caller {
3588 let Some(slot) = storage.first_mut() else {
3589 bail_bug!("no storage for async call status");
3590 };
3591 *slot = MaybeUninit::new(ValRaw::u32(status.pack(waitable)));
3592 } else {
3593 let state = store.0.concurrent_state_mut()?;
3596 let task = state.get_mut(guest_thread.task)?;
3597 if let Some(result) = task.sync_result.take()? {
3598 if let Some(result) = result {
3599 storage[0] = MaybeUninit::new(result);
3600 }
3601
3602 if task.exited && task.ready_to_delete() {
3603 Waitable::Guest(guest_thread.task).delete_from(store.0)?;
3604 }
3605 }
3606 }
3607
3608 Ok(())
3609 }
3610
3611 pub(crate) fn first_poll<T: 'static, R: Send + 'static>(
3627 self,
3628 mut store: StoreContextMut<'_, T>,
3629 host_task: EnteredHostTask,
3630 result_may_require_realloc: bool,
3631 future: impl Future<Output = Result<R>> + Send + 'static,
3632 lower: impl FnOnce(StoreContextMut<T>, Option<R>, bool, Option<TableId<HostTask>>) -> Result<()>
3633 + Send
3634 + 'static,
3635 ) -> Result<u32> {
3636 let token = StoreToken::new(store.as_context_mut());
3637
3638 let (join_handle, future) = JoinHandle::run(future);
3641 let mut future = Box::pin(future);
3642
3643 let poll = tls::set(store.0, || {
3648 future
3649 .as_mut()
3650 .poll(&mut Context::from_waker(&Waker::noop()))
3651 });
3652
3653 match poll {
3654 Poll::Ready(result) => {
3656 let result = result.transpose()?;
3657 let task = store.0.current_materialized_host_task()?;
3660 lower(store.as_context_mut(), result, true, task)?;
3661 return Ok(Status::Returned.pack(None));
3662 }
3663
3664 Poll::Pending => {}
3666 }
3667
3668 let Some(task) = store.0.materialize_host_task_id()? else {
3672 bail_bug!("current thread is not a host thread")
3673 };
3674 {
3675 let state = &mut store.0.concurrent_state_mut()?.get_mut(task)?.state;
3676 assert!(matches!(state, HostTaskState::CalleeStarted));
3677 *state = HostTaskState::CalleeRunning(join_handle);
3678 }
3679
3680 let future = Box::pin(async move {
3688 let result = match run_with_host_task_set(task, future).await? {
3689 Some(result) => Some(result?),
3690 None => None,
3691 };
3692 let on_complete = move |store: &mut dyn VMStore| {
3693 let mut store = token.as_context_mut(store);
3697 let old = store.0.set_thread(task)?;
3698
3699 let status = if result.is_some() {
3700 Status::Returned
3701 } else {
3702 Status::ReturnCancelled
3703 };
3704
3705 lower(store.as_context_mut(), result, false, Some(task))?;
3706 let state = store.0.concurrent_state_mut()?;
3707 match &mut state.get_mut(task)?.state {
3708 pending @ HostTaskState::CalleeCancelling => {
3711 *pending = HostTaskState::CalleeDone { cancelled: true };
3712 }
3713
3714 other => *other = HostTaskState::CalleeDone { cancelled: false },
3716 }
3717 Waitable::Host(task).set_event(state, Some(Event::Subtask { status }))?;
3718
3719 store.0.set_thread(old)?;
3720 Ok(())
3721 };
3722
3723 tls::get(move |store| {
3724 if result_may_require_realloc {
3725 store
3730 .concurrent_state_mut()?
3731 .push_high_priority(WorkItem::WorkerFunction(AlwaysMut::new(Box::new(
3732 on_complete,
3733 ))));
3734 Ok(())
3735 } else {
3736 on_complete(store)
3739 }
3740 })
3741 });
3742
3743 let caller = match host_task {
3746 Some(caller) => caller,
3747 None => bail_bug!("host task wasn't created but should have been"),
3748 };
3749 let state = store.0.concurrent_state_mut()?;
3750 state.push_future(future);
3751 let instance = state.get_mut(caller.task)?.instance;
3752 let handle = store
3753 .0
3754 .instance_state(instance)
3755 .handle_table()
3756 .subtask_insert_host(task.rep())?;
3757 store.0.concurrent_state_mut()?.get_mut(task)?.common.handle = Some(handle);
3758 log::trace!("assign {task:?} handle {handle} for {caller:?} instance {instance:?}");
3759
3760 store.0.set_thread(caller)?;
3764 Ok(Status::Started.pack(Some(handle)))
3765 }
3766
3767 pub(crate) fn task_return(
3770 self,
3771 store: &mut dyn VMStore,
3772 ty: TypeTupleIndex,
3773 options: OptionsIndex,
3774 storage: &[ValRaw],
3775 ) -> Result<()> {
3776 let guest_thread = store.current_guest_thread()?;
3777 let state = store.concurrent_state_mut()?;
3778 if !state.get_mut(guest_thread.task)?.async_lifted {
3779 bail!(Trap::TaskReturnOrCancelSyncLifted);
3780 }
3781 let lift = state
3782 .get_mut(guest_thread.task)?
3783 .lift_result
3784 .take()
3785 .ok_or_else(|| Trap::TaskCancelOrReturnTwice)?;
3786 if !state.get_mut(guest_thread.task)?.result.is_none() {
3787 bail_bug!("task result unexpectedly already set");
3788 }
3789
3790 let CanonicalOptions {
3791 string_encoding,
3792 data_model,
3793 ..
3794 } = &self.id().get(store).component().env_component().options[options];
3795
3796 let invalid = ty != lift.ty
3797 || string_encoding != &lift.string_encoding
3798 || match data_model {
3799 CanonicalOptionsDataModel::LinearMemory(opts) => match opts.memory {
3800 Some(memory) => {
3801 let expected = lift.memory.map(|v| v.as_ptr()).unwrap_or(ptr::null_mut());
3802 let actual = self.id().get(store).runtime_memory(memory);
3803 expected != actual.as_ptr()
3804 }
3805 None => false,
3808 },
3809 CanonicalOptionsDataModel::Gc { .. } => true,
3811 };
3812
3813 if invalid {
3814 bail!(Trap::TaskReturnInvalid);
3815 }
3816
3817 log::trace!("task.return for {guest_thread:?}");
3818
3819 let result = (lift.lift)(store, storage)?;
3820 self.task_complete(store, guest_thread.task, result, Status::Returned)
3821 }
3822
3823 pub(crate) fn task_cancel(self, store: &mut StoreOpaque) -> Result<()> {
3825 let guest_thread = store.current_guest_thread()?;
3826 let state = store.concurrent_state_mut()?;
3827 let task = state.get_mut(guest_thread.task)?;
3828 if !task.async_lifted {
3829 bail!(Trap::TaskReturnOrCancelSyncLifted);
3830 }
3831 if !task.cancel_request_delivered {
3832 bail!(Trap::TaskCancelNotCancelled);
3833 }
3834 _ = task
3835 .lift_result
3836 .take()
3837 .ok_or_else(|| Trap::TaskCancelOrReturnTwice)?;
3838
3839 if !task.result.is_none() {
3840 bail_bug!("task result should not bet set yet");
3841 }
3842
3843 log::trace!("task.cancel for {guest_thread:?}");
3844
3845 self.task_complete(
3846 store,
3847 guest_thread.task,
3848 Box::new(DummyResult),
3849 Status::ReturnCancelled,
3850 )
3851 }
3852
3853 fn task_complete(
3859 self,
3860 store: &mut StoreOpaque,
3861 guest_task: TableId<GuestTask>,
3862 result: Box<dyn Any + Send + Sync>,
3863 status: Status,
3864 ) -> Result<()> {
3865 store
3866 .component_resource_tables(Some(self))?
3867 .validate_scope_exit()?;
3868
3869 let state = store.concurrent_state_mut()?;
3870 let task = state.get_mut(guest_task)?;
3871
3872 task.event = None;
3876
3877 if let Caller::Host { tx, .. } = &mut task.caller {
3878 if let Some(tx) = tx.take() {
3879 _ = tx.send(result);
3880 }
3881 } else {
3882 task.result = Some(result);
3883 Waitable::Guest(guest_task).set_event(state, Some(Event::Subtask { status }))?;
3884 }
3885
3886 Ok(())
3887 }
3888
3889 pub(crate) fn waitable_set_new(
3891 self,
3892 store: &mut StoreOpaque,
3893 caller_instance: RuntimeComponentInstanceIndex,
3894 ) -> Result<u32> {
3895 let set = store.concurrent_state_mut()?.push(WaitableSet::default())?;
3896 let handle = store
3897 .instance_state(self.runtime_instance(caller_instance))
3898 .handle_table()
3899 .waitable_set_insert(set.rep())?;
3900 log::trace!("new waitable set {set:?} (handle {handle})");
3901 Ok(handle)
3902 }
3903
3904 pub(crate) fn waitable_set_drop(
3906 self,
3907 store: &mut StoreOpaque,
3908 caller_instance: RuntimeComponentInstanceIndex,
3909 set: u32,
3910 ) -> Result<()> {
3911 let rep = store
3912 .instance_state(self.runtime_instance(caller_instance))
3913 .handle_table()
3914 .waitable_set_remove(set)?;
3915
3916 log::trace!("drop waitable set {rep} (handle {set})");
3917
3918 let set = store
3922 .concurrent_state_mut()?
3923 .get_mut(TableId::<WaitableSet>::new(rep))?;
3924 if set.num_waiting > 0 {
3925 bail!(Trap::WaitableSetDropHasWaiters);
3926 }
3927
3928 store
3929 .concurrent_state_mut()?
3930 .delete(TableId::<WaitableSet>::new(rep))?;
3931
3932 Ok(())
3933 }
3934
3935 pub(crate) fn waitable_join(
3937 self,
3938 store: &mut StoreOpaque,
3939 caller_instance: RuntimeComponentInstanceIndex,
3940 waitable_handle: u32,
3941 set_handle: u32,
3942 ) -> Result<()> {
3943 let mut instance = self.id().get_mut(store);
3944 let waitable =
3945 Waitable::from_instance(instance.as_mut(), caller_instance, waitable_handle)?;
3946
3947 let set = if set_handle == 0 {
3948 None
3949 } else {
3950 let set = instance.instance_states().0[caller_instance]
3951 .handle_table()
3952 .waitable_set_rep(set_handle)?;
3953
3954 let state = store.concurrent_state_mut()?;
3955 if let Some(old) = waitable.common(state)?.set
3956 && state.get_mut(old)?.is_sync_call_set
3957 {
3958 bail!(Trap::WaitableSyncAndAsync);
3959 }
3960
3961 Some(TableId::<WaitableSet>::new(set))
3962 };
3963
3964 log::trace!(
3965 "waitable {waitable:?} (handle {waitable_handle}) join set {set:?} (handle {set_handle})",
3966 );
3967
3968 waitable.join(store.concurrent_state_mut()?, set)
3969 }
3970
3971 pub(crate) fn subtask_drop(
3973 self,
3974 store: &mut StoreOpaque,
3975 caller_instance: RuntimeComponentInstanceIndex,
3976 task_id: u32,
3977 ) -> Result<()> {
3978 self.waitable_join(store, caller_instance, task_id, 0)?;
3979
3980 let (rep, is_host) = store
3981 .instance_state(self.runtime_instance(caller_instance))
3982 .handle_table()
3983 .subtask_remove(task_id)?;
3984
3985 let concurrent_state = store.concurrent_state_mut()?;
3986 let (waitable, delete) = if is_host {
3987 let id = TableId::<HostTask>::new(rep);
3988 let task = concurrent_state.get_mut(id)?;
3989 match &task.state {
3990 HostTaskState::CalleeRunning(_) | HostTaskState::CalleeCancelling => {
3991 bail!(Trap::SubtaskDropNotResolved)
3992 }
3993 HostTaskState::CalleeDone { .. } => {}
3994 HostTaskState::CalleeStarted | HostTaskState::CalleeFinished(_) => {
3995 bail_bug!("invalid state for callee in `subtask.drop`")
3996 }
3997 }
3998
3999 (Waitable::Host(id), true)
4000 } else {
4001 let id = TableId::<GuestTask>::new(rep);
4002 let task = concurrent_state.get_mut(id)?;
4003 if task.lift_result.is_some() {
4004 bail!(Trap::SubtaskDropNotResolved);
4005 }
4006 (
4007 Waitable::Guest(id),
4008 concurrent_state.get_mut(id)?.ready_to_delete(),
4009 )
4010 };
4011
4012 waitable.common(concurrent_state)?.handle = None;
4013
4014 if waitable.take_event(concurrent_state)?.is_some() {
4017 bail!(Trap::SubtaskDropNotResolved);
4018 }
4019
4020 if delete {
4021 waitable.delete_from(store)?;
4022 }
4023
4024 log::trace!("subtask_drop {waitable:?} (handle {task_id})");
4025 Ok(())
4026 }
4027
4028 pub(crate) fn waitable_set_wait(
4030 self,
4031 store: &mut StoreOpaque,
4032 options: OptionsIndex,
4033 set: u32,
4034 payload: u32,
4035 ) -> Result<u32> {
4036 let &CanonicalOptions {
4037 instance: caller_instance,
4038 ..
4039 } = &self.id().get(store).component().env_component().options[options];
4040 let caller = self.runtime_instance(caller_instance);
4041 let rep = store
4042 .instance_state(self.runtime_instance(caller_instance))
4043 .handle_table()
4044 .waitable_set_rep(set)?;
4045
4046 self.waitable_check(
4047 store,
4048 caller,
4049 WaitableCheck::Wait,
4050 WaitableCheckParams {
4051 set: TableId::new(rep),
4052 options,
4053 payload,
4054 },
4055 )
4056 }
4057
4058 pub(crate) fn waitable_set_poll(
4060 self,
4061 store: &mut StoreOpaque,
4062 options: OptionsIndex,
4063 set: u32,
4064 payload: u32,
4065 ) -> Result<u32> {
4066 let &CanonicalOptions {
4067 instance: caller_instance,
4068 ..
4069 } = &self.id().get(store).component().env_component().options[options];
4070 let caller = self.runtime_instance(caller_instance);
4071 let rep = store
4072 .instance_state(caller)
4073 .handle_table()
4074 .waitable_set_rep(set)?;
4075
4076 self.waitable_check(
4077 store,
4078 caller,
4079 WaitableCheck::Poll,
4080 WaitableCheckParams {
4081 set: TableId::new(rep),
4082 options,
4083 payload,
4084 },
4085 )
4086 }
4087
4088 pub(crate) fn thread_index(&self, store: &mut dyn VMStore) -> Result<u32> {
4090 let thread_id = store.current_guest_thread()?.thread;
4091 match store
4092 .concurrent_state_mut()?
4093 .get_mut(thread_id)?
4094 .instance_rep
4095 {
4096 Some(r) => Ok(r),
4097 None => bail_bug!("thread should have instance_rep by now"),
4098 }
4099 }
4100
4101 pub(crate) fn thread_new_indirect<T: 'static>(
4103 self,
4104 mut store: StoreContextMut<T>,
4105 runtime_instance: RuntimeComponentInstanceIndex,
4106 start_func_ty: ModuleInternedTypeIndex,
4107 start_func_table_idx: RuntimeTableIndex,
4108 start_func_idx: u32,
4109 context: i32,
4110 ) -> Result<u32> {
4111 log::trace!("creating new thread");
4112
4113 let (instance, registry) = self.id().get_mut_and_registry(store.0);
4114 let Some(start_func_ty) = instance.component().signatures().shared_type(start_func_ty)
4115 else {
4116 bail_bug!("thread.new-indirect start function type should be registered");
4117 };
4118 let callee = instance
4119 .index_runtime_func_table(registry, start_func_table_idx, start_func_idx as u64)?
4120 .ok_or_else(|| Trap::ThreadNewIndirectUninitialized)?;
4121 if callee.type_index(store.0) != start_func_ty {
4122 bail!(Trap::ThreadNewIndirectInvalidType);
4123 }
4124
4125 let token = StoreToken::new(store.as_context_mut());
4126 let start_func = Box::new(
4127 move |store: &mut dyn VMStore, guest_thread: QualifiedThreadId| -> Result<()> {
4128 let old_thread = store.set_thread(guest_thread)?;
4129 log::trace!(
4130 "thread start: replaced {old_thread:?} with {guest_thread:?} as current thread"
4131 );
4132
4133 let mut store = token.as_context_mut(store);
4134 let mut params = [ValRaw::i32(context)];
4135 unsafe { callee.call_unchecked(store.as_context_mut(), &mut params)? };
4138
4139 store.0.set_thread(old_thread)?;
4140
4141 let runtime_instance = self.runtime_instance(runtime_instance);
4142
4143 store
4146 .0
4147 .switch_or_trap_if_may_not_suspend(runtime_instance)?;
4148
4149 store
4150 .0
4151 .cleanup_thread(guest_thread, runtime_instance, CleanupTask::Yes)?;
4152
4153 log::trace!("explicit thread {guest_thread:?} completed");
4154 let state = store.0.concurrent_state_mut()?;
4155 if let Some(t) = old_thread.guest() {
4156 state.get_mut(t.thread)?.state = GuestThreadState::Running;
4157 }
4158 log::trace!("thread start: restored {old_thread:?} as current thread");
4159
4160 Ok(())
4161 },
4162 );
4163
4164 let current_thread = store.0.current_guest_thread()?;
4165 let state = store.0.concurrent_state_mut()?;
4166 let parent_task = current_thread.task;
4167
4168 let new_thread = GuestThread::new_explicit(state, parent_task, start_func)?;
4169 let thread_id = state.push(new_thread)?;
4170 state.get_mut(parent_task)?.threads.insert(thread_id);
4171
4172 log::trace!("new thread with id {thread_id:?} created");
4173
4174 self.add_guest_thread_to_instance_table(thread_id, store.0, runtime_instance)
4175 }
4176
4177 pub(crate) fn resume_thread(
4178 self,
4179 store: &mut StoreOpaque,
4180 runtime_instance: RuntimeComponentInstanceIndex,
4181 thread_idx: u32,
4182 how: ResumeThread,
4183 ) -> Result<bool> {
4184 let thread_id =
4185 GuestThread::from_instance(self.id().get_mut(store), runtime_instance, thread_idx)?;
4186 let state = store.concurrent_state_mut()?;
4187 let guest_thread = QualifiedThreadId::qualify(state, thread_id)?;
4188
4189 if store.current_guest_thread()? == guest_thread {
4190 bail!(Trap::CannotResumeThread);
4191 }
4192
4193 let priority = match how {
4194 ResumeThread::Resume => Priority::Switch,
4195 ResumeThread::ResumeLater => Priority::Low,
4196
4197 ResumeThread::Promote => {
4200 let instance = self.runtime_instance(runtime_instance);
4201 let do_not_enter = store
4202 .instance_state(instance)
4203 .concurrent_state()
4204 .do_not_enter;
4205 return store.concurrent_state_mut()?.promote_work_item_matching(
4206 |item: &WorkItem| match item {
4207 WorkItem::ResumeThread { thread, .. }
4208 | WorkItem::ResumeFiber { thread, .. } => *thread == guest_thread,
4209
4210 WorkItem::GuestCall {
4211 call: GuestCall { thread, kind },
4212 ..
4213 } => {
4214 *thread == guest_thread
4215 && match kind {
4216 GuestCallKind::DeliverEvent { .. } => !do_not_enter,
4220
4221 GuestCallKind::StartExplicit(_) => true,
4224
4225 GuestCallKind::StartImplicit(_) => false,
4229 }
4230 }
4231
4232 WorkItem::PushFuture(_) | WorkItem::WorkerFunction(_) => false,
4233 },
4234 );
4235 }
4236 };
4237
4238 let state = store.concurrent_state_mut()?;
4239 let thread = state.get_mut(guest_thread.thread)?;
4240
4241 match &thread.state {
4244 GuestThreadState::NotStartedExplicit(_) | GuestThreadState::Suspended(_) => {}
4245 _ => bail!(Trap::CannotResumeThread),
4246 }
4247
4248 match mem::replace(&mut thread.state, GuestThreadState::Running) {
4249 GuestThreadState::NotStartedExplicit(start_func) => {
4250 log::trace!("starting thread {guest_thread:?}");
4251 let guest_call = WorkItem::GuestCall {
4252 instance: self.runtime_instance(runtime_instance),
4253 call: GuestCall {
4254 thread: guest_thread,
4255 kind: GuestCallKind::StartExplicit(Box::new(move |store| {
4256 start_func(store, guest_thread)
4257 })),
4258 },
4259 };
4260 store
4261 .concurrent_state_mut()?
4262 .push_work_item(guest_call, priority)?;
4263 }
4264 GuestThreadState::Suspended(fiber) => {
4265 log::trace!("resuming thread {thread_id:?} that was suspended");
4266 store.concurrent_state_mut()?.push_work_item(
4267 WorkItem::ResumeFiber {
4268 instance: self.runtime_instance(runtime_instance),
4269 thread: guest_thread,
4270 fiber,
4271 },
4272 priority,
4273 )?;
4274 }
4275 other => {
4276 thread.state = other;
4277 bail_bug!("thread state checked to be resumable above");
4278 }
4279 }
4280 Ok(true)
4281 }
4282
4283 fn add_guest_thread_to_instance_table(
4284 self,
4285 thread_id: TableId<GuestThread>,
4286 store: &mut StoreOpaque,
4287 runtime_instance: RuntimeComponentInstanceIndex,
4288 ) -> Result<u32> {
4289 let guest_id = store
4290 .instance_state(self.runtime_instance(runtime_instance))
4291 .thread_handle_table()
4292 .guest_thread_insert(thread_id.rep())?;
4293 store
4294 .concurrent_state_mut()?
4295 .get_mut(thread_id)?
4296 .instance_rep = Some(guest_id);
4297 Ok(guest_id)
4298 }
4299
4300 pub(crate) fn suspension_intrinsic(
4304 self,
4305 store: &mut StoreOpaque,
4306 caller: RuntimeComponentInstanceIndex,
4307 yielding: bool,
4308 to_thread: SuspensionTarget,
4309 ) -> Result<WaitResult> {
4310 let check_suspend = match to_thread {
4311 SuspensionTarget::Promote(thread) => {
4312 !self.resume_thread(store, caller, thread, ResumeThread::Promote)?
4313 }
4314 SuspensionTarget::Resume(thread) => {
4315 if !self.resume_thread(store, caller, thread, ResumeThread::Resume)? {
4316 bail_bug!(
4317 "`resume_thread` should only ever return false \
4318 when `ResumeThread::Promote` is passed to it"
4319 );
4320 }
4321 false
4322 }
4323 SuspensionTarget::None => true,
4324 };
4325
4326 if check_suspend && !store.switch_if_may_not_suspend(self.runtime_instance(caller))? {
4327 return if yielding {
4328 Ok(WaitResult::Completed)
4329 } else {
4330 Err(Trap::CannotBlockSyncTask.into())
4331 };
4332 }
4333
4334 let guest_thread = store.current_guest_thread()?;
4335
4336 let reason = if yielding {
4337 SuspendReason::Yielding {
4338 thread: guest_thread,
4339 }
4340 } else {
4341 SuspendReason::ExplicitlySuspending {
4342 thread: guest_thread,
4343 }
4344 };
4345
4346 store.suspend(reason)?;
4347
4348 Ok(WaitResult::Completed)
4349 }
4350
4351 fn waitable_check(
4353 self,
4354 store: &mut StoreOpaque,
4355 caller: RuntimeInstance,
4356 check: WaitableCheck,
4357 params: WaitableCheckParams,
4358 ) -> Result<u32> {
4359 let guest_thread = store.current_guest_thread()?;
4360
4361 log::trace!("waitable check for {guest_thread:?}; set {:?}", params.set);
4362
4363 match &check {
4366 WaitableCheck::Wait => {
4367 let set = params.set;
4368
4369 loop {
4374 let state = store.concurrent_state_mut()?;
4375 let task = state.get_mut(guest_thread.task)?;
4376 if !(task.event.is_none() || matches!(task.event, Some(Event::Cancelled)))
4377 || !state.get_mut(set)?.ready.is_empty()
4378 {
4379 break;
4380 }
4381
4382 store.switch_or_trap_if_may_not_suspend(caller)?;
4383
4384 store.suspend(SuspendReason::Waiting {
4385 set,
4386 thread: guest_thread,
4387 })?;
4388 }
4389 }
4390 WaitableCheck::Poll => {}
4391 }
4392
4393 log::trace!(
4394 "waitable check for {guest_thread:?}; set {:?}, part two",
4395 params.set
4396 );
4397
4398 let event = self.get_event(store, guest_thread.task, Some(params.set), false)?;
4400
4401 let (ordinal, handle, result) = match &check {
4402 WaitableCheck::Wait => {
4403 let (event, waitable) = match event {
4404 Some(p) => p,
4405 None => bail_bug!("event expected to be present"),
4406 };
4407 let handle = waitable.map(|(_, v)| v).unwrap_or(0);
4408 let (ordinal, result) = event.parts();
4409 (ordinal, handle, result)
4410 }
4411 WaitableCheck::Poll => {
4412 if let Some((event, waitable)) = event {
4413 let handle = waitable.map(|(_, v)| v).unwrap_or(0);
4414 let (ordinal, result) = event.parts();
4415 (ordinal, handle, result)
4416 } else {
4417 log::trace!(
4418 "no events ready to deliver via waitable-set.poll to {:?}; set {:?}",
4419 guest_thread.task,
4420 params.set
4421 );
4422 let (ordinal, result) = Event::None.parts();
4423 (ordinal, 0, result)
4424 }
4425 }
4426 };
4427 let memory = self.options_memory_mut(store, params.options);
4428 let ptr = crate::component::func::validate_inbounds_dynamic(
4429 &CanonicalAbiInfo::POINTER_PAIR,
4430 memory,
4431 &ValRaw::u32(params.payload),
4432 )?;
4433 memory[ptr + 0..][..4].copy_from_slice(&handle.to_le_bytes());
4434 memory[ptr + 4..][..4].copy_from_slice(&result.to_le_bytes());
4435 Ok(ordinal)
4436 }
4437
4438 pub(crate) fn subtask_cancel(
4440 self,
4441 store: &mut StoreOpaque,
4442 caller_instance: RuntimeComponentInstanceIndex,
4443 async_: bool,
4444 task_id: u32,
4445 ) -> Result<u32> {
4446 let (rep, is_host) = store
4447 .instance_state(self.runtime_instance(caller_instance))
4448 .handle_table()
4449 .subtask_rep(task_id)?;
4450 let waitable = if is_host {
4451 Waitable::Host(TableId::<HostTask>::new(rep))
4452 } else {
4453 Waitable::Guest(TableId::<GuestTask>::new(rep))
4454 };
4455 let concurrent_state = store.concurrent_state_mut()?;
4456
4457 log::trace!("subtask_cancel {waitable:?} (handle {task_id}; async {async_})");
4458
4459 waitable.trap_if_in_waitable_set(concurrent_state)?;
4460
4461 let needs_block;
4462 if let Waitable::Host(host_task) = waitable {
4463 let state = &mut concurrent_state.get_mut(host_task)?.state;
4464 match state {
4465 HostTaskState::CalleeRunning(handle) => {
4472 handle.abort();
4473 *state = HostTaskState::CalleeCancelling;
4474 needs_block = true;
4475 }
4476
4477 HostTaskState::CalleeCancelling | HostTaskState::CalleeDone { cancelled: true } => {
4480 bail!(Trap::SubtaskCancelAfterTerminal);
4481 }
4482 HostTaskState::CalleeDone { cancelled: false } => {
4483 *state = HostTaskState::CalleeDone { cancelled: true };
4486 needs_block = false;
4487 }
4488
4489 HostTaskState::CalleeStarted | HostTaskState::CalleeFinished(_) => {
4492 bail_bug!("invalid states for host callee")
4493 }
4494 }
4495 } else {
4496 let guest_task = TableId::<GuestTask>::new(rep);
4497 let task = concurrent_state.get_mut(guest_task)?;
4498 if !task.already_lowered_parameters() {
4499 store.cancel_guest_subtask_without_lowered_parameters(
4500 self.runtime_instance(caller_instance),
4501 guest_task,
4502 )?;
4503 return Ok(Status::StartCancelled as u32);
4504 } else if !task.returned_or_cancelled() {
4505 task.event = Some(Event::Cancelled);
4513 let runtime_instance = task.instance;
4514 for thread in task.threads.clone() {
4515 let thread = QualifiedThreadId {
4516 task: guest_task,
4517 thread,
4518 };
4519 let concurrent_state = store.concurrent_state_mut()?;
4520 let thread_mut = concurrent_state.get_mut(thread.thread)?;
4521
4522 let yield_ = |store: &mut StoreOpaque| {
4523 let state = store.instance_state(runtime_instance).concurrent_state();
4528 let old_do_not_suspend = state.do_not_suspend;
4529 state.do_not_suspend = false;
4530
4531 let caller = store.current_guest_thread()?;
4532
4533 let state = store.concurrent_state_mut()?;
4538 let set = state.get_mut(caller.thread)?.sync_call_set;
4539 waitable.join(state, Some(set))?;
4540
4541 store.suspend(SuspendReason::YieldingToSubtask { thread: caller })?;
4542
4543 let state = store.concurrent_state_mut()?;
4544 waitable.join(state, None)?;
4545
4546 store
4547 .instance_state(runtime_instance)
4548 .concurrent_state()
4549 .do_not_suspend = old_do_not_suspend;
4550
4551 Ok::<(), crate::Error>(())
4552 };
4553
4554 match thread_mut.wake_on_cancel.take() {
4555 WakeOnCancel::Waiting(set) => {
4556 let item = match concurrent_state.get_mut(set)?.waiting.remove(&thread)
4558 {
4559 Some(WaitMode::Callback(instance)) => WorkItem::GuestCall {
4560 instance: runtime_instance,
4561 call: GuestCall {
4562 thread,
4563 kind: GuestCallKind::DeliverEvent {
4564 instance,
4565 set: Some(set),
4566 },
4567 },
4568 },
4569 other => bail_bug!(
4570 "expected `Some(WaitMode::Callback(_))`; got `{other:?}`"
4571 ),
4572 };
4573 concurrent_state.set_switch_item(item)?;
4574
4575 yield_(store)?;
4576
4577 break;
4578 }
4579 WakeOnCancel::Yielding => {
4580 if concurrent_state.promote_thread_work_item(thread)? {
4581 yield_(store)?;
4582 break;
4583 } else if store
4584 .instance_state(runtime_instance)
4585 .concurrent_state()
4586 .pending
4587 .contains_key(&thread)
4588 {
4589 store
4599 .concurrent_state_mut()?
4600 .get_mut(thread.thread)?
4601 .wake_on_cancel = WakeOnCancel::Yielding;
4602 } else {
4603 bail_bug!("thread with `WakeOnCancel::Yielding` not promotable");
4604 }
4605 }
4606 WakeOnCancel::None => {}
4607 }
4608 }
4609
4610 needs_block = !store
4613 .concurrent_state_mut()?
4614 .get_mut(guest_task)?
4615 .returned_or_cancelled()
4616 } else {
4617 needs_block = false;
4618 }
4619 };
4620
4621 if needs_block {
4625 if async_ {
4626 return Ok(BLOCKED);
4627 }
4628
4629 let old_next_switch_item = {
4632 let state = store.concurrent_state_mut()?;
4633 let item = state.next_switch_item.take();
4634 state.push(item)?
4638 };
4639
4640 store.wait_for_event(self.runtime_instance(caller_instance), waitable)?;
4643
4644 let state = store.concurrent_state_mut()?;
4645 state.next_switch_item = state.delete(old_next_switch_item)?;
4646
4647 }
4649
4650 let event = waitable.take_event(store.concurrent_state_mut()?)?;
4651 if let Some(Event::Subtask {
4652 status: status @ (Status::Returned | Status::ReturnCancelled),
4653 }) = event
4654 {
4655 Ok(status as u32)
4656 } else {
4657 bail!(Trap::SubtaskCancelAfterTerminal);
4658 }
4659 }
4660}
4661
4662pub trait VMComponentAsyncStore {
4670 unsafe fn prepare_call(
4676 &mut self,
4677 instance: Instance,
4678 memory: *mut VMMemoryDefinition,
4679 start: NonNull<VMFuncRef>,
4680 return_: NonNull<VMFuncRef>,
4681 caller_instance: RuntimeComponentInstanceIndex,
4682 callee_instance: RuntimeComponentInstanceIndex,
4683 task_return_type: TypeTupleIndex,
4684 callee_async: bool,
4685 string_encoding: StringEncoding,
4686 result_count: u32,
4687 storage: *mut ValRaw,
4688 storage_len: usize,
4689 ) -> Result<()>;
4690
4691 unsafe fn start_call(
4694 &mut self,
4695 instance: Instance,
4696 callback: *mut VMFuncRef,
4697 post_return: *mut VMFuncRef,
4698 callee: NonNull<VMFuncRef>,
4699 param_count: u32,
4700 result_count: u32,
4701 flags: u32,
4702 storage: *mut MaybeUninit<ValRaw>,
4703 storage_len: usize,
4704 ) -> Result<()>;
4705
4706 fn future_write(
4708 &mut self,
4709 instance: Instance,
4710 caller: RuntimeComponentInstanceIndex,
4711 ty: TypeFutureTableIndex,
4712 options: OptionsIndex,
4713 future: u32,
4714 address: u32,
4715 ) -> Result<u32>;
4716
4717 fn future_read(
4719 &mut self,
4720 instance: Instance,
4721 caller: RuntimeComponentInstanceIndex,
4722 ty: TypeFutureTableIndex,
4723 options: OptionsIndex,
4724 future: u32,
4725 address: u32,
4726 ) -> Result<u32>;
4727
4728 fn future_drop_writable(
4730 &mut self,
4731 instance: Instance,
4732 ty: TypeFutureTableIndex,
4733 writer: u32,
4734 ) -> Result<()>;
4735
4736 fn stream_write(
4738 &mut self,
4739 instance: Instance,
4740 caller: RuntimeComponentInstanceIndex,
4741 ty: TypeStreamTableIndex,
4742 options: OptionsIndex,
4743 stream: u32,
4744 address: u32,
4745 count: u32,
4746 ) -> Result<u32>;
4747
4748 fn stream_read(
4750 &mut self,
4751 instance: Instance,
4752 caller: RuntimeComponentInstanceIndex,
4753 ty: TypeStreamTableIndex,
4754 options: OptionsIndex,
4755 stream: u32,
4756 address: u32,
4757 count: u32,
4758 ) -> Result<u32>;
4759
4760 fn flat_stream_write(
4763 &mut self,
4764 instance: Instance,
4765 caller: RuntimeComponentInstanceIndex,
4766 ty: TypeStreamTableIndex,
4767 options: OptionsIndex,
4768 payload_size: u32,
4769 payload_align: u32,
4770 stream: u32,
4771 address: u32,
4772 count: u32,
4773 ) -> Result<u32>;
4774
4775 fn flat_stream_read(
4778 &mut self,
4779 instance: Instance,
4780 caller: RuntimeComponentInstanceIndex,
4781 ty: TypeStreamTableIndex,
4782 options: OptionsIndex,
4783 payload_size: u32,
4784 payload_align: u32,
4785 stream: u32,
4786 address: u32,
4787 count: u32,
4788 ) -> Result<u32>;
4789
4790 fn stream_drop_writable(
4792 &mut self,
4793 instance: Instance,
4794 ty: TypeStreamTableIndex,
4795 writer: u32,
4796 ) -> Result<()>;
4797
4798 fn error_context_debug_message(
4800 &mut self,
4801 instance: Instance,
4802 ty: TypeComponentLocalErrorContextTableIndex,
4803 options: OptionsIndex,
4804 err_ctx_handle: u32,
4805 debug_msg_address: u32,
4806 ) -> Result<()>;
4807
4808 fn thread_new_indirect(
4810 &mut self,
4811 instance: Instance,
4812 caller: RuntimeComponentInstanceIndex,
4813 func_ty_idx: ModuleInternedTypeIndex,
4814 start_func_table_idx: RuntimeTableIndex,
4815 start_func_idx: u32,
4816 context: i32,
4817 ) -> Result<u32>;
4818}
4819
4820impl<T: 'static> VMComponentAsyncStore for StoreInner<T> {
4822 unsafe fn prepare_call(
4823 &mut self,
4824 instance: Instance,
4825 memory: *mut VMMemoryDefinition,
4826 start: NonNull<VMFuncRef>,
4827 return_: NonNull<VMFuncRef>,
4828 caller_instance: RuntimeComponentInstanceIndex,
4829 callee_instance: RuntimeComponentInstanceIndex,
4830 task_return_type: TypeTupleIndex,
4831 callee_async: bool,
4832 string_encoding: StringEncoding,
4833 result_count_or_max_if_async: u32,
4834 storage: *mut ValRaw,
4835 storage_len: usize,
4836 ) -> Result<()> {
4837 let params = unsafe { core::slice::from_raw_parts(storage, storage_len) }.to_vec();
4841
4842 unsafe {
4843 instance.prepare_call(
4844 StoreContextMut(self),
4845 start,
4846 return_,
4847 caller_instance,
4848 callee_instance,
4849 task_return_type,
4850 callee_async,
4851 memory,
4852 string_encoding,
4853 match result_count_or_max_if_async {
4854 PREPARE_ASYNC_NO_RESULT => CallerInfo::Async {
4855 params,
4856 has_result: false,
4857 },
4858 PREPARE_ASYNC_WITH_RESULT => CallerInfo::Async {
4859 params,
4860 has_result: true,
4861 },
4862 result_count => CallerInfo::Sync {
4863 params,
4864 result_count,
4865 },
4866 },
4867 )
4868 }
4869 }
4870
4871 unsafe fn start_call(
4872 &mut self,
4873 instance: Instance,
4874 callback: *mut VMFuncRef,
4875 post_return: *mut VMFuncRef,
4876 callee: NonNull<VMFuncRef>,
4877 param_count: u32,
4878 result_count: u32,
4879 flags: u32,
4880 storage: *mut MaybeUninit<ValRaw>,
4881 storage_len: usize,
4882 ) -> Result<()> {
4883 unsafe {
4884 instance.start_call(
4885 StoreContextMut(self),
4886 callback,
4887 post_return,
4888 callee,
4889 param_count,
4890 result_count,
4891 flags,
4892 core::slice::from_raw_parts_mut(storage, storage_len),
4896 )
4897 }
4898 }
4899
4900 fn future_write(
4901 &mut self,
4902 instance: Instance,
4903 caller: RuntimeComponentInstanceIndex,
4904 ty: TypeFutureTableIndex,
4905 options: OptionsIndex,
4906 future: u32,
4907 address: u32,
4908 ) -> Result<u32> {
4909 instance
4910 .guest_write(
4911 StoreContextMut(self),
4912 caller,
4913 TransmitIndex::Future(ty),
4914 options,
4915 None,
4916 future,
4917 address,
4918 1,
4919 )
4920 .map(|result| result.encode())
4921 }
4922
4923 fn future_read(
4924 &mut self,
4925 instance: Instance,
4926 caller: RuntimeComponentInstanceIndex,
4927 ty: TypeFutureTableIndex,
4928 options: OptionsIndex,
4929 future: u32,
4930 address: u32,
4931 ) -> Result<u32> {
4932 instance
4933 .guest_read(
4934 StoreContextMut(self),
4935 caller,
4936 TransmitIndex::Future(ty),
4937 options,
4938 None,
4939 future,
4940 address,
4941 1,
4942 )
4943 .map(|result| result.encode())
4944 }
4945
4946 fn stream_write(
4947 &mut self,
4948 instance: Instance,
4949 caller: RuntimeComponentInstanceIndex,
4950 ty: TypeStreamTableIndex,
4951 options: OptionsIndex,
4952 stream: u32,
4953 address: u32,
4954 count: u32,
4955 ) -> Result<u32> {
4956 instance
4957 .guest_write(
4958 StoreContextMut(self),
4959 caller,
4960 TransmitIndex::Stream(ty),
4961 options,
4962 None,
4963 stream,
4964 address,
4965 count,
4966 )
4967 .map(|result| result.encode())
4968 }
4969
4970 fn stream_read(
4971 &mut self,
4972 instance: Instance,
4973 caller: RuntimeComponentInstanceIndex,
4974 ty: TypeStreamTableIndex,
4975 options: OptionsIndex,
4976 stream: u32,
4977 address: u32,
4978 count: u32,
4979 ) -> Result<u32> {
4980 instance
4981 .guest_read(
4982 StoreContextMut(self),
4983 caller,
4984 TransmitIndex::Stream(ty),
4985 options,
4986 None,
4987 stream,
4988 address,
4989 count,
4990 )
4991 .map(|result| result.encode())
4992 }
4993
4994 fn future_drop_writable(
4995 &mut self,
4996 instance: Instance,
4997 ty: TypeFutureTableIndex,
4998 writer: u32,
4999 ) -> Result<()> {
5000 instance.guest_drop_writable(self, TransmitIndex::Future(ty), writer)
5001 }
5002
5003 fn flat_stream_write(
5004 &mut self,
5005 instance: Instance,
5006 caller: RuntimeComponentInstanceIndex,
5007 ty: TypeStreamTableIndex,
5008 options: OptionsIndex,
5009 payload_size: u32,
5010 payload_align: u32,
5011 stream: u32,
5012 address: u32,
5013 count: u32,
5014 ) -> Result<u32> {
5015 instance
5016 .guest_write(
5017 StoreContextMut(self),
5018 caller,
5019 TransmitIndex::Stream(ty),
5020 options,
5021 Some(FlatAbi {
5022 size: payload_size,
5023 align: payload_align,
5024 }),
5025 stream,
5026 address,
5027 count,
5028 )
5029 .map(|result| result.encode())
5030 }
5031
5032 fn flat_stream_read(
5033 &mut self,
5034 instance: Instance,
5035 caller: RuntimeComponentInstanceIndex,
5036 ty: TypeStreamTableIndex,
5037 options: OptionsIndex,
5038 payload_size: u32,
5039 payload_align: u32,
5040 stream: u32,
5041 address: u32,
5042 count: u32,
5043 ) -> Result<u32> {
5044 instance
5045 .guest_read(
5046 StoreContextMut(self),
5047 caller,
5048 TransmitIndex::Stream(ty),
5049 options,
5050 Some(FlatAbi {
5051 size: payload_size,
5052 align: payload_align,
5053 }),
5054 stream,
5055 address,
5056 count,
5057 )
5058 .map(|result| result.encode())
5059 }
5060
5061 fn stream_drop_writable(
5062 &mut self,
5063 instance: Instance,
5064 ty: TypeStreamTableIndex,
5065 writer: u32,
5066 ) -> Result<()> {
5067 instance.guest_drop_writable(self, TransmitIndex::Stream(ty), writer)
5068 }
5069
5070 fn error_context_debug_message(
5071 &mut self,
5072 instance: Instance,
5073 ty: TypeComponentLocalErrorContextTableIndex,
5074 options: OptionsIndex,
5075 err_ctx_handle: u32,
5076 debug_msg_address: u32,
5077 ) -> Result<()> {
5078 instance.error_context_debug_message(
5079 StoreContextMut(self),
5080 ty,
5081 options,
5082 err_ctx_handle,
5083 debug_msg_address,
5084 )
5085 }
5086
5087 fn thread_new_indirect(
5088 &mut self,
5089 instance: Instance,
5090 caller: RuntimeComponentInstanceIndex,
5091 func_ty_idx: ModuleInternedTypeIndex,
5092 start_func_table_idx: RuntimeTableIndex,
5093 start_func_idx: u32,
5094 context: i32,
5095 ) -> Result<u32> {
5096 instance.thread_new_indirect(
5097 StoreContextMut(self),
5098 caller,
5099 func_ty_idx,
5100 start_func_table_idx,
5101 start_func_idx,
5102 context,
5103 )
5104 }
5105}
5106
5107type HostTaskFuture = Pin<Box<dyn Future<Output = Result<()>> + Send + 'static>>;
5108
5109async fn run_with_host_task_set<F>(task: TableId<HostTask>, future: F) -> Result<F::Output>
5112where
5113 F: Future,
5114{
5115 let mut future = pin!(future);
5116 future::poll_fn(|cx| {
5117 let old_thread = match tls::get(|store| store.set_thread(task)) {
5118 Ok(thread) => thread,
5119 Err(error) => return Poll::Ready(Err(error)),
5120 };
5121 let result = future.as_mut().poll(cx);
5122 match tls::get(|store| store.set_thread(old_thread)) {
5123 Ok(_) => result.map(Ok),
5124 Err(error) => Poll::Ready(Err(error)),
5125 }
5126 })
5127 .await
5128}
5129
5130pub(crate) struct HostTask {
5134 common: WaitableCommon,
5135
5136 call_context: CallContext,
5139
5140 state: HostTaskState,
5141
5142 group: TaskGroupId,
5143}
5144
5145enum HostTaskState {
5146 CalleeStarted,
5151
5152 CalleeRunning(JoinHandle),
5157
5158 CalleeCancelling,
5162
5163 CalleeFinished(LiftedResult),
5167
5168 CalleeDone { cancelled: bool },
5171}
5172
5173impl HostTask {
5174 fn new(
5175 concurrent_state: &mut ConcurrentState,
5176 state: HostTaskState,
5177 caller: QualifiedThreadId,
5178 ) -> Result<Self> {
5179 let group = concurrent_state.get_mut(caller.task)?.group;
5180 concurrent_state.increment_group_ref_count(group)?;
5181
5182 Ok(Self {
5183 common: WaitableCommon::default(),
5184 call_context: CallContext::default(),
5185 state,
5186 group,
5187 })
5188 }
5189}
5190
5191impl TableDebug for HostTask {
5192 fn type_name() -> &'static str {
5193 "HostTask"
5194 }
5195}
5196
5197type CallbackFn = Box<dyn Fn(&mut dyn VMStore, Event, u32) -> Result<u32> + Send + Sync + 'static>;
5198
5199enum Caller {
5201 Host {
5203 tx: Option<oneshot::Sender<LiftedResult>>,
5205 host_future_present: bool,
5208 caller: Option<TableId<HostTask>>,
5212 },
5213 Guest {
5215 thread: QualifiedThreadId,
5217 },
5218}
5219
5220struct LiftResult {
5223 lift: RawLift,
5224 ty: TypeTupleIndex,
5225 memory: Option<SendSyncPtr<VMMemoryDefinition>>,
5226 string_encoding: StringEncoding,
5227}
5228
5229#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
5234pub(crate) struct QualifiedThreadId {
5235 task: TableId<GuestTask>,
5236 thread: TableId<GuestThread>,
5237}
5238
5239impl QualifiedThreadId {
5240 fn qualify(
5241 state: &mut ConcurrentState,
5242 thread: TableId<GuestThread>,
5243 ) -> Result<QualifiedThreadId> {
5244 Ok(QualifiedThreadId {
5245 task: state.get_mut(thread)?.parent_task,
5246 thread,
5247 })
5248 }
5249}
5250
5251impl fmt::Debug for QualifiedThreadId {
5252 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
5253 f.debug_tuple("QualifiedThreadId")
5254 .field(&self.task.rep())
5255 .field(&self.thread.rep())
5256 .finish()
5257 }
5258}
5259
5260enum GuestThreadState {
5261 NotStartedImplicit,
5262 NotStartedExplicit(
5263 Box<dyn FnOnce(&mut dyn VMStore, QualifiedThreadId) -> Result<()> + Send + Sync>,
5264 ),
5265 Running,
5266 Suspended(StoreFiber<'static>),
5267 Ready {
5268 fiber: StoreFiber<'static>,
5269 },
5270 Completed,
5271}
5272
5273impl fmt::Debug for GuestThreadState {
5274 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
5275 match self {
5276 Self::NotStartedImplicit => f.debug_tuple("NotStartedImplicit").finish(),
5277 Self::NotStartedExplicit(_) => f.debug_tuple("NotStartedExplicit").finish(),
5278 Self::Running => f.debug_tuple("Running").finish(),
5279 Self::Suspended(_) => f.debug_tuple("Suspended").finish(),
5280 Self::Ready { .. } => f.debug_struct("Ready").finish(),
5281 Self::Completed => f.debug_tuple("Completed").finish(),
5282 }
5283 }
5284}
5285
5286#[derive(Copy, Clone, PartialEq, Eq, Debug)]
5287enum WakeOnCancel {
5288 None,
5289 Waiting(TableId<WaitableSet>),
5290 Yielding,
5291}
5292
5293impl WakeOnCancel {
5294 fn is_none(self) -> bool {
5295 matches!(self, WakeOnCancel::None)
5296 }
5297
5298 fn replace(&mut self, other: WakeOnCancel) -> Self {
5299 let old = *self;
5300 *self = other;
5301 old
5302 }
5303
5304 fn take(&mut self) -> Self {
5305 self.replace(WakeOnCancel::None)
5306 }
5307}
5308
5309pub struct GuestThread {
5310 context: [u32; NUM_COMPONENT_CONTEXT_SLOTS],
5313 parent_task: TableId<GuestTask>,
5315 wake_on_cancel: WakeOnCancel,
5318 state: GuestThreadState,
5320 instance_rep: Option<u32>,
5323 sync_call_set: TableId<WaitableSet>,
5325 old_do_not_suspend: Option<bool>,
5328}
5329
5330impl GuestThread {
5331 fn from_instance(
5334 state: Pin<&mut ComponentInstance>,
5335 caller_instance: RuntimeComponentInstanceIndex,
5336 guest_thread: u32,
5337 ) -> Result<TableId<Self>> {
5338 let rep = state.instance_states().0[caller_instance]
5339 .thread_handle_table()
5340 .guest_thread_rep(guest_thread)?;
5341 Ok(TableId::new(rep))
5342 }
5343
5344 fn new_implicit(state: &mut ConcurrentState, parent_task: TableId<GuestTask>) -> Result<Self> {
5345 let sync_call_set = state.push(WaitableSet {
5346 is_sync_call_set: true,
5347 ..WaitableSet::default()
5348 })?;
5349 Ok(Self {
5350 context: [0; NUM_COMPONENT_CONTEXT_SLOTS],
5351 parent_task,
5352 wake_on_cancel: WakeOnCancel::None,
5353 state: GuestThreadState::NotStartedImplicit,
5354 instance_rep: None,
5355 sync_call_set,
5356 old_do_not_suspend: None,
5357 })
5358 }
5359
5360 fn new_explicit(
5361 state: &mut ConcurrentState,
5362 parent_task: TableId<GuestTask>,
5363 start_func: Box<
5364 dyn FnOnce(&mut dyn VMStore, QualifiedThreadId) -> Result<()> + Send + Sync,
5365 >,
5366 ) -> Result<Self> {
5367 let sync_call_set = state.push(WaitableSet {
5368 is_sync_call_set: true,
5369 ..WaitableSet::default()
5370 })?;
5371 Ok(Self {
5372 context: [0; NUM_COMPONENT_CONTEXT_SLOTS],
5373 parent_task,
5374 wake_on_cancel: WakeOnCancel::None,
5375 state: GuestThreadState::NotStartedExplicit(start_func),
5376 instance_rep: None,
5377 sync_call_set,
5378 old_do_not_suspend: None,
5379 })
5380 }
5381}
5382
5383impl TableDebug for GuestThread {
5384 fn type_name() -> &'static str {
5385 "GuestThread"
5386 }
5387}
5388
5389enum SyncResult {
5390 NotProduced,
5391 Produced(Option<ValRaw>),
5392 Taken,
5393}
5394
5395impl SyncResult {
5396 fn take(&mut self) -> Result<Option<Option<ValRaw>>> {
5397 Ok(match mem::replace(self, SyncResult::Taken) {
5398 SyncResult::NotProduced => None,
5399 SyncResult::Produced(val) => Some(val),
5400 SyncResult::Taken => {
5401 bail_bug!("attempted to take a synchronous result that was already taken")
5402 }
5403 })
5404 }
5405}
5406
5407#[derive(Debug)]
5408enum HostFutureState {
5409 NotApplicable,
5410 Live,
5411 Dropped,
5412}
5413
5414pub(crate) struct GuestTask {
5416 common: WaitableCommon,
5418 lower_params: Option<RawLower>,
5420 lift_result: Option<LiftResult>,
5422 result: Option<LiftedResult>,
5425 callback: Option<CallbackFn>,
5428 caller: Caller,
5430 call_context: CallContext,
5435 sync_result: SyncResult,
5438 cancel_request_delivered: bool,
5442 starting_sent: bool,
5445 instance: RuntimeInstance,
5452 event: Option<Event>,
5454 exited: bool,
5456 threads: HashSet<TableId<GuestThread>>,
5458 host_future_state: HostFutureState,
5461 async_typed: bool,
5464 async_lifted: bool,
5467
5468 decremented_interesting_task_count: bool,
5469
5470 group: TaskGroupId,
5471}
5472
5473impl GuestTask {
5474 fn already_lowered_parameters(&self) -> bool {
5475 self.lower_params.is_none()
5477 }
5478
5479 fn returned_or_cancelled(&self) -> bool {
5480 self.lift_result.is_none()
5482 }
5483
5484 fn ready_to_delete(&self) -> bool {
5485 let threads_completed = self.threads.is_empty();
5486 let has_sync_result = matches!(self.sync_result, SyncResult::Produced(_));
5487 let pending_completion_event = matches!(
5488 self.common.event,
5489 Some(Event::Subtask {
5490 status: Status::Returned | Status::ReturnCancelled
5491 })
5492 );
5493 let ready = threads_completed
5494 && !has_sync_result
5495 && !pending_completion_event
5496 && !matches!(self.host_future_state, HostFutureState::Live);
5497 log::trace!(
5498 "ready to delete? {ready} (threads_completed: {}, has_sync_result: {}, pending_completion_event: {}, host_future_state: {:?})",
5499 threads_completed,
5500 has_sync_result,
5501 pending_completion_event,
5502 self.host_future_state
5503 );
5504 ready
5505 }
5506
5507 fn new(
5508 state: &mut ConcurrentState,
5509 lower_params: RawLower,
5510 lift_result: LiftResult,
5511 caller: Caller,
5512 callback: Option<CallbackFn>,
5513 instance: RuntimeInstance,
5514 async_typed: bool,
5515 async_lifted: bool,
5516 ) -> Result<QualifiedThreadId> {
5517 let host_future_state = match &caller {
5518 Caller::Guest { .. } => HostFutureState::NotApplicable,
5519 Caller::Host {
5520 host_future_present,
5521 ..
5522 } => {
5523 if *host_future_present {
5524 HostFutureState::Live
5525 } else {
5526 HostFutureState::NotApplicable
5527 }
5528 }
5529 };
5530
5531 let group = match caller {
5532 Caller::Guest { thread } => {
5533 let group = state.get_mut(thread.task)?.group;
5534 state.increment_group_ref_count(group)?;
5535 group
5536 }
5537 Caller::Host { .. } => state.make_task_group()?,
5538 };
5539
5540 let task = state.push(Self {
5541 common: WaitableCommon::default(),
5542 lower_params: Some(lower_params),
5543 lift_result: Some(lift_result),
5544 result: None,
5545 callback,
5546 caller,
5547 call_context: CallContext::default(),
5548 sync_result: SyncResult::NotProduced,
5549 cancel_request_delivered: false,
5550 starting_sent: false,
5551 instance,
5552 event: None,
5553 exited: false,
5554 threads: HashSet::new(),
5555 host_future_state,
5556 async_typed,
5557 async_lifted,
5558 decremented_interesting_task_count: false,
5559 group,
5560 })?;
5561 let new_thread = GuestThread::new_implicit(state, task)?;
5562 let thread = state.push(new_thread)?;
5563 state.get_mut(task)?.threads.insert(thread);
5564 state.interesting_tasks += 1;
5565 let thread = QualifiedThreadId { task, thread };
5566 log::trace!("new implicit thread {thread:?} for instance {instance:?}");
5567 Ok(thread)
5568 }
5569}
5570
5571impl TableDebug for GuestTask {
5572 fn type_name() -> &'static str {
5573 "GuestTask"
5574 }
5575}
5576
5577#[derive(Default)]
5579struct WaitableCommon {
5580 event: Option<Event>,
5582 set: Option<TableId<WaitableSet>>,
5584 handle: Option<u32>,
5586}
5587
5588#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
5590enum Waitable {
5591 Host(TableId<HostTask>),
5593 Guest(TableId<GuestTask>),
5595 Transmit(TableId<TransmitHandle>),
5597}
5598
5599impl Waitable {
5600 fn from_instance(
5603 state: Pin<&mut ComponentInstance>,
5604 caller_instance: RuntimeComponentInstanceIndex,
5605 waitable: u32,
5606 ) -> Result<Self> {
5607 use crate::runtime::vm::component::Waitable;
5608
5609 let (waitable, kind) = state.instance_states().0[caller_instance]
5610 .handle_table()
5611 .waitable_rep(waitable)?;
5612
5613 Ok(match kind {
5614 Waitable::Subtask { is_host: true } => Self::Host(TableId::new(waitable)),
5615 Waitable::Subtask { is_host: false } => Self::Guest(TableId::new(waitable)),
5616 Waitable::Stream | Waitable::Future => Self::Transmit(TableId::new(waitable)),
5617 })
5618 }
5619
5620 fn rep(&self) -> u32 {
5622 match self {
5623 Self::Host(id) => id.rep(),
5624 Self::Guest(id) => id.rep(),
5625 Self::Transmit(id) => id.rep(),
5626 }
5627 }
5628
5629 fn join(&self, state: &mut ConcurrentState, set: Option<TableId<WaitableSet>>) -> Result<()> {
5633 log::trace!("waitable {self:?} join set {set:?}");
5634
5635 let old = mem::replace(&mut self.common(state)?.set, set);
5636
5637 if let Some(old) = old {
5638 match *self {
5639 Waitable::Host(id) => state.remove_child(id, old),
5640 Waitable::Guest(id) => state.remove_child(id, old),
5641 Waitable::Transmit(id) => state.remove_child(id, old),
5642 }?;
5643
5644 state.get_mut(old)?.ready.remove(self);
5645 }
5646
5647 if let Some(set) = set {
5648 match *self {
5649 Waitable::Host(id) => state.add_child(id, set),
5650 Waitable::Guest(id) => state.add_child(id, set),
5651 Waitable::Transmit(id) => state.add_child(id, set),
5652 }?;
5653
5654 if self.common(state)?.event.is_some() {
5655 self.mark_ready(state)?;
5656 }
5657 }
5658
5659 Ok(())
5660 }
5661
5662 fn common<'a>(&self, state: &'a mut ConcurrentState) -> Result<&'a mut WaitableCommon> {
5664 Ok(match self {
5665 Self::Host(id) => &mut state.get_mut(*id)?.common,
5666 Self::Guest(id) => &mut state.get_mut(*id)?.common,
5667 Self::Transmit(id) => &mut state.get_mut(*id)?.common,
5668 })
5669 }
5670
5671 fn trap_if_in_waitable_set(&self, state: &mut ConcurrentState) -> Result<()> {
5677 if self.common(state)?.set.is_some() {
5678 bail!(Trap::WaitableSyncAndAsync);
5679 }
5680 Ok(())
5681 }
5682
5683 fn set_event(&self, state: &mut ConcurrentState, event: Option<Event>) -> Result<()> {
5687 log::trace!("set event for {self:?}: {event:?}");
5688 self.common(state)?.event = event;
5689 self.mark_ready(state)
5690 }
5691
5692 fn take_event(&self, state: &mut ConcurrentState) -> Result<Option<Event>> {
5694 let common = self.common(state)?;
5695 let event = common.event.take();
5696 if let Some(set) = self.common(state)?.set {
5697 state.get_mut(set)?.ready.remove(self);
5698 }
5699
5700 Ok(event)
5701 }
5702
5703 fn mark_ready(&self, state: &mut ConcurrentState) -> Result<()> {
5707 if let Some(set) = self.common(state)?.set {
5708 state.get_mut(set)?.ready.insert(*self);
5709 state.wake_waiter(set)?;
5710 }
5711 Ok(())
5712 }
5713
5714 fn delete_from(&self, store: &mut StoreOpaque) -> Result<()> {
5716 match self {
5717 Self::Host(task) => {
5718 log::trace!("delete host task {task:?}");
5719 let state = store.concurrent_state_mut()?;
5720 let task = state.delete(*task)?;
5721
5722 state.decrement_group_ref_count(task.group)?;
5723 }
5724 Self::Guest(task) => {
5725 log::trace!("delete guest task {task:?}");
5726 let state = store.concurrent_state_mut()?;
5727 let task = state.delete(*task)?;
5728
5729 state.decrement_group_ref_count(task.group)?;
5730
5731 debug_assert!(task.decremented_interesting_task_count);
5738 }
5739 Self::Transmit(task) => {
5740 store.concurrent_state_mut()?.delete(*task)?;
5741 }
5742 }
5743
5744 Ok(())
5745 }
5746}
5747
5748impl fmt::Debug for Waitable {
5749 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
5750 match self {
5751 Self::Host(id) => write!(f, "{id:?}"),
5752 Self::Guest(id) => write!(f, "{id:?}"),
5753 Self::Transmit(id) => write!(f, "{id:?}"),
5754 }
5755 }
5756}
5757
5758#[derive(Default)]
5760struct WaitableSet {
5761 ready: BTreeSet<Waitable>,
5763 waiting: BTreeMap<QualifiedThreadId, WaitMode>,
5765 num_waiting: usize,
5768 is_sync_call_set: bool,
5771}
5772
5773impl WaitableSet {
5774 fn stop_waiting(&mut self) -> Result<()> {
5776 self.num_waiting = match self.num_waiting.checked_sub(1) {
5777 Some(n) => n,
5778 None => bail_bug!("waiter not accounted for in waitable set"),
5779 };
5780 Ok(())
5781 }
5782}
5783
5784impl TableDebug for WaitableSet {
5785 fn type_name() -> &'static str {
5786 "WaitableSet"
5787 }
5788}
5789
5790type RawLower =
5792 Box<dyn FnOnce(&mut dyn VMStore, &mut [MaybeUninit<ValRaw>]) -> Result<()> + Send + Sync>;
5793
5794type RawLift = Box<
5796 dyn FnOnce(&mut dyn VMStore, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>> + Send + Sync,
5797>;
5798
5799type LiftedResult = Box<dyn Any + Send + Sync>;
5803
5804struct DummyResult;
5807
5808#[derive(Default)]
5810pub struct ConcurrentInstanceState {
5811 backpressure: u16,
5813 do_not_enter: bool,
5815 do_not_suspend: bool,
5818 pending: BTreeMap<QualifiedThreadId, GuestCallKind>,
5821}
5822
5823impl ConcurrentInstanceState {
5824 pub fn pending_is_empty(&self) -> bool {
5825 self.pending.is_empty()
5826 }
5827}
5828
5829#[derive(Debug, Copy, Clone)]
5830pub(crate) enum CurrentThread {
5831 Guest(QualifiedThreadId),
5834 Host(TableId<HostTask>),
5836 DeferredHost(QualifiedThreadId),
5839 None,
5842}
5843
5844impl CurrentThread {
5845 fn guest(&self) -> Option<&QualifiedThreadId> {
5846 match self {
5847 Self::Guest(id) => Some(id),
5848 _ => None,
5849 }
5850 }
5851
5852 fn guest_task(&self) -> Option<TableId<GuestTask>> {
5853 match self {
5854 Self::Guest(id) => Some(id.task),
5855 _ => None,
5856 }
5857 }
5858
5859 fn is_none(&self) -> bool {
5860 matches!(self, Self::None)
5861 }
5862}
5863
5864impl From<QualifiedThreadId> for CurrentThread {
5865 fn from(id: QualifiedThreadId) -> Self {
5866 Self::Guest(id)
5867 }
5868}
5869
5870impl From<TableId<HostTask>> for CurrentThread {
5871 fn from(id: TableId<HostTask>) -> Self {
5872 Self::Host(id)
5873 }
5874}
5875
5876enum Priority {
5877 Switch,
5878 High,
5879 Low,
5880}
5881
5882pub struct ConcurrentState {
5884 unforced_current_thread: CurrentThread,
5890
5891 deferred_host_call_context: Option<CallContext>,
5897
5898 futures: AlwaysMut<Option<FuturesUnordered<HostTaskFuture>>>,
5903 table: AlwaysMut<ResourceTable>,
5905 switch_item: Option<WorkItem>,
5913 next_switch_item: Option<WorkItem>,
5919 high_priority: VecDeque<WorkItem>,
5921 low_priority: VecDeque<WorkItem>,
5923 suspend_reason: Option<SuspendReason>,
5927 worker: Option<StoreFiber<'static>>,
5931 worker_item: Option<WorkerItem>,
5933
5934 global_error_context_ref_counts:
5947 BTreeMap<TypeComponentGlobalErrorContextTableIndex, GlobalErrorContextRefCount>,
5948
5949 interesting_tasks: usize,
5962
5963 interesting_tasks_empty_waker: Option<Waker>,
5967
5968 ready_for_concurrent_call_waker: Option<Waker>,
5973
5974 event_loop_running: bool,
5976
5977 saved_next_switch_items: Vec<Option<WorkItem>>,
5980
5981 #[cfg(feature = "task-group-hook")]
5983 task_group_hook: Option<Box<dyn TaskGroupHook>>,
5984}
5985
5986impl Default for ConcurrentState {
5987 fn default() -> Self {
5988 Self {
5989 unforced_current_thread: CurrentThread::None,
5990 deferred_host_call_context: None,
5991 table: AlwaysMut::new(ResourceTable::new()),
5992 futures: AlwaysMut::new(Some(FuturesUnordered::new())),
5993 switch_item: None,
5994 next_switch_item: None,
5995 high_priority: VecDeque::new(),
5996 low_priority: VecDeque::new(),
5997 suspend_reason: None,
5998 worker: None,
5999 worker_item: None,
6000 global_error_context_ref_counts: BTreeMap::new(),
6001 interesting_tasks: 0,
6002 interesting_tasks_empty_waker: None,
6003 ready_for_concurrent_call_waker: None,
6004 event_loop_running: false,
6005 saved_next_switch_items: Vec::new(),
6006 #[cfg(feature = "task-group-hook")]
6007 task_group_hook: None,
6008 }
6009 }
6010}
6011
6012impl ConcurrentState {
6013 pub(crate) fn take_fibers_and_futures(
6030 &mut self,
6031 fibers: &mut Vec<StoreFiber<'static>>,
6032 futures: &mut Vec<FuturesUnordered<HostTaskFuture>>,
6033 ) {
6034 let mut items = Vec::new();
6035 for (_, entry) in self.table.get_mut().iter_mut() {
6036 if let Some(set) = entry.downcast_mut::<WaitableSet>() {
6037 for mode in mem::take(&mut set.waiting).into_values() {
6038 match mode {
6039 WaitMode::Fiber(fiber) => {
6040 fibers.push(fiber);
6041 }
6042 WaitMode::Callback(_) => {}
6043 }
6044 }
6045 } else if let Some(thread) = entry.downcast_mut::<GuestThread>() {
6046 if let GuestThreadState::Suspended(fiber) | GuestThreadState::Ready { fiber, .. } =
6047 mem::replace(&mut thread.state, GuestThreadState::Completed)
6048 {
6049 fibers.push(fiber);
6050 }
6051 } else if let Some(item) = entry.downcast_mut::<Option<WorkItem>>() {
6052 if let Some(item) = item.take() {
6053 items.push(item);
6054 }
6055 }
6056 }
6057
6058 if let Some(fiber) = self.worker.take() {
6059 fibers.push(fiber);
6060 }
6061
6062 let mut handle_item = |item| match item {
6063 WorkItem::ResumeFiber { fiber, .. } => {
6064 fibers.push(fiber);
6065 }
6066 WorkItem::PushFuture(future) => {
6067 self.futures
6068 .get_mut()
6069 .as_mut()
6070 .unwrap()
6071 .push(future.into_inner());
6072 }
6073 WorkItem::ResumeThread { .. }
6074 | WorkItem::GuestCall { .. }
6075 | WorkItem::WorkerFunction(_) => {}
6076 };
6077
6078 for item in items {
6079 handle_item(item);
6080 }
6081 if let Some(item) = self.switch_item.take() {
6082 handle_item(item);
6083 }
6084 if let Some(item) = self.next_switch_item.take() {
6085 handle_item(item);
6086 }
6087 for item in mem::take(&mut self.high_priority) {
6088 handle_item(item);
6089 }
6090 for item in mem::take(&mut self.low_priority) {
6091 handle_item(item);
6092 }
6093 for item in mem::take(&mut self.saved_next_switch_items)
6094 .into_iter()
6095 .filter_map(|v| v)
6096 {
6097 handle_item(item);
6098 }
6099
6100 if let Some(them) = self.futures.get_mut().take() {
6101 futures.push(them);
6102 }
6103 }
6104
6105 #[cfg(feature = "gc")]
6106 pub(crate) fn trace_fiber_roots(
6107 &mut self,
6108 modules: &ModuleRegistry,
6109 unwind: &dyn Unwind,
6110 gc_roots_list: &mut GcRootsList,
6111 ) {
6112 let ConcurrentState {
6113 table,
6114 worker,
6115 switch_item,
6116 next_switch_item,
6117 high_priority,
6118 low_priority,
6119 saved_next_switch_items,
6120
6121 futures: _,
6125
6126 worker_item: _,
6128 unforced_current_thread: _,
6129 deferred_host_call_context: _,
6130 suspend_reason: _,
6131 global_error_context_ref_counts: _,
6132 interesting_tasks: _,
6133 interesting_tasks_empty_waker: _,
6134 ready_for_concurrent_call_waker: _,
6135 event_loop_running: _,
6136 #[cfg(feature = "task-group-hook")]
6137 task_group_hook: _,
6138 } = self;
6139
6140 for (_, entry) in table.get_mut().iter_mut() {
6141 if let Some(set) = entry.downcast_mut::<WaitableSet>() {
6142 for mode in set.waiting.values_mut() {
6143 match mode {
6144 WaitMode::Fiber(fiber) => {
6145 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6146 }
6147 WaitMode::Callback(_) => {}
6148 }
6149 }
6150 } else if let Some(thread) = entry.downcast_mut::<GuestThread>() {
6151 if let GuestThreadState::Suspended(fiber) | GuestThreadState::Ready { fiber, .. } =
6152 &mut thread.state
6153 {
6154 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6155 }
6156 } else if let Some(Some(WorkItem::ResumeFiber { fiber, .. })) =
6157 entry.downcast_mut::<Option<WorkItem>>()
6158 {
6159 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6160 }
6161 }
6162
6163 if let Some(fiber) = worker {
6164 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6165 }
6166
6167 let mut handle_item = |item: &mut WorkItem| match item {
6168 WorkItem::ResumeFiber { fiber, .. } => {
6169 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6170 }
6171 WorkItem::PushFuture(_future) => {
6172 }
6175 WorkItem::ResumeThread { .. }
6176 | WorkItem::GuestCall { .. }
6177 | WorkItem::WorkerFunction(_) => {}
6178 };
6179
6180 if let Some(item) = switch_item {
6181 handle_item(item);
6182 }
6183 if let Some(item) = next_switch_item {
6184 handle_item(item);
6185 }
6186 for item in high_priority {
6187 handle_item(item);
6188 }
6189 for item in low_priority {
6190 handle_item(item);
6191 }
6192 for item in saved_next_switch_items
6193 .iter_mut()
6194 .filter_map(|v| v.as_mut())
6195 {
6196 handle_item(item);
6197 }
6198 }
6199
6200 fn push<V: Send + Sync + 'static>(
6201 &mut self,
6202 value: V,
6203 ) -> Result<TableId<V>, ResourceTableError> {
6204 self.table.get_mut().push(value).map(TableId::from)
6205 }
6206
6207 fn get_mut<V: 'static>(&mut self, id: TableId<V>) -> Result<&mut V, ResourceTableError> {
6208 self.table.get_mut().get_mut(&Resource::from(id))
6209 }
6210
6211 pub fn add_child<T: 'static, U: 'static>(
6212 &mut self,
6213 child: TableId<T>,
6214 parent: TableId<U>,
6215 ) -> Result<(), ResourceTableError> {
6216 self.table
6217 .get_mut()
6218 .add_child(Resource::from(child), Resource::from(parent))
6219 }
6220
6221 pub fn remove_child<T: 'static, U: 'static>(
6222 &mut self,
6223 child: TableId<T>,
6224 parent: TableId<U>,
6225 ) -> Result<(), ResourceTableError> {
6226 self.table
6227 .get_mut()
6228 .remove_child(Resource::from(child), Resource::from(parent))
6229 }
6230
6231 fn delete<V: 'static>(&mut self, id: TableId<V>) -> Result<V, ResourceTableError> {
6232 self.table.get_mut().delete(Resource::from(id))
6233 }
6234
6235 fn push_future(&mut self, future: HostTaskFuture) {
6236 self.push_high_priority(WorkItem::PushFuture(AlwaysMut::new(future)));
6243 }
6244
6245 fn set_switch_item(&mut self, item: WorkItem) -> Result<()> {
6246 log::trace!("set switch item: {item:?}");
6247
6248 if self.switch_item.is_some() {
6249 bail_bug!("switch item already set");
6250 }
6251
6252 self.switch_item = Some(item);
6253
6254 Ok(())
6255 }
6256
6257 fn take_next_switch_item(&mut self) -> Result<()> {
6258 if let Some(item) = self.next_switch_item.take() {
6259 self.set_switch_item(item)?;
6260 }
6261 Ok(())
6262 }
6263
6264 fn wake_waiter(&mut self, set: TableId<WaitableSet>) -> Result<()> {
6267 let Some((thread, mode)) = self.get_mut(set)?.waiting.pop_first() else {
6268 return Ok(());
6269 };
6270 let wake_on_cancel = self.get_mut(thread.thread)?.wake_on_cancel.take();
6271 assert!(wake_on_cancel.is_none() || wake_on_cancel == WakeOnCancel::Waiting(set));
6272
6273 let instance = self.get_mut(thread.task)?.instance;
6274 let item = match mode {
6275 WaitMode::Fiber(fiber) => WorkItem::ResumeFiber {
6276 instance,
6277 thread,
6278 fiber,
6279 },
6280 WaitMode::Callback(callback_instance) => WorkItem::GuestCall {
6281 instance,
6282 call: GuestCall {
6283 thread,
6284 kind: GuestCallKind::DeliverEvent {
6285 instance: callback_instance,
6286 set: Some(set),
6287 },
6288 },
6289 },
6290 };
6291 self.push_high_priority(item);
6292 Ok(())
6293 }
6294
6295 fn push_high_priority(&mut self, item: WorkItem) {
6296 log::trace!("push high priority: {item:?}");
6297 self.high_priority.push_front(item);
6298 }
6299
6300 fn push_low_priority(&mut self, item: WorkItem) {
6301 log::trace!("push low priority: {item:?}");
6302 self.low_priority.push_front(item);
6303 }
6304
6305 fn push_work_item(&mut self, item: WorkItem, priority: Priority) -> Result<()> {
6306 match priority {
6307 Priority::Switch => self.set_switch_item(item)?,
6308 Priority::High => self.push_high_priority(item),
6309 Priority::Low => self.push_low_priority(item),
6310 }
6311
6312 Ok(())
6313 }
6314
6315 fn promote_instance_local_thread_work_item(
6316 &mut self,
6317 current_instance: RuntimeInstance,
6318 ) -> Result<bool> {
6319 log::trace!("promote thread work items for {current_instance:?}");
6320
6321 self.promote_work_item_matching(|item: &WorkItem| {
6322 let result = match item {
6323 WorkItem::ResumeThread { instance, .. }
6324 | WorkItem::ResumeFiber { instance, .. }
6325 | WorkItem::GuestCall { instance, .. } => *instance == current_instance,
6326 _ => false,
6327 };
6328
6329 log::trace!("candidate {item:?}: {result}");
6330 result
6331 })
6332 }
6333
6334 fn promote_thread_work_item(&mut self, thread: QualifiedThreadId) -> Result<bool> {
6335 self.promote_work_item_matching(|item: &WorkItem| match item {
6336 WorkItem::ResumeThread {
6337 thread: item_thread,
6338 ..
6339 }
6340 | WorkItem::GuestCall {
6341 call:
6342 GuestCall {
6343 thread: item_thread,
6344 ..
6345 },
6346 ..
6347 } => *item_thread == thread,
6348 _ => false,
6349 })
6350 }
6351
6352 fn promote_work_item_matching<F>(&mut self, mut predicate: F) -> Result<bool>
6353 where
6354 F: FnMut(&WorkItem) -> bool,
6355 {
6356 for item in mem::take(&mut self.high_priority).into_iter().rev() {
6361 if self.switch_item.is_none() && predicate(&item) {
6362 self.set_switch_item(item)?;
6363 } else {
6364 self.push_high_priority(item);
6365 }
6366 }
6367
6368 if self.switch_item.is_none() {
6369 for item in mem::take(&mut self.low_priority).into_iter().rev() {
6370 if self.switch_item.is_none() && predicate(&item) {
6371 self.set_switch_item(item)?;
6372 } else {
6373 self.push_low_priority(item);
6374 }
6375 }
6376 }
6377
6378 Ok(self.switch_item.is_some())
6379 }
6380
6381 pub fn call_context(&mut self, task: Scope) -> Result<&mut CallContext> {
6384 match task {
6385 Scope::HostId(task) => {
6386 let task: TableId<HostTask> = TableId::new(task);
6387 Ok(&mut self.get_mut(task)?.call_context)
6388 }
6389 Scope::Id(task) => {
6390 let task: TableId<GuestTask> = TableId::new(task);
6391 Ok(&mut self.get_mut(task)?.call_context)
6392 }
6393 }
6394 }
6395
6396 pub(crate) fn deferred_host_call_context(&mut self) -> Option<&mut CallContext> {
6397 self.deferred_host_call_context.as_mut()
6398 }
6399
6400 fn futures_mut(&mut self) -> Result<&mut FuturesUnordered<HostTaskFuture>> {
6401 match self.futures.get_mut().as_mut() {
6402 Some(f) => Ok(f),
6403 None => bail_bug!("futures field of concurrent state is currently taken"),
6404 }
6405 }
6406
6407 pub(crate) fn table(&mut self) -> &mut ResourceTable {
6408 self.table.get_mut()
6409 }
6410
6411 fn debug_assert_deferred_host_invariant(&self) {
6412 debug_assert_eq!(
6413 self.deferred_host_call_context.is_some(),
6414 matches!(self.unforced_current_thread, CurrentThread::DeferredHost(_)),
6415 "a deferred host thread and call context must exist together",
6416 );
6417 }
6418
6419 fn materialize_host_task(&mut self) -> Result<CurrentThread> {
6420 self.debug_assert_deferred_host_invariant();
6421 let caller = match self.unforced_current_thread {
6422 CurrentThread::DeferredHost(caller) => caller,
6423 thread => return Ok(thread),
6424 };
6425
6426 let task = HostTask::new(self, HostTaskState::CalleeStarted, caller)?;
6428 let task = self.push(task)?;
6429 let call_context = self
6430 .deferred_host_call_context
6431 .take()
6432 .expect("deferred host call context should be present");
6433 self.get_mut(task)
6434 .expect("newly inserted host task should be present")
6435 .call_context = call_context;
6436 self.unforced_current_thread = CurrentThread::Host(task);
6437 self.debug_assert_deferred_host_invariant();
6438 log::trace!("new host task materialized {task:?}");
6439 Ok(CurrentThread::Host(task))
6440 }
6441
6442 fn materialize_current_host_task_id(&mut self) -> Result<Option<TableId<HostTask>>> {
6443 match self.materialize_host_task()? {
6444 CurrentThread::Host(id) => Ok(Some(id)),
6445 CurrentThread::None => Ok(None),
6446 CurrentThread::Guest(_) => {
6447 bail_bug!("tried to materialize a host task id from a guest thread")
6448 }
6449 CurrentThread::DeferredHost(_) => {
6450 bail_bug!(
6451 "current thread is a deferred host thread which should have been materialized"
6452 )
6453 }
6454 }
6455 }
6456
6457 pub(crate) fn materialize_current_scope(&mut self) -> Result<Scope> {
6458 match self.materialize_host_task()? {
6459 CurrentThread::Host(id) => Ok(Scope::HostId(id.rep())),
6460 _ => bail_bug!("current scope is not a deferred host scope"),
6461 }
6462 }
6463}
6464
6465fn for_any_lower<
6468 F: FnOnce(&mut dyn VMStore, &mut [MaybeUninit<ValRaw>]) -> Result<()> + Send + Sync,
6469>(
6470 fun: F,
6471) -> F {
6472 fun
6473}
6474
6475fn for_any_lift<
6477 F: FnOnce(&mut dyn VMStore, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>> + Send + Sync,
6478>(
6479 fun: F,
6480) -> F {
6481 fun
6482}
6483
6484fn check_ambient_store(id: StoreId) {
6485 let message = "\
6486 `Future`s which depend on asynchronous component tasks, streams, or \
6487 futures to complete may only be polled from the event loop of the \
6488 store to which they belong. Please use \
6489 `StoreContextMut::{run_concurrent,spawn}` to poll or await them.\
6490 ";
6491 tls::try_get(|store| {
6492 let matched = match store {
6493 tls::TryGet::Some(store) => store.id() == id,
6494 tls::TryGet::Taken | tls::TryGet::None => false,
6495 };
6496
6497 if !matched {
6498 panic!("{message}")
6499 }
6500 });
6501}
6502
6503fn unpack_callback_code(code: u32) -> (u32, u32) {
6504 (code & 0xF, code >> 4)
6505}
6506
6507struct WaitableCheckParams {
6511 set: TableId<WaitableSet>,
6512 options: OptionsIndex,
6513 payload: u32,
6514}
6515
6516enum WaitableCheck {
6519 Wait,
6520 Poll,
6521}
6522
6523pub(crate) struct PreparedCall<R> {
6525 handle: Func,
6527 thread: QualifiedThreadId,
6529 param_count: usize,
6531 rx: oneshot::Receiver<LiftedResult>,
6534 runtime_instance: RuntimeInstance,
6536 _phantom: PhantomData<R>,
6537}
6538
6539impl<R> PreparedCall<R> {
6540 pub(crate) fn task_id(&self) -> TaskId {
6542 TaskId {
6543 task: self.thread.task,
6544 runtime_instance: self.runtime_instance,
6545 }
6546 }
6547}
6548
6549pub(crate) struct TaskId {
6551 task: TableId<GuestTask>,
6552 runtime_instance: RuntimeInstance,
6553}
6554
6555impl TaskId {
6556 pub(crate) fn host_future_dropped(&self, store: &mut StoreOpaque) -> Result<()> {
6562 let task = store.concurrent_state_mut()?.get_mut(self.task)?;
6563 let delete = if !task.already_lowered_parameters() {
6564 store.cancel_guest_subtask_without_lowered_parameters(
6565 self.runtime_instance,
6566 self.task,
6567 )?;
6568 true
6569 } else {
6570 task.host_future_state = HostFutureState::Dropped;
6571 task.ready_to_delete()
6572 };
6573 if delete {
6574 Waitable::Guest(self.task).delete_from(store)?
6575 }
6576 Ok(())
6577 }
6578}
6579
6580pub(crate) fn prepare_call<T, R>(
6586 mut store: StoreContextMut<T>,
6587 handle: Func,
6588 param_count: usize,
6589 host_future_present: bool,
6590 lower_params: impl FnOnce(StoreContextMut<T>, &mut [MaybeUninit<ValRaw>]) -> Result<()>
6591 + Send
6592 + Sync
6593 + 'static,
6594 lift_result: impl FnOnce(&mut StoreOpaque, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>>
6595 + Send
6596 + Sync
6597 + 'static,
6598) -> Result<PreparedCall<R>> {
6599 if !store.0.may_enter() {
6600 bail!(Trap::CannotEnterComponent);
6601 }
6602
6603 let (options, _flags, ty, raw_options) = handle.abi_info(store.0);
6604
6605 let instance = handle.instance().id().get(store.0);
6606 let options = &instance.component().env_component().options[options];
6607 let ty = &instance.component().types()[ty];
6608 let async_typed = ty.async_;
6609 let async_lifted = raw_options.async_;
6610 let task_return_type = ty.results;
6611 let component_instance = raw_options.instance;
6612 let callback = options.callback.map(|i| instance.runtime_callback(i));
6613 let memory = options
6614 .memory()
6615 .map(|i| instance.runtime_memory(i))
6616 .map(SendSyncPtr::new);
6617 let string_encoding = options.string_encoding;
6618 let token = StoreToken::new(store.as_context_mut());
6619 let caller = store.0.materialize_host_task_id()?;
6620 let state = store.0.concurrent_state_mut()?;
6621
6622 let (tx, rx) = oneshot::channel();
6623
6624 let instance = handle.instance().runtime_instance(component_instance);
6625 let thread = GuestTask::new(
6626 state,
6627 Box::new(for_any_lower(move |store, params| {
6628 lower_params(token.as_context_mut(store), params)
6629 })),
6630 LiftResult {
6631 lift: Box::new(for_any_lift(move |store, result| {
6632 lift_result(store, result)
6633 })),
6634 ty: task_return_type,
6635 memory,
6636 string_encoding,
6637 },
6638 Caller::Host {
6639 tx: Some(tx),
6640 host_future_present,
6641 caller,
6642 },
6643 callback.map(|callback| {
6644 let callback = SendSyncPtr::new(callback);
6645 let instance = handle.instance();
6646 Box::new(move |store: &mut dyn VMStore, event, handle| {
6647 let store = token.as_context_mut(store);
6648 unsafe { instance.call_callback(store, callback, event, handle) }
6651 }) as CallbackFn
6652 }),
6653 instance,
6654 async_typed,
6655 async_lifted,
6656 )?;
6657
6658 Ok(PreparedCall {
6659 handle,
6660 thread,
6661 param_count,
6662 runtime_instance: instance,
6663 rx,
6664 _phantom: PhantomData,
6665 })
6666}
6667
6668pub(crate) struct StagedCall<R> {
6669 store: StoreId,
6670 rx: oneshot::Receiver<LiftedResult>,
6671 _marker: PhantomData<fn() -> R>,
6672 group: TaskGroupId,
6673}
6674
6675impl<R> StagedCall<R> {
6676 pub(crate) fn new<T: 'static>(
6683 mut store: StoreContextMut<T>,
6684 prepared: PreparedCall<R>,
6685 ) -> Result<StagedCall<R>> {
6686 let PreparedCall {
6687 handle,
6688 thread,
6689 param_count,
6690 rx,
6691 ..
6692 } = prepared;
6693
6694 stage_call0(store.as_context_mut(), handle, thread, param_count)?;
6695
6696 Ok(StagedCall {
6697 store: store.0.id(),
6698 rx,
6699 _marker: PhantomData,
6700 group: store.0.concurrent_state_mut()?.get_mut(thread.task)?.group,
6701 })
6702 }
6703}
6704
6705impl<R> Future for StagedCall<R>
6706where
6707 R: 'static,
6708{
6709 type Output = Result<R>;
6710
6711 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
6712 check_ambient_store(self.store);
6713 Pin::new(&mut self.rx).poll(cx).map(|result| match result {
6714 Ok(r) => match r.downcast() {
6715 Ok(r) => Ok(*r),
6716 Err(_) => bail_bug!("wrong type of value produced"),
6717 },
6718 Err(oneshot::Canceled) => bail_bug!("channel erroneously dropped"),
6719 })
6720 }
6721}
6722
6723fn stage_call0<T: 'static>(
6726 store: StoreContextMut<T>,
6727 handle: Func,
6728 guest_thread: QualifiedThreadId,
6729 param_count: usize,
6730) -> Result<()> {
6731 let (_options, _, _ty, raw_options) = handle.abi_info(store.0);
6732 let is_concurrent = raw_options.async_;
6733 let callback = raw_options.callback;
6734 let instance = handle.instance();
6735 let callee = handle.lifted_core_func(store.0);
6736 let post_return = raw_options
6737 .post_return
6738 .map(|i| instance.id().get(store.0).runtime_post_return(i));
6739 let callback = callback.map(|i| {
6740 let instance = instance.id().get(store.0);
6741 SendSyncPtr::new(instance.runtime_callback(i))
6742 });
6743
6744 log::trace!("queueing call {guest_thread:?}");
6745
6746 unsafe {
6750 instance.stage_call(
6751 store,
6752 guest_thread,
6753 SendSyncPtr::new(callee),
6754 param_count,
6755 1,
6756 is_concurrent,
6757 callback,
6758 post_return.map(SendSyncPtr::new),
6759 true,
6760 )
6761 }
6762}
6763
6764#[cfg(all(test, feature = "cranelift", feature = "wat"))]
6765mod tests {
6766 use super::*;
6767 use crate::component::{Component, Linker};
6768 use crate::store::AsStoreOpaque;
6769 use crate::{Config, Engine};
6770
6771 fn host_subtask(
6772 state: HostTaskState,
6773 event: Option<Event>,
6774 ) -> Result<(Store<()>, Instance, TableId<HostTask>, u32)> {
6775 let mut config = Config::new();
6776 config.wasm_component_model_async(true);
6777 let engine = Engine::new(&config)?;
6778 let component = Component::new(&engine, "(component)")?;
6779 let mut store = Store::new(&engine, ());
6780 let instance = Linker::new(&engine).instantiate(&mut store, &component)?;
6781 let store_opaque = store.as_store_opaque();
6782 let concurrent_state = store_opaque.concurrent_state_mut()?;
6783 let group = concurrent_state.make_task_group()?;
6786 let task = concurrent_state.push(HostTask {
6787 common: WaitableCommon::default(),
6788 call_context: CallContext::default(),
6789 state,
6790 group,
6791 })?;
6792 let handle = store_opaque
6793 .instance_state(instance.runtime_instance(RuntimeComponentInstanceIndex::from_u32(0)))
6794 .handle_table()
6795 .subtask_insert_host(task.rep())?;
6796 let common = &mut store_opaque.concurrent_state_mut()?.get_mut(task)?.common;
6797 common.handle = Some(handle);
6798 common.event = event;
6799 Ok((store, instance, task, handle))
6800 }
6801
6802 #[test]
6803 fn host_subtask_drop_during_cancellation() -> Result<()> {
6804 for abort_completed in [false, true] {
6805 let (handle, future) = JoinHandle::run(future::pending::<()>());
6806 let mut future = pin!(future);
6807 let (mut store, instance, task, handle) =
6808 host_subtask(HostTaskState::CalleeRunning(handle), None)?;
6809 let store = store.as_store_opaque();
6810 let caller = RuntimeComponentInstanceIndex::from_u32(0);
6811 assert_eq!(
6812 instance.subtask_cancel(store, caller, true, handle)?,
6813 BLOCKED
6814 );
6815 if abort_completed {
6816 assert!(matches!(
6819 future
6820 .as_mut()
6821 .poll(&mut Context::from_waker(Waker::noop())),
6822 Poll::Ready(None),
6823 ));
6824 }
6825 for async_ in [false, true] {
6826 let err = instance
6827 .subtask_cancel(store, caller, async_, handle)
6828 .unwrap_err();
6829 assert_eq!(err.downcast::<Trap>()?, Trap::SubtaskCancelAfterTerminal);
6830 }
6831 let err = instance.subtask_drop(store, caller, handle).unwrap_err();
6832 assert_eq!(err.downcast::<Trap>()?, Trap::SubtaskDropNotResolved);
6833 assert!(store.concurrent_state_mut()?.get_mut(task).is_ok());
6834 }
6835 Ok(())
6836 }
6837
6838 #[test]
6839 fn host_subtask_cancel_after_completion() -> Result<()> {
6840 for async_ in [false, true] {
6841 let (mut store, instance, task, handle) = host_subtask(
6842 HostTaskState::CalleeDone { cancelled: false },
6843 Some(Event::Subtask {
6844 status: Status::Returned,
6845 }),
6846 )?;
6847 let store = store.as_store_opaque();
6848 let caller = RuntimeComponentInstanceIndex::from_u32(0);
6849 assert_eq!(
6850 instance.subtask_cancel(store, caller, async_, handle)?,
6851 Status::Returned as u32,
6852 );
6853 let err = instance
6854 .subtask_cancel(store, caller, async_, handle)
6855 .unwrap_err();
6856 assert_eq!(err.downcast::<Trap>()?, Trap::SubtaskCancelAfterTerminal);
6857 instance.subtask_drop(store, caller, handle)?;
6858 assert!(store.concurrent_state_mut()?.get_mut(task).is_err());
6859 }
6860 Ok(())
6861 }
6862
6863 #[test]
6864 fn host_subtask_drop_requires_terminal_event_delivery() -> Result<()> {
6865 for (cancelled, status) in [
6866 (false, Status::Returned),
6867 (true, Status::Returned),
6868 (true, Status::ReturnCancelled),
6869 ] {
6870 for delivered in [false, true] {
6871 let event = if delivered {
6872 None
6873 } else {
6874 Some(Event::Subtask { status })
6875 };
6876 let (mut store, instance, task, handle) =
6877 host_subtask(HostTaskState::CalleeDone { cancelled }, event)?;
6878 let store = store.as_store_opaque();
6879 let result = instance.subtask_drop(
6880 store,
6881 RuntimeComponentInstanceIndex::from_u32(0),
6882 handle,
6883 );
6884 if delivered {
6885 result?;
6886 assert!(store.concurrent_state_mut()?.get_mut(task).is_err());
6887 } else {
6888 let err = result.unwrap_err();
6889 assert_eq!(err.downcast::<Trap>()?, Trap::SubtaskDropNotResolved);
6890 assert!(store.concurrent_state_mut()?.get_mut(task).is_ok());
6891 }
6892 }
6893 }
6894 Ok(())
6895 }
6896}