1use self::error_contexts::GlobalErrorContextRefCount;
54use crate::component::func::{Func, call_post_return};
55use crate::component::{
56 HasData, HasSelf, Instance, Resource, ResourceTable, ResourceTableError, RuntimeInstance,
57};
58use crate::fiber::{self, StoreFiber, StoreFiberYield};
59use crate::hash_set::HashSet;
60#[cfg(feature = "gc")]
61use crate::module::ModuleRegistry;
62use crate::prelude::*;
63use crate::store::{Store, StoreId, StoreInner, StoreOpaque, StoreToken};
64#[cfg(feature = "gc")]
65use crate::vm::GcRootsList;
66use crate::vm::component::{CallContext, ComponentInstance, CurrentScope, InstanceState, Scope};
67use crate::vm::{AlwaysMut, SendSyncPtr, VMFuncRef, VMLazyThread, VMMemoryDefinition, VMStore};
68use crate::{
69 AsContext, AsContextMut, FuncType, Result, StoreContext, StoreContextMut, ValRaw, ValType, bail,
70};
71use crate::{Instance as ModuleInstance, bail_bug};
72use alloc::borrow::ToOwned;
73use alloc::collections::{BTreeMap, BTreeSet, VecDeque};
74use core::any::Any;
75use core::cell::UnsafeCell;
76use core::fmt;
77use core::future;
78use core::future::Future;
79use core::marker::PhantomData;
80use core::mem::{self, ManuallyDrop, MaybeUninit};
81use core::ops::DerefMut;
82use core::pin::{Pin, pin};
83use core::ptr::{self, NonNull};
84use core::task::{Context, Poll, Waker};
85use futures::channel::oneshot;
86use futures::stream::{FuturesUnordered, StreamExt};
87use futures_and_streams::{FlatAbi, ReturnCode, TransmitHandle, TransmitIndex};
88use table::{TableDebug, TableId};
89use wasmtime_environ::component::{
90 CanonicalAbiInfo, CanonicalOptions, CanonicalOptionsDataModel, MAX_FLAT_PARAMS,
91 MAX_FLAT_RESULTS, OptionsIndex, PREPARE_ASYNC_NO_RESULT, PREPARE_ASYNC_WITH_RESULT,
92 RuntimeComponentInstanceIndex, RuntimeTableIndex, StringEncoding,
93 TypeComponentGlobalErrorContextTableIndex, TypeComponentLocalErrorContextTableIndex,
94 TypeFuncIndex, TypeFutureTableIndex, TypeStreamTableIndex, TypeTupleIndex,
95};
96use wasmtime_environ::packed_option::ReservedValue;
97use wasmtime_environ::{NUM_COMPONENT_CONTEXT_SLOTS, Trap};
98#[cfg(feature = "gc")]
99use wasmtime_unwinder::Unwind;
100
101pub use abort::JoinHandle;
102pub use func::{FuncCallConcurrent, TypedFuncCallConcurrent};
103pub use future_stream_any::{FutureAny, StreamAny};
104pub use futures_and_streams::{
105 Destination, DirectDestination, DirectSource, ErrorContext, FutureConsumer, FutureProducer,
106 FutureReader, GuardedFutureReader, GuardedStreamReader, ReadBuffer, Source, StreamConsumer,
107 StreamProducer, StreamReader, StreamResult, VecBuffer, WriteBuffer,
108};
109pub(crate) use futures_and_streams::{ResourcePair, lower_error_context_to_index};
110#[cfg(feature = "task-group-hook")]
111pub use task_group_hook::TaskGroupHook;
112pub use task_group_hook::TaskGroupId;
113
114mod abort;
115mod error_contexts;
116mod func;
117mod future_stream_any;
118mod futures_and_streams;
119pub(crate) mod table;
120#[cfg(feature = "task-group-hook")]
121mod task_group_hook;
122#[cfg(not(feature = "task-group-hook"))]
123mod task_group_hook_disabled;
124#[cfg(not(feature = "task-group-hook"))]
125use task_group_hook_disabled as task_group_hook;
126pub(crate) mod tls;
127
128const BLOCKED: u32 = 0xffff_ffff;
131
132#[derive(Clone, Copy, Eq, PartialEq, Debug)]
134pub enum Status {
135 Starting = 0,
136 Started = 1,
137 Returned = 2,
138 StartCancelled = 3,
139 ReturnCancelled = 4,
140}
141
142impl Status {
143 pub fn pack(self, waitable: Option<u32>) -> u32 {
149 assert!(matches!(self, Status::Returned) == waitable.is_none());
150 let waitable = waitable.unwrap_or(0);
151 assert!(waitable < (1 << 28));
152 (waitable << 4) | (self as u32)
153 }
154}
155
156#[derive(Clone, Copy, Debug)]
159enum Event {
160 None,
161 Subtask {
162 status: Status,
163 },
164 StreamRead {
165 code: ReturnCode,
166 pending: Option<(TypeStreamTableIndex, u32)>,
167 },
168 StreamWrite {
169 code: ReturnCode,
170 pending: Option<(TypeStreamTableIndex, u32)>,
171 },
172 FutureRead {
173 code: ReturnCode,
174 pending: Option<(TypeFutureTableIndex, u32)>,
175 },
176 FutureWrite {
177 code: ReturnCode,
178 pending: Option<(TypeFutureTableIndex, u32)>,
179 },
180 Cancelled,
181}
182
183impl Event {
184 fn parts(self) -> (u32, u32) {
189 const EVENT_NONE: u32 = 0;
190 const EVENT_SUBTASK: u32 = 1;
191 const EVENT_STREAM_READ: u32 = 2;
192 const EVENT_STREAM_WRITE: u32 = 3;
193 const EVENT_FUTURE_READ: u32 = 4;
194 const EVENT_FUTURE_WRITE: u32 = 5;
195 const EVENT_CANCELLED: u32 = 6;
196 match self {
197 Event::None => (EVENT_NONE, 0),
198 Event::Cancelled => (EVENT_CANCELLED, 0),
199 Event::Subtask { status } => (EVENT_SUBTASK, status as u32),
200 Event::StreamRead { code, .. } => (EVENT_STREAM_READ, code.encode()),
201 Event::StreamWrite { code, .. } => (EVENT_STREAM_WRITE, code.encode()),
202 Event::FutureRead { code, .. } => (EVENT_FUTURE_READ, code.encode()),
203 Event::FutureWrite { code, .. } => (EVENT_FUTURE_WRITE, code.encode()),
204 }
205 }
206}
207
208mod callback_code {
210 pub const EXIT: u32 = 0;
211 pub const YIELD: u32 = 1;
212 pub const WAIT: u32 = 2;
213}
214
215const START_FLAG_ASYNC_CALLEE: u32 = wasmtime_environ::component::START_FLAG_ASYNC_CALLEE as u32;
219
220pub struct Access<'a, T: 'static, D: HasData + ?Sized = HasSelf<T>> {
226 store: StoreContextMut<'a, T>,
227 get_data: fn(&mut T) -> D::Data<'_>,
228}
229
230impl<'a, T, D> Access<'a, T, D>
231where
232 D: HasData + ?Sized,
233 T: 'static,
234{
235 pub fn new(store: StoreContextMut<'a, T>, get_data: fn(&mut T) -> D::Data<'_>) -> Self {
237 Self { store, get_data }
238 }
239
240 pub fn data_mut(&mut self) -> &mut T {
242 self.store.data_mut()
243 }
244
245 pub fn get(&mut self) -> D::Data<'_> {
247 (self.get_data)(self.data_mut())
248 }
249
250 pub fn spawn(&mut self, task: impl for<'fut> AccessorTask<'fut, T, D>) -> Result<JoinHandle>
254 where
255 T: 'static,
256 {
257 let accessor = Accessor {
258 get_data: self.get_data,
259 token: StoreToken::new(self.store.as_context_mut()),
260 };
261 self.store
262 .as_context_mut()
263 .spawn_with_accessor(accessor, task)
264 }
265
266 pub fn getter(&self) -> fn(&mut T) -> D::Data<'_> {
269 self.get_data
270 }
271}
272
273impl<'a, T, D> AsContext for Access<'a, T, D>
274where
275 D: HasData + ?Sized,
276 T: 'static,
277{
278 type Data = T;
279
280 fn as_context(&self) -> StoreContext<'_, T> {
281 self.store.as_context()
282 }
283}
284
285impl<'a, T, D> AsContextMut for Access<'a, T, D>
286where
287 D: HasData + ?Sized,
288 T: 'static,
289{
290 fn as_context_mut(&mut self) -> StoreContextMut<'_, T> {
291 self.store.as_context_mut()
292 }
293}
294
295pub struct Accessor<T: 'static, D = HasSelf<T>>
355where
356 D: HasData + ?Sized,
357{
358 token: StoreToken<T>,
359 get_data: fn(&mut T) -> D::Data<'_>,
360}
361
362pub trait AsAccessor {
379 type Data: 'static;
381
382 type AccessorData: HasData + ?Sized;
385
386 fn as_accessor(&self) -> &Accessor<Self::Data, Self::AccessorData>;
388}
389
390impl<T: AsAccessor + ?Sized> AsAccessor for &T {
391 type Data = T::Data;
392 type AccessorData = T::AccessorData;
393
394 fn as_accessor(&self) -> &Accessor<Self::Data, Self::AccessorData> {
395 T::as_accessor(self)
396 }
397}
398
399impl<T, D: HasData + ?Sized> AsAccessor for Accessor<T, D> {
400 type Data = T;
401 type AccessorData = D;
402
403 fn as_accessor(&self) -> &Accessor<T, D> {
404 self
405 }
406}
407
408const _: () = {
431 const fn assert<T: Send + Sync>() {}
432 assert::<Accessor<UnsafeCell<u32>>>();
433};
434
435impl<T> Accessor<T> {
436 pub(crate) fn new(token: StoreToken<T>) -> Self {
445 Self {
446 token,
447 get_data: |x| x,
448 }
449 }
450}
451
452impl<T, D> Accessor<T, D>
453where
454 D: HasData + ?Sized,
455{
456 pub fn with<R>(&self, fun: impl FnOnce(Access<'_, T, D>) -> R) -> R {
474 tls::get(|vmstore| {
475 fun(Access {
476 store: self.token.as_context_mut(vmstore),
477 get_data: self.get_data,
478 })
479 })
480 }
481
482 pub fn getter(&self) -> fn(&mut T) -> D::Data<'_> {
485 self.get_data
486 }
487
488 pub fn with_getter<D2: HasData>(
505 &self,
506 get_data: fn(&mut T) -> D2::Data<'_>,
507 ) -> Accessor<T, D2> {
508 Accessor {
509 token: self.token,
510 get_data,
511 }
512 }
513
514 pub fn spawn(&self, task: impl for<'fut> AccessorTask<'fut, T, D>) -> Result<JoinHandle>
530 where
531 T: 'static,
532 {
533 let accessor = self.clone_for_spawn();
534 self.with(|mut access| access.as_context_mut().spawn_with_accessor(accessor, task))
535 }
536
537 fn clone_for_spawn(&self) -> Self {
538 Self {
539 token: self.token,
540 get_data: self.get_data,
541 }
542 }
543
544 pub fn poll_no_interesting_tasks(&self, cx: &mut Context<'_>) -> Poll<()> {
580 self.with(|mut access| {
581 let store = access.as_context_mut().0;
582 let state = store.concurrent_state_mut_without_forcing_current_thread();
583 if state.interesting_tasks == 0 {
584 Poll::Ready(())
585 } else {
586 state.interesting_tasks_empty_waker = Some(cx.waker().clone());
587 Poll::Pending
588 }
589 })
590 }
591
592 pub fn poll_ready_for_concurrent_call(&self, func: Func, cx: &mut Context<'_>) -> Poll<()> {
609 self.with(|mut access| {
610 let store = access.as_context_mut().0;
611 let (_, _, _, raw_options) = func.abi_info(store);
612 let instance = func.instance().runtime_instance(raw_options.instance);
613 let state = store.instance_state(instance).concurrent_state();
614 if state.backpressure == 0 {
615 Poll::Ready(())
616 } else {
617 store
618 .concurrent_state_mut_without_forcing_current_thread()
619 .ready_for_concurrent_call_waker = Some(cx.waker().clone());
620 Poll::Pending
621 }
622 })
623 }
624}
625
626pub trait AccessorTask<'fut, T, D = HasSelf<T>>:
648 AsyncFnOnce(&Accessor<T, D>) -> Result<()> + Send + 'static
649where
650 D: HasData + ?Sized,
651{
652 fn run(self, accessor: &'fut Accessor<T, D>) -> impl Future<Output = Result<()>> + Send + 'fut;
654}
655
656impl<'fut, F, Fut, T, D> AccessorTask<'fut, T, D> for F
657where
658 T: 'static,
659 F: AsyncFnOnce(&Accessor<T, D>) -> Result<()>,
660 F: FnOnce(&'fut Accessor<T, D>) -> Fut + Send + 'static,
661 Fut: Future<Output = Result<()>> + Send + 'fut,
662 D: HasData,
663{
664 fn run(self, accessor: &'fut Accessor<T, D>) -> impl Future<Output = Result<()>> + Send + 'fut {
665 (self)(accessor)
666 }
667}
668
669enum CallerInfo {
672 Async {
674 params: Vec<ValRaw>,
675 has_result: bool,
676 },
677 Sync {
679 params: Vec<ValRaw>,
680 result_count: u32,
681 },
682}
683
684enum WaitMode {
686 Fiber(StoreFiber<'static>),
688 Callback(Instance),
691}
692
693impl fmt::Debug for WaitMode {
694 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
695 match self {
696 Self::Fiber(_) => f.debug_tuple("Fiber").finish(),
697 Self::Callback(instance) => f.debug_tuple("Callback").field(instance).finish(),
698 }
699 }
700}
701
702#[derive(Debug)]
704enum SuspendReason {
705 Waiting {
708 set: TableId<WaitableSet>,
709 thread: QualifiedThreadId,
710 },
711 YieldingToSubtask { thread: QualifiedThreadId },
714 NeedWork,
717 Yielding { thread: QualifiedThreadId },
720 ExplicitlySuspending { thread: QualifiedThreadId },
723}
724
725enum GuestCallKind {
727 DeliverEvent {
730 instance: Instance,
732 set: Option<TableId<WaitableSet>>,
737 },
738 StartImplicit(Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send + Sync>),
744 StartExplicit(Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send + Sync>),
745}
746
747impl fmt::Debug for GuestCallKind {
748 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
749 match self {
750 Self::DeliverEvent { instance, set } => f
751 .debug_struct("DeliverEvent")
752 .field("instance", instance)
753 .field("set", set)
754 .finish(),
755 Self::StartImplicit(_) => f.debug_tuple("StartImplicit").finish(),
756 Self::StartExplicit(_) => f.debug_tuple("StartExplicit").finish(),
757 }
758 }
759}
760
761#[derive(Copy, Clone, Debug)]
763pub enum SuspensionTarget {
764 Resume(u32),
765 Promote(u32),
766 None,
767}
768
769#[derive(Copy, Clone, Debug)]
771pub enum ResumeThread {
772 Promote,
773 Resume,
774 ResumeLater,
775}
776
777#[derive(Debug)]
779struct GuestCall {
780 thread: QualifiedThreadId,
781 kind: GuestCallKind,
782}
783
784impl GuestCall {
785 fn is_ready(&self, store: &mut StoreOpaque) -> Result<bool> {
795 let task = store.concurrent_state_mut()?.get_mut(self.thread.task)?;
796 let async_typed = task.async_typed;
797 let instance = task.instance;
798 let state = store.instance_state(instance).concurrent_state();
799
800 let ready = match &self.kind {
801 GuestCallKind::DeliverEvent { .. } => !state.do_not_enter,
802 GuestCallKind::StartImplicit(_) => {
803 !async_typed || !(state.do_not_enter || state.backpressure > 0)
804 }
805 GuestCallKind::StartExplicit(_) => true,
806 };
807 log::trace!(
808 "call {self:?} ready? {ready} (do_not_enter: {}; backpressure: {})",
809 state.do_not_enter,
810 state.backpressure
811 );
812 Ok(ready)
813 }
814}
815
816enum WorkerItem {
818 GuestCall(GuestCall),
819 Function(AlwaysMut<Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send>>),
820}
821
822enum WorkItem {
825 PushFuture(AlwaysMut<HostTaskFuture>),
827 ResumeFiber {
829 instance: RuntimeInstance,
830 thread: QualifiedThreadId,
831 fiber: StoreFiber<'static>,
832 },
833 ResumeThread {
835 instance: RuntimeInstance,
836 thread: QualifiedThreadId,
837 },
838 GuestCall {
840 instance: RuntimeInstance,
841 call: GuestCall,
842 },
843 WorkerFunction(AlwaysMut<Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send>>),
845}
846
847impl fmt::Debug for WorkItem {
848 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
849 match self {
850 Self::PushFuture(_) => f.debug_tuple("PushFuture").finish(),
851 Self::ResumeFiber {
852 instance, thread, ..
853 } => f
854 .debug_struct("ResumeFiber")
855 .field("instance", instance)
856 .field("thread", thread)
857 .finish(),
858 Self::ResumeThread { instance, thread } => f
859 .debug_struct("ResumeThread")
860 .field("instance", instance)
861 .field("thread", thread)
862 .finish(),
863 Self::GuestCall { instance, call } => f
864 .debug_struct("GuestCall")
865 .field("instance", instance)
866 .field("call", call)
867 .finish(),
868 Self::WorkerFunction(_) => f.debug_tuple("WorkerFunction").finish(),
869 }
870 }
871}
872
873#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
875pub(crate) enum WaitResult {
876 Cancelled,
877 Completed,
878}
879
880pub(crate) fn poll_and_block<R: Send + Sync + 'static>(
888 store: &mut dyn VMStore,
889 host_task: EnteredHostTask,
890 future: impl Future<Output = Result<R>> + Send + 'static,
891) -> Result<R> {
892 let mut future = Box::pin(future);
899 let poll = tls::set(store, || {
900 future
901 .as_mut()
902 .poll(&mut Context::from_waker(&Waker::noop()))
903 });
904
905 let caller = match host_task {
906 Some(caller) => caller,
907 None => bail_bug!("host task wasn't created but should have been"),
908 };
909
910 let task = match poll {
911 Poll::Ready(result) => return result,
913
914 Poll::Pending => {
919 let Some(task) = store.materialize_host_task_id()? else {
920 bail_bug!("current thread is not a host thread")
921 };
922
923 let future = Box::pin(async move {
926 let result = run_with_host_task_set(task, future).await??;
927 tls::get(move |store| {
928 let state = store.concurrent_state_mut()?;
929 let host_state = &mut state.get_mut(task)?.state;
930 assert!(matches!(host_state, HostTaskState::CalleeStarted));
931 *host_state = HostTaskState::CalleeFinished(Box::new(result));
932
933 Waitable::Host(task).set_event(
934 state,
935 Some(Event::Subtask {
936 status: Status::Returned,
937 }),
938 )?;
939
940 Ok(())
941 })
942 }) as HostTaskFuture;
943
944 let caller_instance = store.concurrent_state_mut()?.get_mut(caller.task)?.instance;
945 store.switch_or_trap_if_may_not_suspend(caller_instance)?;
946
947 let state = store.concurrent_state_mut()?;
948 state.push_future(future);
949
950 let set = state.get_mut(caller.thread)?.sync_call_set;
951 Waitable::Host(task).join(state, Some(set))?;
952
953 store.suspend(SuspendReason::Waiting {
954 set,
955 thread: caller,
956 })?;
957
958 Waitable::Host(task).join(store.concurrent_state_mut()?, None)?;
962 task
963 }
964 };
965
966 let host_state = &mut store.concurrent_state_mut()?.get_mut(task)?.state;
968 match mem::replace(host_state, HostTaskState::CalleeDone { cancelled: false }) {
969 HostTaskState::CalleeFinished(result) => Ok(match result.downcast() {
970 Ok(result) => *result,
971 Err(_) => bail_bug!("host task finished with wrong type of result"),
972 }),
973 _ => bail_bug!("unexpected host task state after completion"),
974 }
975}
976
977fn handle_guest_call(store: &mut dyn VMStore, call: GuestCall) -> Result<()> {
979 match call.kind {
980 GuestCallKind::DeliverEvent { instance, set } => {
981 let (event, waitable) = match instance.get_event(store, call.thread.task, set, true)? {
982 Some(pair) => pair,
983 None => bail_bug!("delivering non-present event"),
984 };
985 let state = store.concurrent_state_mut()?;
986 let task = state.get_mut(call.thread.task)?;
987 let runtime_instance = task.instance;
988 let handle = waitable.map(|(_, v)| v).unwrap_or(0);
989
990 log::trace!(
991 "use callback to deliver event {event:?} to {:?} for {waitable:?}",
992 call.thread,
993 );
994
995 let old_thread = store.set_thread(call.thread)?;
996 log::trace!(
997 "GuestCallKind::DeliverEvent: replaced {old_thread:?} with {:?} as current thread",
998 call.thread
999 );
1000
1001 store.enter_instance(runtime_instance);
1002
1003 let Some(callback) = store
1004 .concurrent_state_mut()?
1005 .get_mut(call.thread.task)?
1006 .callback
1007 .take()
1008 else {
1009 bail_bug!("guest task callback field not present")
1010 };
1011
1012 let code = callback(store, event, handle)?;
1013
1014 store
1015 .concurrent_state_mut()?
1016 .get_mut(call.thread.task)?
1017 .callback = Some(callback);
1018
1019 store.exit_instance(runtime_instance)?;
1020
1021 store.set_thread(old_thread)?;
1022
1023 instance.handle_callback_code(store, call.thread, runtime_instance.index, code)?;
1024
1025 log::trace!("GuestCallKind::DeliverEvent: restored {old_thread:?} as current thread");
1026 }
1027 GuestCallKind::StartImplicit(fun) => {
1028 fun(store)?;
1029 }
1030 GuestCallKind::StartExplicit(fun) => {
1031 fun(store)?;
1032 }
1033 }
1034
1035 Ok(())
1036}
1037
1038impl<T> Store<T> {
1039 pub async fn run_concurrent<R>(&mut self, fun: impl AsyncFnOnce(&Accessor<T>) -> R) -> Result<R>
1041 where
1042 T: Send + 'static,
1043 {
1044 ensure!(
1045 self.as_context().0.concurrency_support(),
1046 "cannot use `run_concurrent` when Config::concurrency_support disabled",
1047 );
1048 self.as_context_mut().run_concurrent(fun).await
1049 }
1050
1051 #[doc(hidden)]
1052 pub fn assert_concurrent_state_empty(&mut self) {
1053 self.as_context_mut().assert_concurrent_state_empty();
1054 }
1055
1056 #[doc(hidden)]
1057 pub fn concurrent_state_table_size(&mut self) -> usize {
1058 self.as_context_mut().concurrent_state_table_size()
1059 }
1060
1061 pub fn spawn(
1063 &mut self,
1064 task: impl for<'fut> AccessorTask<'fut, T, HasSelf<T>>,
1065 ) -> Result<JoinHandle>
1066 where
1067 T: 'static,
1068 {
1069 self.as_context_mut().spawn(task)
1070 }
1071}
1072
1073impl<T> StoreContextMut<'_, T> {
1074 #[doc(hidden)]
1085 pub fn assert_concurrent_state_empty(self) {
1086 let store = self.0;
1087 store
1088 .store_data_mut()
1089 .components
1090 .assert_instance_states_empty();
1091 let state = store.concurrent_state_mut().unwrap();
1092 assert!(
1093 state.table.get_mut().is_empty(),
1094 "non-empty table: {:?}",
1095 state.table.get_mut()
1096 );
1097 assert!(state.switch_item.is_none());
1098 assert!(state.next_switch_item.is_none());
1099 assert!(state.high_priority.is_empty());
1100 assert!(state.low_priority.is_empty());
1101 assert!(state.unforced_current_thread.is_none());
1102 assert!(state.deferred_host_call_context.is_none());
1103 assert!(state.futures_mut().unwrap().is_empty());
1104 assert!(state.global_error_context_ref_counts.is_empty());
1105 }
1106
1107 #[doc(hidden)]
1112 pub fn concurrent_state_table_size(&mut self) -> usize {
1113 self.0
1114 .concurrent_state_mut()
1115 .unwrap()
1116 .table
1117 .get_mut()
1118 .iter_mut()
1119 .count()
1120 }
1121
1122 pub fn spawn(mut self, task: impl for<'fut> AccessorTask<'fut, T>) -> Result<JoinHandle>
1132 where
1133 T: 'static,
1134 {
1135 let accessor = Accessor::new(StoreToken::new(self.as_context_mut()));
1136 self.spawn_with_accessor(accessor, task)
1137 }
1138
1139 fn spawn_with_accessor<D>(
1142 self,
1143 accessor: Accessor<T, D>,
1144 task: impl for<'fut> AccessorTask<'fut, T, D>,
1145 ) -> Result<JoinHandle>
1146 where
1147 T: 'static,
1148 D: HasData + ?Sized,
1149 {
1150 let (handle, future) = JoinHandle::run(async move { task.run(&accessor).await });
1154 self.0
1155 .concurrent_state_mut()?
1156 .push_future(Box::pin(async move { future.await.unwrap_or(Ok(())) }));
1157 Ok(handle)
1158 }
1159
1160 pub async fn run_concurrent<R>(self, fun: impl AsyncFnOnce(&Accessor<T>) -> R) -> Result<R>
1244 where
1245 T: Send + 'static,
1246 {
1247 ensure!(
1248 self.0.concurrency_support(),
1249 "cannot use `run_concurrent` when Config::concurrency_support disabled",
1250 );
1251 self.do_run_concurrent(fun, false).await
1252 }
1253
1254 pub(super) async fn run_concurrent_trap_on_idle<R>(
1255 self,
1256 fun: impl AsyncFnOnce(&Accessor<T>) -> R,
1257 ) -> Result<R> {
1258 self.do_run_concurrent(fun, true).await
1259 }
1260
1261 async fn do_run_concurrent<R>(
1262 mut self,
1263 fun: impl AsyncFnOnce(&Accessor<T>) -> R,
1264 trap_on_idle: bool,
1265 ) -> Result<R> {
1266 debug_assert!(self.0.concurrency_support());
1267 let already_running = self
1268 .0
1269 .concurrent_state_mut_already_forced_current_thread()
1270 .event_loop_running;
1271 if already_running {
1272 bail!("Recursive `StoreContextMut::run_concurrent` calls not supported")
1273 }
1274 let token = StoreToken::new(self.as_context_mut());
1275
1276 struct Dropper<'a, T: 'static, V> {
1277 store: StoreContextMut<'a, T>,
1278 value: ManuallyDrop<V>,
1279 }
1280
1281 impl<'a, T, V> Drop for Dropper<'a, T, V> {
1282 fn drop(&mut self) {
1283 self.store
1284 .0
1285 .concurrent_state_mut_already_forced_current_thread()
1286 .event_loop_running = false;
1287
1288 tls::set(self.store.0, || {
1289 unsafe { ManuallyDrop::drop(&mut self.value) }
1294 });
1295 }
1296 }
1297
1298 let accessor = &Accessor::new(token);
1299 self.0
1300 .concurrent_state_mut_already_forced_current_thread()
1301 .event_loop_running = true;
1302 let dropper = &mut Dropper {
1303 store: self,
1304 value: ManuallyDrop::new(fun(accessor)),
1305 };
1306 let future = unsafe { Pin::new_unchecked(dropper.value.deref_mut()) };
1308
1309 let result = dropper
1310 .store
1311 .as_context_mut()
1312 .poll_until(future, trap_on_idle)
1313 .await;
1314
1315 if result.is_err() {
1316 dropper.store.0.set_trapped();
1317 }
1318
1319 result
1320 }
1321
1322 async fn poll_until<R>(
1328 mut self,
1329 mut future: Pin<&mut impl Future<Output = R>>,
1330 trap_on_idle: bool,
1331 ) -> Result<R> {
1332 struct Reset<'a, T: 'static> {
1333 store: StoreContextMut<'a, T>,
1334 futures: Option<FuturesUnordered<HostTaskFuture>>,
1335 }
1336
1337 impl<'a, T> Drop for Reset<'a, T> {
1338 fn drop(&mut self) {
1339 if let Some(futures) = self.futures.take() {
1340 *self
1341 .store
1342 .0
1343 .concurrent_state_mut_already_forced_current_thread()
1344 .futures
1345 .get_mut() = Some(futures);
1346 }
1347 }
1348 }
1349
1350 const MAX_TURNS_WITHOUT_YIELD: usize = 128;
1354 let mut turns_without_yield = 0;
1355
1356 loop {
1357 let futures = self.0.concurrent_state_mut()?.futures.get_mut().take();
1361 let mut reset = Reset {
1362 store: self.as_context_mut(),
1363 futures,
1364 };
1365 let mut next = match reset.futures.as_mut() {
1366 Some(f) => pin!(f.next()),
1367 None => bail_bug!("concurrent state missing futures field"),
1368 };
1369
1370 enum PollResult<R> {
1371 Complete(R),
1372 ProcessWork {
1373 ready: Option<WorkItem>,
1374 low_priority: bool,
1375 },
1376 }
1377
1378 let result = future::poll_fn(|cx| {
1379 if let Poll::Ready(value) = tls::set(reset.store.0, || future.as_mut().poll(cx)) {
1382 return Poll::Ready(Ok(PollResult::Complete(value)));
1383 }
1384
1385 let next = match tls::set(reset.store.0, || next.as_mut().poll(cx)) {
1389 Poll::Ready(Some(output)) => {
1390 match output {
1391 Err(e) => return Poll::Ready(Err(e)),
1392 Ok(()) => {}
1393 }
1394 Poll::Ready(true)
1395 }
1396 Poll::Ready(None) => Poll::Ready(false),
1397 Poll::Pending => Poll::Pending,
1398 };
1399
1400 let state = reset.store.0.concurrent_state_mut()?;
1415 let mut ready = state.switch_item.take();
1416 let mut low_priority = false;
1417 if ready.is_none() {
1418 ready = state.high_priority.pop_back();
1419 if ready.is_none() {
1420 ready = state.low_priority.pop_back();
1421 low_priority = true;
1422 }
1423 }
1424 if ready.is_some() {
1425 return Poll::Ready(Ok(PollResult::ProcessWork {
1426 ready,
1427 low_priority,
1428 }));
1429 }
1430
1431 return match next {
1435 Poll::Ready(true) => {
1436 Poll::Ready(Ok(PollResult::ProcessWork {
1442 ready: None,
1443 low_priority: false,
1444 }))
1445 }
1446 Poll::Ready(false) => {
1447 if let Poll::Ready(value) =
1451 tls::set(reset.store.0, || future.as_mut().poll(cx))
1452 {
1453 Poll::Ready(Ok(PollResult::Complete(value)))
1454 } else {
1455 if trap_on_idle {
1461 Poll::Ready(Err(if reset.store.0.any_may_not_suspend()? {
1468 Trap::CannotBlockSyncTask.into()
1469 } else {
1470 Trap::AsyncDeadlock.into()
1472 }))
1473 } else {
1474 Poll::Pending
1478 }
1479 }
1480 }
1481 Poll::Pending => Poll::Pending,
1486 };
1487 })
1488 .await;
1489
1490 drop(reset);
1494
1495 match result? {
1496 PollResult::Complete(value) => break Ok(value),
1499 PollResult::ProcessWork {
1502 ready,
1503 low_priority,
1504 } => {
1505 struct Dispose<'a, T: 'static> {
1506 store: StoreContextMut<'a, T>,
1507 ready: Option<WorkItem>,
1508 }
1509
1510 impl<'a, T> Drop for Dispose<'a, T> {
1511 fn drop(&mut self) {
1512 if let Some(item) = self.ready.take() {
1513 match item {
1514 WorkItem::ResumeFiber { mut fiber, .. } => {
1515 fiber.dispose(self.store.0)
1516 }
1517 WorkItem::PushFuture(future) => {
1518 tls::set(self.store.0, move || drop(future))
1519 }
1520 _ => {}
1521 }
1522 }
1523 }
1524 }
1525
1526 let mut dispose = Dispose {
1527 store: self.as_context_mut(),
1528 ready,
1529 };
1530
1531 if low_priority {
1553 dispose.store.0.yield_now().await;
1554 turns_without_yield = 0;
1555 }
1556
1557 if let Some(item) = dispose.ready.take() {
1558 dispose
1559 .store
1560 .as_context_mut()
1561 .handle_work_item(item)
1562 .await?;
1563 }
1564
1565 turns_without_yield += 1;
1566 if turns_without_yield == MAX_TURNS_WITHOUT_YIELD {
1567 turns_without_yield = 0;
1568 dispose.store.0.yield_now().await;
1569 }
1570 }
1571 }
1572 }
1573 }
1574
1575 async fn handle_work_item(self, item: WorkItem) -> Result<()> {
1577 log::trace!("handle work item {item:?}");
1578 match item {
1579 WorkItem::PushFuture(future) => {
1580 self.0
1581 .concurrent_state_mut()?
1582 .futures_mut()?
1583 .push(future.into_inner());
1584 }
1585 WorkItem::ResumeFiber { fiber, .. } => {
1586 self.0.resume_fiber(fiber).await?;
1587 }
1588 WorkItem::ResumeThread { thread, .. } => {
1589 if let GuestThreadState::Ready { fiber, .. } = mem::replace(
1590 &mut self.0.concurrent_state_mut()?.get_mut(thread.thread)?.state,
1591 GuestThreadState::Running,
1592 ) {
1593 self.0.resume_fiber(fiber).await?;
1594 } else {
1595 bail_bug!("cannot resume non-pending thread {thread:?}");
1596 }
1597 }
1598 WorkItem::GuestCall { call, .. } => {
1599 if call.is_ready(self.0)? {
1600 self.0
1601 .concurrent_state_mut()?
1602 .get_mut(call.thread.thread)?
1603 .wake_on_cancel = WakeOnCancel::None;
1604 self.run_on_worker(WorkerItem::GuestCall(call)).await?;
1605 } else {
1606 let state = self.0.concurrent_state_mut()?;
1607 let task = state.get_mut(call.thread.task)?;
1608 if !task.starting_sent {
1609 task.starting_sent = true;
1610 if let GuestCallKind::StartImplicit(_) = &call.kind {
1611 Waitable::Guest(call.thread.task).set_event(
1612 state,
1613 Some(Event::Subtask {
1614 status: Status::Starting,
1615 }),
1616 )?;
1617 }
1618 }
1619
1620 let instance = state.get_mut(call.thread.task)?.instance;
1621 self.0
1622 .instance_state(instance)
1623 .concurrent_state()
1624 .pending
1625 .insert(call.thread, call.kind);
1626
1627 self.0.concurrent_state_mut()?.take_next_switch_item()?;
1631 }
1632 }
1633 WorkItem::WorkerFunction(fun) => {
1634 self.run_on_worker(WorkerItem::Function(fun)).await?;
1635 }
1636 }
1637
1638 Ok(())
1639 }
1640
1641 async fn run_on_worker(self, item: WorkerItem) -> Result<()> {
1643 let worker = if let Some(fiber) = self.0.concurrent_state_mut()?.worker.take() {
1644 fiber
1645 } else {
1646 unsafe {
1665 fiber::make_fiber_unchecked(self.0, move |store| {
1666 loop {
1667 let Some(item) = store.concurrent_state_mut()?.worker_item.take() else {
1668 bail_bug!("worker_item not present when resuming fiber")
1669 };
1670 match item {
1671 WorkerItem::GuestCall(call) => handle_guest_call(store, call)?,
1672 WorkerItem::Function(fun) => fun.into_inner()(store)?,
1673 }
1674
1675 store.suspend(SuspendReason::NeedWork)?;
1676 }
1677 })?
1678 }
1679 };
1680
1681 let worker_item = &mut self.0.concurrent_state_mut()?.worker_item;
1682 assert!(worker_item.is_none());
1683 *worker_item = Some(item);
1684
1685 self.0.resume_fiber(worker).await
1686 }
1687
1688 pub(crate) fn wrap_call<F, R>(self, closure: F) -> impl Future<Output = Result<R>> + 'static
1693 where
1694 T: 'static,
1695 F: FnOnce(&Accessor<T>) -> Pin<Box<dyn Future<Output = Result<R>> + Send + '_>>
1696 + Send
1697 + Sync
1698 + 'static,
1699 R: Send + Sync + 'static,
1700 {
1701 let token = StoreToken::new(self);
1702 async move {
1703 let mut accessor = Accessor::new(token);
1704 closure(&mut accessor).await
1705 }
1706 }
1707
1708 pub(crate) async fn start_instance(
1709 &mut self,
1710 instance: ModuleInstance,
1711 ) -> Result<ModuleInstance> {
1712 let (tx, rx) = oneshot::channel();
1713 let token = StoreToken::new(self.as_context_mut());
1714 self.0.queue_task(move |store| {
1715 _ = tx.send(
1716 instance
1717 .start_raw(&mut token.as_context_mut(store))
1718 .map(|()| instance),
1719 );
1720 Ok(())
1721 })?;
1722 self.as_context_mut()
1723 .run_concurrent_trap_on_idle(async |_| {
1724 rx.await
1725 .map_err(|_| format_err!("oneshot channel canceled"))
1726 })
1727 .await??
1728 }
1729}
1730
1731pub type EnteredHostTask = Option<QualifiedThreadId>;
1738
1739impl StoreOpaque {
1740 #[inline]
1744 pub(crate) fn current_thread(&mut self) -> Result<CurrentThread> {
1745 if !self.concurrency_support() {
1747 return Ok(CurrentThread::None);
1748 }
1749
1750 if !self
1753 .vm_store_context_mut()
1754 .current_thread_mut()
1755 .is_deferred()
1756 {
1757 return Ok(self
1758 .concurrent_state_mut_already_forced_current_thread()
1759 .unforced_current_thread);
1760 }
1761
1762 self.force_deferred_current_thread()
1763 }
1764
1765 #[cold]
1768 fn force_deferred_current_thread(&mut self) -> Result<CurrentThread> {
1769 let state = self.concurrent_state_mut_without_forcing_current_thread();
1778 let id = match state.unforced_current_thread.guest_task() {
1779 Some(task) => state.get_mut(task)?.instance.instance,
1780 None => bail_bug!("deferred component-model thread with non-guest base"),
1781 };
1782
1783 let mut frames = Vec::new();
1786 let mut cur = *self.vm_store_context_mut().current_thread_mut();
1787 while let Some(ptr) = cur.as_deferred() {
1788 let deferred = unsafe { ptr.as_non_null().as_ref() };
1793 frames.push((
1794 deferred.callee_async != 0,
1795 deferred.callee_instance,
1796 deferred.saved_context,
1797 ));
1798 cur = deferred.parent;
1799 }
1800
1801 *self.vm_store_context_mut().current_thread_mut() = VMLazyThread::forced();
1805
1806 let current_context = *self.vm_store_context_mut().component_context_mut();
1809
1810 for (callee_async, callee_instance, saved_context) in frames.into_iter().rev() {
1814 *self.vm_store_context_mut().component_context_mut() = saved_context;
1818 let callee = RuntimeInstance {
1819 instance: id,
1820 index: RuntimeComponentInstanceIndex::from_u32(callee_instance),
1821 };
1822 self.enter_guest_sync_call(callee_async, callee)?;
1823 }
1824
1825 *self.vm_store_context_mut().component_context_mut() = current_context;
1827
1828 Ok(self
1829 .concurrent_state_mut_without_forcing_current_thread()
1830 .unforced_current_thread)
1831 }
1832
1833 fn current_guest_thread(&mut self) -> Result<QualifiedThreadId> {
1834 match self.current_thread()?.guest() {
1835 Some(id) => Ok(*id),
1836 None => bail_bug!("current thread is not a guest thread"),
1837 }
1838 }
1839
1840 pub(crate) fn current_materialized_host_task(&mut self) -> Result<Option<TableId<HostTask>>> {
1844 match self.current_thread()? {
1845 CurrentThread::Host(id) => Ok(Some(id)),
1846 CurrentThread::DeferredHost(_) | CurrentThread::None => Ok(None),
1847 _ => bail_bug!("current thread is not a host thread"),
1848 }
1849 }
1850
1851 fn materialize_host_task_id(&mut self) -> Result<Option<TableId<HostTask>>> {
1854 Ok(self
1855 .concurrent_state_mut()?
1856 .materialize_current_host_task_id()?)
1857 }
1858
1859 fn enter_sync_call(&mut self, callee: RuntimeInstance) -> Result<()> {
1860 log::trace!("enter sync-typed call {callee:?}");
1861 let state = self.instance_state(callee).concurrent_state();
1862 let old_do_not_suspend = state.do_not_suspend;
1863 state.do_not_suspend = true;
1864
1865 let thread = self.current_guest_thread()?;
1866 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
1867 if thread.old_do_not_suspend.is_some() {
1868 bail_bug!("current thread already has `old_do_not_suspend` value");
1869 }
1870
1871 thread.old_do_not_suspend = Some(old_do_not_suspend);
1872
1873 Ok(())
1874 }
1875
1876 fn exit_sync_call(&mut self, callee: RuntimeInstance) -> Result<()> {
1877 log::trace!("exit sync-typed call {callee:?}");
1878 let thread = self.current_guest_thread()?;
1879 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
1880 let Some(old_do_not_suspend) = thread.old_do_not_suspend.take() else {
1881 bail_bug!("current thread missing `old_do_not_suspend` value");
1882 };
1883 let state = self.instance_state(callee).concurrent_state();
1884 state.do_not_suspend = old_do_not_suspend;
1885 Ok(())
1886 }
1887
1888 pub(crate) fn enter_guest_sync_call(
1900 &mut self,
1901 callee_async_typed: bool,
1902 callee: RuntimeInstance,
1903 ) -> Result<()> {
1904 log::trace!("enter sync-lifted call {callee:?}");
1905 if !self.concurrency_support() {
1906 return self.enter_call_not_concurrent();
1907 }
1908
1909 let thread = self.current_thread()?;
1910 let caller = if let Some(thread) = thread.guest() {
1911 Caller::Guest { thread: *thread }
1912 } else {
1913 Caller::Host {
1914 tx: None,
1915 host_future_present: false,
1916 caller: self.materialize_host_task_id()?,
1917 }
1918 };
1919
1920 let state = self.concurrent_state_mut()?;
1921 let guest_thread = GuestTask::new(
1922 state,
1923 Box::new(move |_, _| bail_bug!("cannot lower params in sync call")),
1924 LiftResult {
1925 lift: Box::new(move |_, _| bail_bug!("cannot lift result in sync call")),
1926 ty: TypeTupleIndex::reserved_value(),
1927 memory: None,
1928 string_encoding: StringEncoding::Utf8,
1929 },
1930 caller,
1931 None,
1932 callee,
1933 callee_async_typed,
1934 true,
1935 )?;
1936
1937 Instance::from_wasmtime(self, callee.instance).add_guest_thread_to_instance_table(
1938 guest_thread.thread,
1939 self,
1940 callee.index,
1941 )?;
1942 self.set_thread(guest_thread)?;
1943
1944 if !callee_async_typed {
1945 self.enter_sync_call(callee)?;
1946 }
1947
1948 Ok(())
1949 }
1950
1951 pub(crate) fn exit_guest_sync_call(&mut self) -> Result<()> {
1959 if !self.concurrency_support() {
1960 return Ok(self.exit_call_not_concurrent());
1961 }
1962
1963 let thread = match self.current_thread()?.guest() {
1964 Some(t) => *t,
1965 None => bail_bug!("expected task when exiting"),
1966 };
1967 let task = self.concurrent_state_mut()?.get_mut(thread.task)?;
1968 let instance = task.instance;
1969
1970 let caller = match &task.caller {
1971 &Caller::Guest { thread } => thread.into(),
1972 &Caller::Host { caller, .. } => caller
1973 .map(CurrentThread::Host)
1974 .unwrap_or(CurrentThread::None),
1975 };
1976 task.lift_result = None;
1977 task.exited = true;
1978 let async_typed = task.async_typed;
1979
1980 if !async_typed {
1981 self.exit_sync_call(instance)?;
1982 }
1983
1984 self.set_thread(caller)?;
1985
1986 log::trace!("exit sync-lifted call {instance:?}");
1987
1988 if async_typed {
1989 self.switch_or_trap_if_may_not_suspend(instance)?;
1994 }
1995
1996 self.cleanup_thread(thread, instance, CleanupTask::Yes)?;
1997
1998 Ok(())
1999 }
2000
2001 pub(crate) fn host_task_create(&mut self) -> Result<EnteredHostTask> {
2008 if !self.concurrency_support() {
2009 self.enter_call_not_concurrent()?;
2010 return Ok(None);
2011 }
2012 let caller = self.current_guest_thread()?;
2013 log::trace!("new deferred host task with caller {caller:?}");
2014
2015 self.set_thread(CurrentThread::DeferredHost(caller))?;
2016 let state = self.concurrent_state_mut()?;
2017 debug_assert!(state.deferred_host_call_context.is_none());
2018 state.deferred_host_call_context = Some(CallContext::default());
2019 state.debug_assert_deferred_host_invariant();
2020 Ok(Some(caller))
2021 }
2022
2023 pub(crate) fn host_task_delete(
2030 &mut self,
2031 original_task: EnteredHostTask,
2032 materialized_task: Option<TableId<HostTask>>,
2033 ) -> Result<()> {
2034 match original_task {
2035 Some(caller) => {
2036 self.set_thread(caller)?;
2037 if materialized_task.is_none() {
2038 let state = self.concurrent_state_mut()?;
2039 let context = state
2040 .deferred_host_call_context
2041 .take()
2042 .expect("deferred host call context should be present");
2043 debug_assert!(context.is_empty());
2044 state.debug_assert_deferred_host_invariant();
2045 }
2046 log::trace!(
2047 "delete host task with caller {original_task:?} and materialized as {materialized_task:?}"
2048 );
2049 if let Some(task) = materialized_task {
2050 Waitable::Host(task).delete_from(self)?;
2051 }
2052 }
2053 None => {
2054 debug_assert!(materialized_task.is_none());
2055 self.exit_call_not_concurrent();
2056 }
2057 }
2058 Ok(())
2059 }
2060
2061 fn instance_state(&mut self, instance: RuntimeInstance) -> &mut InstanceState {
2064 self.component_instance_mut(instance.instance)
2065 .instance_state(instance.index)
2066 }
2067
2068 pub(crate) fn set_thread(&mut self, thread: impl Into<CurrentThread>) -> Result<CurrentThread> {
2074 let thread = thread.into();
2075 let state = self.concurrent_state_mut()?;
2076 state.debug_assert_deferred_host_invariant();
2077 let old_thread = mem::replace(&mut state.unforced_current_thread, thread);
2078
2079 state.handle_thread_switch(old_thread, thread)?;
2080
2081 if let Some(old_thread) = old_thread.guest() {
2089 let old_context = *self.vm_store_context_mut().component_context_mut();
2090 self.concurrent_state_mut()?
2091 .get_mut(old_thread.thread)?
2092 .context = old_context;
2093 }
2094 if cfg!(debug_assertions) {
2095 *self.vm_store_context_mut().component_context_mut() =
2096 [u32::MAX; NUM_COMPONENT_CONTEXT_SLOTS];
2097 }
2098 if let Some(thread) = thread.guest() {
2099 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
2100 let context = thread.context;
2101 if cfg!(debug_assertions) {
2102 thread.context = [u32::MAX; NUM_COMPONENT_CONTEXT_SLOTS];
2103 }
2104 *self.vm_store_context_mut().component_context_mut() = context;
2105 }
2106
2107 *self.vm_store_context_mut().current_thread_mut() = if thread.is_none() {
2109 VMLazyThread::none()
2110 } else {
2111 VMLazyThread::forced()
2112 };
2113
2114 Ok(old_thread)
2115 }
2116
2117 fn switch_or_trap_if_may_not_suspend(&mut self, instance: RuntimeInstance) -> Result<()> {
2119 if self.switch_if_may_not_suspend(instance)? {
2120 Ok(())
2121 } else {
2122 Err(Trap::CannotBlockSyncTask.into())
2123 }
2124 }
2125
2126 fn switch_if_may_not_suspend(&mut self, instance: RuntimeInstance) -> Result<bool> {
2130 self.concurrent_state_mut()?;
2134
2135 Ok(!self.concurrency_support()
2136 || !self
2137 .instance_state(instance)
2138 .concurrent_state()
2139 .do_not_suspend
2140 || self
2141 .concurrent_state_mut()?
2142 .promote_instance_local_thread_work_item(instance)?)
2143 }
2144
2145 fn enter_instance(&mut self, instance: RuntimeInstance) {
2149 log::trace!("enter {instance:?}");
2150 self.instance_state(instance)
2151 .concurrent_state()
2152 .do_not_enter = true;
2153 }
2154
2155 fn exit_instance(&mut self, instance: RuntimeInstance) -> Result<()> {
2159 log::trace!("exit {instance:?}");
2160 self.instance_state(instance)
2161 .concurrent_state()
2162 .do_not_enter = false;
2163 self.partition_pending(instance)
2164 }
2165
2166 fn partition_pending(&mut self, instance: RuntimeInstance) -> Result<()> {
2174 for (thread, kind) in
2175 mem::take(&mut self.instance_state(instance).concurrent_state().pending).into_iter()
2176 {
2177 let call = GuestCall { thread, kind };
2178 if call.is_ready(self)? {
2179 self.concurrent_state_mut()?
2180 .push_high_priority(WorkItem::GuestCall { instance, call });
2181 } else {
2182 self.instance_state(instance)
2183 .concurrent_state()
2184 .pending
2185 .insert(call.thread, call.kind);
2186 }
2187 }
2188
2189 if let Some(waker) = self
2190 .concurrent_state_mut()?
2191 .ready_for_concurrent_call_waker
2192 .take()
2193 {
2194 waker.wake();
2195 }
2196
2197 Ok(())
2198 }
2199
2200 pub(crate) fn backpressure_modify(
2202 &mut self,
2203 caller_instance: RuntimeInstance,
2204 modify: impl FnOnce(u16) -> Option<u16>,
2205 ) -> Result<()> {
2206 let state = self.instance_state(caller_instance).concurrent_state();
2207 let old = state.backpressure;
2208 let new = modify(old).ok_or_else(|| Trap::BackpressureOverflow)?;
2209 state.backpressure = new;
2210
2211 if old > 0 && new == 0 {
2212 self.partition_pending(caller_instance)?;
2215 }
2216
2217 Ok(())
2218 }
2219
2220 async fn resume_fiber(&mut self, fiber: StoreFiber<'static>) -> Result<()> {
2223 let old_thread = self.current_thread()?;
2224 log::trace!("resume_fiber: save current thread {old_thread:?}");
2225
2226 let fiber = fiber::resolve_or_release(self, fiber).await?;
2227
2228 self.set_thread(old_thread)?;
2229
2230 let state = self.concurrent_state_mut()?;
2231
2232 if let Some(ot) = old_thread.guest() {
2233 state.get_mut(ot.thread)?.state = GuestThreadState::Running;
2234 }
2235 log::trace!("resume_fiber: restore current thread {old_thread:?}");
2236
2237 if let Some(mut fiber) = fiber {
2238 log::trace!("resume_fiber: suspend reason {:?}", &state.suspend_reason);
2239 let reason = match state.suspend_reason.take() {
2241 Some(r) => r,
2242 None => bail_bug!("suspend reason missing when resuming fiber"),
2243 };
2244 match reason {
2245 SuspendReason::NeedWork => {
2246 if state.worker.is_none() {
2247 state.worker = Some(fiber);
2248 } else {
2249 fiber.dispose(self);
2250 }
2251 }
2252 SuspendReason::Yielding { thread } => {
2253 state.get_mut(thread.thread)?.state = GuestThreadState::Ready { fiber };
2254 let instance = state.get_mut(thread.task)?.instance;
2255 state.push_low_priority(WorkItem::ResumeThread { instance, thread });
2256 }
2257 SuspendReason::ExplicitlySuspending { thread } => {
2258 state.get_mut(thread.thread)?.state = GuestThreadState::Suspended(fiber);
2259 }
2260 SuspendReason::Waiting { set, thread } => {
2261 let old = state
2262 .get_mut(set)?
2263 .waiting
2264 .insert(thread, WaitMode::Fiber(fiber));
2265 assert!(old.is_none());
2266 }
2267 SuspendReason::YieldingToSubtask { thread } => {
2268 let item = WorkItem::ResumeFiber {
2277 instance: state.get_mut(thread.task)?.instance,
2278 thread,
2279 fiber,
2280 };
2281
2282 if state.next_switch_item.replace(item).is_some() {
2283 bail_bug!(
2286 "`ConcurrentState::next_switch_item` was already `Some(_)` when \
2287 a thread wanted to wait on a subtask"
2288 );
2289 }
2290 }
2291 };
2292 } else {
2293 log::trace!("resume_fiber: fiber has exited");
2294 }
2295
2296 Ok(())
2297 }
2298
2299 fn suspend(&mut self, reason: SuspendReason) -> Result<()> {
2305 log::trace!("suspend fiber: {reason:?}");
2306
2307 let state = self.concurrent_state_mut()?;
2308
2309 let (save_and_restore_thread, save_and_restore_next_switch_item) = match &reason {
2316 SuspendReason::Yielding { .. }
2317 | SuspendReason::Waiting { .. }
2318 | SuspendReason::ExplicitlySuspending { .. } => {
2319 if state.switch_item.is_none() {
2322 state.take_next_switch_item()?;
2323 }
2324
2325 (true, false)
2326 }
2327 SuspendReason::YieldingToSubtask { .. } => (true, true),
2328 SuspendReason::NeedWork => (false, false),
2329 };
2330
2331 let old_next_switch_item = if save_and_restore_next_switch_item {
2332 let item = state.next_switch_item.take();
2333 Some(state.push(item)?)
2337 } else {
2338 None
2339 };
2340
2341 let old_guest_thread = if save_and_restore_thread {
2342 self.current_thread()?
2343 } else {
2344 CurrentThread::None
2345 };
2346
2347 let suspend_reason = &mut self.concurrent_state_mut()?.suspend_reason;
2348 assert!(suspend_reason.is_none());
2349 *suspend_reason = Some(reason);
2350
2351 if !self.fiber_async_state_mut().can_block() {
2354 return Err(format_err!("future dropped"));
2355 }
2356
2357 self.with_blocking(|_, cx| cx.suspend(StoreFiberYield::ReleaseStore))?;
2358
2359 if save_and_restore_thread {
2360 self.set_thread(old_guest_thread)?;
2361 }
2362
2363 if let Some(item) = old_next_switch_item {
2364 let state = self.concurrent_state_mut()?;
2365 state.next_switch_item = state.delete(item)?;
2366 }
2367
2368 Ok(())
2369 }
2370
2371 fn wait_for_event(
2372 &mut self,
2373 caller_instance: RuntimeInstance,
2374 waitable: Waitable,
2375 ) -> Result<()> {
2376 let caller = self.current_guest_thread()?;
2377 let state = self.concurrent_state_mut()?;
2378
2379 waitable.trap_if_in_waitable_set(state)?;
2380
2381 let set = state.get_mut(caller.thread)?.sync_call_set;
2382 waitable.join(state, Some(set))?;
2383
2384 self.switch_or_trap_if_may_not_suspend(caller_instance)?;
2385
2386 self.suspend(SuspendReason::Waiting {
2387 set,
2388 thread: caller,
2389 })?;
2390 let state = self.concurrent_state_mut()?;
2391
2392 waitable.join(state, None)
2393 }
2394
2395 fn cleanup_thread(
2417 &mut self,
2418 guest_thread: QualifiedThreadId,
2419 runtime_instance: RuntimeInstance,
2420 cleanup_task: CleanupTask,
2421 ) -> Result<()> {
2422 let state = self.concurrent_state_mut()?;
2423 state.take_next_switch_item()?;
2426 let thread_data = state.get_mut(guest_thread.thread)?;
2427 let sync_call_set = thread_data.sync_call_set;
2428 if let Some(guest_id) = thread_data.instance_rep {
2429 self.instance_state(runtime_instance)
2430 .thread_handle_table()
2431 .guest_thread_remove(guest_id)?;
2432 }
2433 let state = self.concurrent_state_mut()?;
2434
2435 for waitable in mem::take(&mut state.get_mut(sync_call_set)?.ready) {
2437 if let Some(Event::Subtask {
2438 status: Status::Returned | Status::ReturnCancelled,
2439 }) = waitable.common(self.concurrent_state_mut()?)?.event
2440 {
2441 waitable.delete_from(self)?;
2442 }
2443 }
2444
2445 let state = self.concurrent_state_mut()?;
2446 state.delete(guest_thread.thread)?;
2447 state.delete(sync_call_set)?;
2448 let task = state.get_mut(guest_thread.task)?;
2449 task.threads.remove(&guest_thread.thread);
2450
2451 if task.threads.is_empty() && !task.returned_or_cancelled() {
2452 bail!(Trap::NoAsyncResult);
2453 }
2454 let ready_to_delete = task.ready_to_delete();
2455
2456 if !task.decremented_interesting_task_count && task.exited && task.returned_or_cancelled() {
2457 task.decremented_interesting_task_count = true;
2458
2459 debug_assert!(state.interesting_tasks > 0);
2460 state.interesting_tasks -= 1;
2461 if state.interesting_tasks == 0
2462 && let Some(waker) = state.interesting_tasks_empty_waker.take()
2463 {
2464 waker.wake();
2465 }
2466 }
2467
2468 match cleanup_task {
2469 CleanupTask::Yes => {
2470 if ready_to_delete {
2471 Waitable::Guest(guest_thread.task).delete_from(self)?;
2472 }
2473 }
2474 CleanupTask::No => {}
2475 }
2476
2477 Ok(())
2478 }
2479
2480 fn cancel_guest_subtask_without_lowered_parameters(
2493 &mut self,
2494 caller_instance: RuntimeInstance,
2495 guest_task: TableId<GuestTask>,
2496 ) -> Result<()> {
2497 let concurrent_state = self.concurrent_state_mut()?;
2498 let task = concurrent_state.get_mut(guest_task)?;
2499 assert!(!task.already_lowered_parameters());
2500 task.lower_params = None;
2504 task.lift_result = None;
2505 task.exited = true;
2506 let instance = task.instance;
2507
2508 assert_eq!(1, task.threads.len());
2511 let thread = *task.threads.iter().next().unwrap();
2512 self.cleanup_thread(
2513 QualifiedThreadId {
2514 task: guest_task,
2515 thread,
2516 },
2517 caller_instance,
2518 CleanupTask::No,
2519 )?;
2520
2521 let pending = &mut self.instance_state(instance).concurrent_state().pending;
2523 let pending_count = pending.len();
2524 pending.retain(|thread, _| thread.task != guest_task);
2525 if pending.len() == pending_count {
2527 bail!(Trap::SubtaskCancelAfterTerminal);
2528 }
2529 Ok(())
2530 }
2531
2532 pub(crate) fn current_scope(&mut self) -> Result<Option<CurrentScope>> {
2535 if !self.concurrency_support() {
2536 return Ok(self
2537 .current_scope_id_not_concurrent()?
2538 .map(|id| CurrentScope::Id(Scope::Id(id))));
2539 }
2540
2541 Ok(match self.current_thread()? {
2542 CurrentThread::Guest(id) => Some(CurrentScope::Id(Scope::Id(id.task.rep()))),
2543 CurrentThread::Host(id) => Some(CurrentScope::Id(Scope::HostId(id.rep()))),
2544 CurrentThread::DeferredHost(_) => Some(CurrentScope::DeferredHost),
2545 CurrentThread::None => return Ok(None),
2546 })
2547 }
2548
2549 pub(crate) fn queue_task(
2550 &mut self,
2551 task: impl FnOnce(&mut dyn VMStore) -> Result<()> + Send + 'static,
2552 ) -> Result<()> {
2553 self.concurrent_state_mut()?
2554 .push_high_priority(WorkItem::WorkerFunction(AlwaysMut::new(Box::new(task))));
2555 Ok(())
2556 }
2557
2558 fn any_may_not_suspend(&mut self) -> Result<bool> {
2567 Ok(self
2575 .concurrent_state_mut()?
2576 .table
2577 .get_mut()
2578 .iter_mut()
2579 .filter_map(|(_, entry)| {
2580 if let Some(task) = entry.downcast_ref::<GuestTask>() {
2581 Some(task.instance)
2582 } else {
2583 None
2584 }
2585 })
2586 .collect::<Vec<_>>()
2587 .into_iter()
2588 .any(|instance| {
2589 self.instance_state(instance)
2590 .concurrent_state()
2591 .do_not_suspend
2592 }))
2593 }
2594}
2595
2596enum CleanupTask {
2597 Yes,
2598 No,
2599}
2600
2601impl Instance {
2602 fn get_event(
2605 self,
2606 store: &mut StoreOpaque,
2607 guest_task: TableId<GuestTask>,
2608 set: Option<TableId<WaitableSet>>,
2609 cancellable: bool,
2610 ) -> Result<Option<(Event, Option<(Waitable, u32)>)>> {
2611 let state = store.concurrent_state_mut()?;
2612
2613 let task = state.get_mut(guest_task)?;
2614 let event = &mut task.event;
2615 if let Some(ev) = event
2616 && (cancellable || !matches!(ev, Event::Cancelled))
2617 {
2618 log::trace!("deliver event {ev:?} to {guest_task:?}");
2619
2620 if matches!(ev, Event::Cancelled) {
2621 task.cancel_request_delivered = true;
2622 }
2623
2624 let ev = *ev;
2625 *event = None;
2626 return Ok(Some((ev, None)));
2627 }
2628
2629 let set = match set {
2630 Some(set) => set,
2631 None => return Ok(None),
2632 };
2633 let waitable = match state.get_mut(set)?.ready.pop_first() {
2634 Some(v) => v,
2635 None => return Ok(None),
2636 };
2637
2638 let common = waitable.common(state)?;
2639 let handle = match common.handle {
2640 Some(h) => h,
2641 None => bail_bug!("handle not set when delivering event"),
2642 };
2643 let event = match common.event.take() {
2644 Some(e) => e,
2645 None => bail_bug!("event not set when delivering event"),
2646 };
2647
2648 log::trace!(
2649 "deliver event {event:?} to {guest_task:?} for {waitable:?} (handle {handle}); set {set:?}"
2650 );
2651
2652 waitable.on_delivery(store, self, event)?;
2653
2654 Ok(Some((event, Some((waitable, handle)))))
2655 }
2656
2657 fn handle_callback_code(
2663 self,
2664 store: &mut StoreOpaque,
2665 guest_thread: QualifiedThreadId,
2666 runtime_instance: RuntimeComponentInstanceIndex,
2667 code: u32,
2668 ) -> Result<()> {
2669 let (code, set) = unpack_callback_code(code);
2670
2671 log::trace!("received callback code from {guest_thread:?}: {code} (set: {set})");
2672
2673 let state = store.concurrent_state_mut()?;
2674
2675 state.take_next_switch_item()?;
2676
2677 let get_set = |store: &mut StoreOpaque, handle| -> Result<_> {
2678 let set = store
2679 .instance_state(self.runtime_instance(runtime_instance))
2680 .handle_table()
2681 .waitable_set_rep(handle)?;
2682
2683 Ok(TableId::<WaitableSet>::new(set))
2684 };
2685
2686 match code {
2687 callback_code::EXIT => {
2688 log::trace!("implicit thread {guest_thread:?} completed");
2689 let task = store.concurrent_state_mut()?.get_mut(guest_thread.task)?;
2690 task.exited = true;
2691 task.callback = None;
2692
2693 let runtime_instance = self.runtime_instance(runtime_instance);
2694
2695 store.switch_or_trap_if_may_not_suspend(runtime_instance)?;
2700
2701 store.cleanup_thread(guest_thread, runtime_instance, CleanupTask::Yes)?;
2702 }
2703 callback_code::YIELD => {
2704 let old = state
2707 .get_mut(guest_thread.thread)?
2708 .wake_on_cancel
2709 .replace(WakeOnCancel::Yielding);
2710 if !old.is_none() {
2711 bail_bug!("thread unexpectedly had wake_on_cancel set");
2712 }
2713
2714 let task = state.get_mut(guest_thread.task)?;
2715 if let Some(event) = task.event {
2720 assert!(matches!(event, Event::None | Event::Cancelled));
2721 } else {
2722 task.event = Some(Event::None);
2723 }
2724 let call = GuestCall {
2725 thread: guest_thread,
2726 kind: GuestCallKind::DeliverEvent {
2727 instance: self,
2728 set: None,
2729 },
2730 };
2731 state.push_low_priority(WorkItem::GuestCall {
2734 instance: self.runtime_instance(runtime_instance),
2735 call,
2736 });
2737 }
2738 callback_code::WAIT => {
2739 let set = get_set(store, set)?;
2740 let state = store.concurrent_state_mut()?;
2741
2742 if state.get_mut(guest_thread.task)?.event.is_some()
2743 || !state.get_mut(set)?.ready.is_empty()
2744 {
2745 state.push_high_priority(WorkItem::GuestCall {
2747 instance: self.runtime_instance(runtime_instance),
2748 call: GuestCall {
2749 thread: guest_thread,
2750 kind: GuestCallKind::DeliverEvent {
2751 instance: self,
2752 set: Some(set),
2753 },
2754 },
2755 });
2756 } else {
2757 let old = state
2765 .get_mut(guest_thread.thread)?
2766 .wake_on_cancel
2767 .replace(WakeOnCancel::Waiting(set));
2768 if !old.is_none() {
2769 bail_bug!("thread unexpectedly had wake_on_cancel set");
2770 }
2771 let old = state
2772 .get_mut(set)?
2773 .waiting
2774 .insert(guest_thread, WaitMode::Callback(self));
2775 if !old.is_none() {
2776 bail_bug!("set's waiting set already had this thread registered");
2777 }
2778 }
2779 }
2780 _ => bail!(Trap::UnsupportedCallbackCode),
2781 }
2782
2783 Ok(())
2784 }
2785
2786 unsafe fn stage_call<T: 'static>(
2793 self,
2794 mut store: StoreContextMut<T>,
2795 guest_thread: QualifiedThreadId,
2796 callee: SendSyncPtr<VMFuncRef>,
2797 param_count: usize,
2798 result_count: usize,
2799 async_: bool,
2800 callback: Option<SendSyncPtr<VMFuncRef>>,
2801 post_return: Option<SendSyncPtr<VMFuncRef>>,
2802 host_caller: bool,
2803 ) -> Result<()> {
2804 unsafe fn make_call<T: 'static>(
2819 store: StoreContextMut<T>,
2820 guest_thread: QualifiedThreadId,
2821 callee: SendSyncPtr<VMFuncRef>,
2822 param_count: usize,
2823 result_count: usize,
2824 ) -> impl FnOnce(&mut dyn VMStore) -> Result<[MaybeUninit<ValRaw>; MAX_FLAT_PARAMS]>
2825 + Send
2826 + Sync
2827 + 'static
2828 + use<T> {
2829 let token = StoreToken::new(store);
2830 move |store: &mut dyn VMStore| {
2831 let mut storage = [MaybeUninit::uninit(); MAX_FLAT_PARAMS];
2832
2833 store
2834 .concurrent_state_mut()?
2835 .get_mut(guest_thread.thread)?
2836 .state = GuestThreadState::Running;
2837 let task = store.concurrent_state_mut()?.get_mut(guest_thread.task)?;
2838 let lower = match task.lower_params.take() {
2839 Some(l) => l,
2840 None => bail_bug!("lower_params missing"),
2841 };
2842
2843 lower(store, &mut storage[..param_count])?;
2844
2845 let mut store = token.as_context_mut(store);
2846
2847 unsafe {
2850 crate::Func::call_unchecked_raw(
2851 &mut store,
2852 callee.as_non_null(),
2853 NonNull::new(
2854 &mut storage[..param_count.max(result_count)]
2855 as *mut [MaybeUninit<ValRaw>] as _,
2856 )
2857 .unwrap(),
2858 )?;
2859 }
2860
2861 Ok(storage)
2862 }
2863 }
2864
2865 let call = unsafe {
2869 make_call(
2870 store.as_context_mut(),
2871 guest_thread,
2872 callee,
2873 param_count,
2874 result_count,
2875 )
2876 };
2877
2878 let callee_instance = store
2879 .0
2880 .concurrent_state_mut()?
2881 .get_mut(guest_thread.task)?
2882 .instance;
2883
2884 let fun = if callback.is_some() {
2885 assert!(async_);
2886
2887 Box::new(move |store: &mut dyn VMStore| {
2888 self.add_guest_thread_to_instance_table(
2889 guest_thread.thread,
2890 store,
2891 callee_instance.index,
2892 )?;
2893 let old_thread = store.set_thread(guest_thread)?;
2894 log::trace!(
2895 "stackless call: replaced {old_thread:?} with {guest_thread:?} as current thread"
2896 );
2897
2898 store.enter_instance(callee_instance);
2899
2900 let storage = call(store)?;
2907
2908 store.exit_instance(callee_instance)?;
2909
2910 store.set_thread(old_thread)?;
2911 let state = store.concurrent_state_mut()?;
2912 if let Some(t) = old_thread.guest() {
2913 state.get_mut(t.thread)?.state = GuestThreadState::Running;
2914 }
2915 log::trace!("stackless call: restored {old_thread:?} as current thread");
2916
2917 let code = unsafe { storage[0].assume_init() }.get_i32() as u32;
2920
2921 self.handle_callback_code(store, guest_thread, callee_instance.index, code)
2922 }) as Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send + Sync>
2923 } else {
2924 let token = StoreToken::new(store.as_context_mut());
2925 Box::new(move |store: &mut dyn VMStore| {
2926 self.add_guest_thread_to_instance_table(
2927 guest_thread.thread,
2928 store,
2929 callee_instance.index,
2930 )?;
2931 let old_thread = store.set_thread(guest_thread)?;
2932 log::trace!(
2933 "sync/async-stackful call: replaced {old_thread:?} with {guest_thread:?} as current thread",
2934 );
2935 let flags = self.id().get(store).instance_flags(callee_instance.index);
2936
2937 let callee_async_typed = store
2938 .concurrent_state_mut()?
2939 .get_mut(guest_thread.task)?
2940 .async_typed;
2941
2942 if !async_ && callee_async_typed {
2946 store.enter_instance(callee_instance);
2947 }
2948
2949 if !callee_async_typed {
2950 store.enter_sync_call(callee_instance)?;
2951 }
2952
2953 let storage = call(store)?;
2960
2961 if !callee_async_typed {
2962 store.exit_sync_call(callee_instance)?;
2963 }
2964
2965 if !async_ {
2966 if callee_async_typed {
2972 store.exit_instance(callee_instance)?;
2973 }
2974
2975 let lift = {
2976 let state = store.concurrent_state_mut()?;
2977 if !state.get_mut(guest_thread.task)?.result.is_none() {
2978 bail_bug!("task has already produced a result");
2979 }
2980
2981 match state.get_mut(guest_thread.task)?.lift_result.take() {
2982 Some(lift) => lift,
2983 None => bail_bug!("lift_result field is missing"),
2984 }
2985 };
2986
2987 let result = (lift.lift)(store, unsafe {
2990 mem::transmute::<&[MaybeUninit<ValRaw>], &[ValRaw]>(
2991 &storage[..result_count],
2992 )
2993 })?;
2994
2995 let post_return_arg = match result_count {
2996 0 => ValRaw::i32(0),
2997 1 => unsafe { storage[0].assume_init() },
3000 _ => unreachable!(),
3001 };
3002
3003 unsafe {
3004 call_post_return(
3005 token.as_context_mut(store),
3006 post_return.map(|v| v.as_non_null()),
3007 post_return_arg,
3008 flags,
3009 )?;
3010 }
3011
3012 self.task_complete(store, guest_thread.task, result, Status::Returned)?;
3013 }
3014
3015 store.set_thread(old_thread)?;
3016
3017 store
3018 .concurrent_state_mut()?
3019 .get_mut(guest_thread.task)?
3020 .exited = true;
3021
3022 log::trace!(
3023 "clean up thread; async lifted? {async_} async typed? {callee_async_typed}"
3024 );
3025
3026 if callee_async_typed {
3027 store.switch_or_trap_if_may_not_suspend(callee_instance)?;
3032 }
3033
3034 store.cleanup_thread(guest_thread, callee_instance, CleanupTask::Yes)?;
3036 Ok(())
3037 })
3038 };
3039
3040 store.0.concurrent_state_mut()?.push_work_item(
3041 WorkItem::GuestCall {
3042 instance: callee_instance,
3043 call: GuestCall {
3044 thread: guest_thread,
3045 kind: GuestCallKind::StartImplicit(fun),
3046 },
3047 },
3048 if host_caller {
3049 Priority::High
3050 } else {
3051 Priority::Switch
3052 },
3053 )?;
3054
3055 Ok(())
3056 }
3057
3058 unsafe fn prepare_call<T: 'static>(
3071 self,
3072 mut store: StoreContextMut<T>,
3073 start: NonNull<VMFuncRef>,
3074 return_: NonNull<VMFuncRef>,
3075 caller_instance: RuntimeComponentInstanceIndex,
3076 callee_instance: RuntimeComponentInstanceIndex,
3077 task_return_type: TypeTupleIndex,
3078 callee_async_typed: bool,
3079 memory: *mut VMMemoryDefinition,
3080 string_encoding: StringEncoding,
3081 caller_info: CallerInfo,
3082 ) -> Result<()> {
3083 enum ResultInfo {
3084 Heap { results: u32 },
3085 Stack { result_count: u32 },
3086 }
3087
3088 let result_info = match &caller_info {
3089 CallerInfo::Async {
3090 has_result: true,
3091 params,
3092 } => ResultInfo::Heap {
3093 results: match params.last() {
3094 Some(r) => r.get_u32(),
3095 None => bail_bug!("retptr missing"),
3096 },
3097 },
3098 CallerInfo::Async {
3099 has_result: false, ..
3100 } => ResultInfo::Stack { result_count: 0 },
3101 CallerInfo::Sync {
3102 result_count,
3103 params,
3104 } if *result_count > u32::try_from(MAX_FLAT_RESULTS)? => ResultInfo::Heap {
3105 results: match params.last() {
3106 Some(r) => r.get_u32(),
3107 None => bail_bug!("arg ptr missing"),
3108 },
3109 },
3110 CallerInfo::Sync { result_count, .. } => ResultInfo::Stack {
3111 result_count: *result_count,
3112 },
3113 };
3114
3115 let sync_caller = matches!(caller_info, CallerInfo::Sync { .. });
3116
3117 let start = SendSyncPtr::new(start);
3121 let return_ = SendSyncPtr::new(return_);
3122 let token = StoreToken::new(store.as_context_mut());
3123 let old_thread = store.0.current_guest_thread()?;
3124
3125 let state = store.0.concurrent_state_mut()?;
3126
3127 debug_assert_eq!(
3128 state.get_mut(old_thread.task)?.instance,
3129 self.runtime_instance(caller_instance)
3130 );
3131
3132 let guest_thread = GuestTask::new(
3133 state,
3134 Box::new(move |store, dst| {
3135 let mut store = token.as_context_mut(store);
3136 assert!(dst.len() <= MAX_FLAT_PARAMS);
3137 let mut src = [MaybeUninit::uninit(); MAX_FLAT_PARAMS + 1];
3139 let count = match caller_info {
3140 CallerInfo::Async { params, has_result } => {
3144 let params = ¶ms[..params.len() - usize::from(has_result)];
3145 for (param, src) in params.iter().zip(&mut src) {
3146 src.write(*param);
3147 }
3148 params.len()
3149 }
3150
3151 CallerInfo::Sync { params, .. } => {
3153 for (param, src) in params.iter().zip(&mut src) {
3154 src.write(*param);
3155 }
3156 params.len()
3157 }
3158 };
3159 unsafe {
3166 crate::Func::call_unchecked_raw(
3167 &mut store,
3168 start.as_non_null(),
3169 NonNull::new(
3170 &mut src[..count.max(dst.len())] as *mut [MaybeUninit<ValRaw>] as _,
3171 )
3172 .unwrap(),
3173 )?;
3174 }
3175 dst.copy_from_slice(&src[..dst.len()]);
3176 let task = store.0.current_guest_thread()?.task;
3177 let state = store.0.concurrent_state_mut()?;
3178 Waitable::Guest(task).set_event(
3179 state,
3180 Some(Event::Subtask {
3181 status: Status::Started,
3182 }),
3183 )?;
3184 Ok(())
3185 }),
3186 LiftResult {
3187 lift: Box::new(move |store, src| {
3188 let mut store = token.as_context_mut(store);
3191 let mut my_src = src.to_owned(); if let ResultInfo::Heap { results } = &result_info {
3193 my_src.push(ValRaw::u32(*results));
3194 }
3195
3196 unsafe {
3203 crate::Func::call_unchecked_raw(
3204 &mut store,
3205 return_.as_non_null(),
3206 my_src.as_mut_slice().into(),
3207 )?;
3208 }
3209
3210 let thread = store.0.current_guest_thread()?;
3211 let state = store.0.concurrent_state_mut()?;
3212 if sync_caller {
3213 state.get_mut(thread.task)?.sync_result = SyncResult::Produced(
3214 if let ResultInfo::Stack { result_count } = &result_info {
3215 match result_count {
3216 0 => None,
3217 1 => Some(my_src[0]),
3218 _ => unreachable!(),
3219 }
3220 } else {
3221 None
3222 },
3223 );
3224 }
3225 Ok(Box::new(DummyResult) as Box<dyn Any + Send + Sync>)
3226 }),
3227 ty: task_return_type,
3228 memory: NonNull::new(memory).map(SendSyncPtr::new),
3229 string_encoding,
3230 },
3231 Caller::Guest { thread: old_thread },
3232 None,
3233 self.runtime_instance(callee_instance),
3234 callee_async_typed,
3235 false,
3238 )?;
3239
3240 store.0.set_thread(guest_thread)?;
3243 log::trace!("pushed {guest_thread:?} as current thread; old thread was {old_thread:?}");
3244
3245 Ok(())
3246 }
3247
3248 unsafe fn call_callback<T>(
3253 self,
3254 mut store: StoreContextMut<T>,
3255 function: SendSyncPtr<VMFuncRef>,
3256 event: Event,
3257 handle: u32,
3258 ) -> Result<u32> {
3259 let (ordinal, result) = event.parts();
3260 let params = &mut [
3261 ValRaw::u32(ordinal),
3262 ValRaw::u32(handle),
3263 ValRaw::u32(result),
3264 ];
3265 unsafe {
3270 crate::Func::call_unchecked_raw(
3271 &mut store,
3272 function.as_non_null(),
3273 params.as_mut_slice().into(),
3274 )?;
3275 }
3276 Ok(params[0].get_u32())
3277 }
3278
3279 unsafe fn start_call<T: 'static>(
3292 self,
3293 mut store: StoreContextMut<T>,
3294 callback: *mut VMFuncRef,
3295 post_return: *mut VMFuncRef,
3296 callee: NonNull<VMFuncRef>,
3297 param_count: u32,
3298 result_count: u32,
3299 flags: u32,
3300 storage: Option<&mut [MaybeUninit<ValRaw>]>,
3301 ) -> Result<u32> {
3302 let token = StoreToken::new(store.as_context_mut());
3303 let async_caller = storage.is_none();
3304 let guest_thread = store.0.current_guest_thread()?;
3305 let state = store.0.concurrent_state_mut()?;
3306
3307 if !state.event_loop_running {
3308 bail_bug!("Instance::start_call called without a running event loop");
3309 }
3310
3311 let callee = SendSyncPtr::new(callee);
3312 let param_count = usize::try_from(param_count)?;
3313 assert!(param_count <= MAX_FLAT_PARAMS);
3314 let result_count = usize::try_from(result_count)?;
3315 assert!(result_count <= MAX_FLAT_RESULTS);
3316
3317 let task = state.get_mut(guest_thread.task)?;
3318 let callee_async_typed = task.async_typed;
3319 let callee_instance = task.instance;
3320
3321 task.async_lifted = (flags & START_FLAG_ASYNC_CALLEE) != 0;
3322
3323 if let Some(callback) = NonNull::new(callback) {
3324 let callback = SendSyncPtr::new(callback);
3328 task.callback = Some(Box::new(move |store, event, handle| {
3329 let store = token.as_context_mut(store);
3330 unsafe { self.call_callback::<T>(store, callback, event, handle) }
3331 }));
3332 }
3333
3334 let Caller::Guest { thread: caller } = &task.caller else {
3335 bail_bug!("start_call unexpectedly invoked for host->guest call");
3338 };
3339 let caller = *caller;
3340 let caller_instance = state.get_mut(caller.task)?.instance;
3341
3342 unsafe {
3344 self.stage_call(
3345 store.as_context_mut(),
3346 guest_thread,
3347 callee,
3348 param_count,
3349 result_count,
3350 (flags & START_FLAG_ASYNC_CALLEE) != 0,
3351 NonNull::new(callback).map(SendSyncPtr::new),
3352 NonNull::new(post_return).map(SendSyncPtr::new),
3353 false,
3354 )?;
3355 }
3356
3357 let old_do_not_suspend = if callee_async_typed {
3358 let state = store.0.instance_state(callee_instance).concurrent_state();
3365 let old_do_not_suspend = state.do_not_suspend;
3366 state.do_not_suspend = false;
3367 Some(old_do_not_suspend)
3368 } else {
3369 None
3370 };
3371
3372 let state = store.0.concurrent_state_mut()?;
3373
3374 let guest_waitable = Waitable::Guest(guest_thread.task);
3377 let old_set = guest_waitable.common(state)?.set;
3378 let set = state.get_mut(caller.thread)?.sync_call_set;
3379 guest_waitable.join(state, Some(set))?;
3380
3381 store.0.set_thread(CurrentThread::None)?;
3382
3383 let mut yielded = false;
3399 let (status, waitable) = loop {
3400 store.0.suspend(if yielded {
3401 SuspendReason::Waiting {
3402 set,
3403 thread: caller,
3404 }
3405 } else {
3406 yielded = true;
3407 SuspendReason::YieldingToSubtask { thread: caller }
3408 })?;
3409
3410 if let Some(old_do_not_suspend) = old_do_not_suspend {
3411 store
3412 .0
3413 .instance_state(callee_instance)
3414 .concurrent_state()
3415 .do_not_suspend = old_do_not_suspend;
3416 }
3417
3418 let state = store.0.concurrent_state_mut()?;
3419
3420 log::trace!("taking event for {:?}", guest_thread.task);
3421 let event = guest_waitable.take_event(state)?;
3422 let Some(Event::Subtask { status }) = event else {
3423 bail_bug!("subtasks should only get subtask events, got {event:?}")
3424 };
3425
3426 log::trace!("status {status:?} for {:?}", guest_thread.task);
3427
3428 if status == Status::Returned {
3429 break (status, None);
3431 } else if async_caller {
3432 let handle = store
3436 .0
3437 .instance_state(caller_instance)
3438 .handle_table()
3439 .subtask_insert_guest(guest_thread.task.rep())?;
3440 store
3441 .0
3442 .concurrent_state_mut()?
3443 .get_mut(guest_thread.task)?
3444 .common
3445 .handle = Some(handle);
3446 break (status, Some(handle));
3447 } else {
3448 store.0.switch_or_trap_if_may_not_suspend(caller_instance)?;
3452 }
3453 };
3454
3455 guest_waitable.join(store.0.concurrent_state_mut()?, old_set)?;
3456
3457 store.0.set_thread(caller)?;
3459 store
3460 .0
3461 .concurrent_state_mut()?
3462 .get_mut(caller.thread)?
3463 .state = GuestThreadState::Running;
3464 log::trace!("popped current thread {guest_thread:?}; new thread is {caller:?}");
3465
3466 if let Some(storage) = storage {
3467 let state = store.0.concurrent_state_mut()?;
3471 let task = state.get_mut(guest_thread.task)?;
3472 if let Some(result) = task.sync_result.take()? {
3473 if let Some(result) = result {
3474 storage[0] = MaybeUninit::new(result);
3475 }
3476
3477 if task.exited && task.ready_to_delete() {
3478 Waitable::Guest(guest_thread.task).delete_from(store.0)?;
3479 }
3480 }
3481 }
3482
3483 Ok(status.pack(waitable))
3484 }
3485
3486 pub(crate) fn first_poll<T: 'static, R: Send + 'static>(
3502 self,
3503 mut store: StoreContextMut<'_, T>,
3504 host_task: EnteredHostTask,
3505 future: impl Future<Output = Result<R>> + Send + 'static,
3506 lower: impl FnOnce(StoreContextMut<T>, Option<R>, bool, Option<TableId<HostTask>>) -> Result<()>
3507 + Send
3508 + 'static,
3509 ) -> Result<u32> {
3510 let token = StoreToken::new(store.as_context_mut());
3511
3512 let (join_handle, future) = JoinHandle::run(future);
3515 let mut future = Box::pin(future);
3516
3517 let poll = tls::set(store.0, || {
3522 future
3523 .as_mut()
3524 .poll(&mut Context::from_waker(&Waker::noop()))
3525 });
3526
3527 match poll {
3528 Poll::Ready(result) => {
3530 let result = result.transpose()?;
3531 let task = store.0.current_materialized_host_task()?;
3534 lower(store.as_context_mut(), result, true, task)?;
3535 return Ok(Status::Returned.pack(None));
3536 }
3537
3538 Poll::Pending => {}
3540 }
3541
3542 let Some(task) = store.0.materialize_host_task_id()? else {
3546 bail_bug!("current thread is not a host thread")
3547 };
3548 {
3549 let state = &mut store.0.concurrent_state_mut()?.get_mut(task)?.state;
3550 assert!(matches!(state, HostTaskState::CalleeStarted));
3551 *state = HostTaskState::CalleeRunning(join_handle);
3552 }
3553
3554 let future = Box::pin(async move {
3562 let result = match run_with_host_task_set(task, future).await? {
3563 Some(result) => Some(result?),
3564 None => None,
3565 };
3566 let on_complete = move |store: &mut dyn VMStore| {
3567 let mut store = token.as_context_mut(store);
3571 let old = store.0.set_thread(task)?;
3572
3573 let status = if result.is_some() {
3574 Status::Returned
3575 } else {
3576 Status::ReturnCancelled
3577 };
3578
3579 lower(store.as_context_mut(), result, false, Some(task))?;
3580 let state = store.0.concurrent_state_mut()?;
3581 match &mut state.get_mut(task)?.state {
3582 HostTaskState::CalleeDone { .. } => {}
3585
3586 other => *other = HostTaskState::CalleeDone { cancelled: false },
3588 }
3589 Waitable::Host(task).set_event(state, Some(Event::Subtask { status }))?;
3590
3591 store.0.set_thread(old)?;
3592 Ok(())
3593 };
3594
3595 tls::get(move |store| {
3600 store
3601 .concurrent_state_mut()?
3602 .push_high_priority(WorkItem::WorkerFunction(AlwaysMut::new(Box::new(
3603 on_complete,
3604 ))));
3605 Ok(())
3606 })
3607 });
3608
3609 let caller = match host_task {
3612 Some(caller) => caller,
3613 None => bail_bug!("host task wasn't created but should have been"),
3614 };
3615 let state = store.0.concurrent_state_mut()?;
3616 state.push_future(future);
3617 let instance = state.get_mut(caller.task)?.instance;
3618 let handle = store
3619 .0
3620 .instance_state(instance)
3621 .handle_table()
3622 .subtask_insert_host(task.rep())?;
3623 store.0.concurrent_state_mut()?.get_mut(task)?.common.handle = Some(handle);
3624 log::trace!("assign {task:?} handle {handle} for {caller:?} instance {instance:?}");
3625
3626 store.0.set_thread(caller)?;
3630 Ok(Status::Started.pack(Some(handle)))
3631 }
3632
3633 pub(crate) fn task_return(
3636 self,
3637 store: &mut dyn VMStore,
3638 ty: TypeTupleIndex,
3639 options: OptionsIndex,
3640 storage: &[ValRaw],
3641 ) -> Result<()> {
3642 let guest_thread = store.current_guest_thread()?;
3643 let state = store.concurrent_state_mut()?;
3644 let lift = state
3645 .get_mut(guest_thread.task)?
3646 .lift_result
3647 .take()
3648 .ok_or_else(|| Trap::TaskCancelOrReturnTwice)?;
3649 if !state.get_mut(guest_thread.task)?.result.is_none() {
3650 bail_bug!("task result unexpectedly already set");
3651 }
3652
3653 let CanonicalOptions {
3654 string_encoding,
3655 data_model,
3656 ..
3657 } = &self.id().get(store).component().env_component().options[options];
3658
3659 let invalid = ty != lift.ty
3660 || string_encoding != &lift.string_encoding
3661 || match data_model {
3662 CanonicalOptionsDataModel::LinearMemory(opts) => match opts.memory {
3663 Some(memory) => {
3664 let expected = lift.memory.map(|v| v.as_ptr()).unwrap_or(ptr::null_mut());
3665 let actual = self.id().get(store).runtime_memory(memory);
3666 expected != actual.as_ptr()
3667 }
3668 None => false,
3671 },
3672 CanonicalOptionsDataModel::Gc { .. } => true,
3674 };
3675
3676 if invalid {
3677 bail!(Trap::TaskReturnInvalid);
3678 }
3679
3680 log::trace!("task.return for {guest_thread:?}");
3681
3682 let result = (lift.lift)(store, storage)?;
3683 self.task_complete(store, guest_thread.task, result, Status::Returned)
3684 }
3685
3686 pub(crate) fn task_cancel(self, store: &mut StoreOpaque) -> Result<()> {
3688 let guest_thread = store.current_guest_thread()?;
3689 let state = store.concurrent_state_mut()?;
3690 let task = state.get_mut(guest_thread.task)?;
3691 if !task.cancel_request_delivered {
3692 bail!(Trap::TaskCancelNotCancelled);
3693 }
3694 _ = task
3695 .lift_result
3696 .take()
3697 .ok_or_else(|| Trap::TaskCancelOrReturnTwice)?;
3698
3699 if !task.result.is_none() {
3700 bail_bug!("task result should not bet set yet");
3701 }
3702
3703 log::trace!("task.cancel for {guest_thread:?}");
3704
3705 self.task_complete(
3706 store,
3707 guest_thread.task,
3708 Box::new(DummyResult),
3709 Status::ReturnCancelled,
3710 )
3711 }
3712
3713 fn task_complete(
3719 self,
3720 store: &mut StoreOpaque,
3721 guest_task: TableId<GuestTask>,
3722 result: Box<dyn Any + Send + Sync>,
3723 status: Status,
3724 ) -> Result<()> {
3725 store
3726 .component_resource_tables(Some(self))?
3727 .validate_scope_exit()?;
3728
3729 let state = store.concurrent_state_mut()?;
3730 let task = state.get_mut(guest_task)?;
3731
3732 if let Caller::Host { tx, .. } = &mut task.caller {
3733 if let Some(tx) = tx.take() {
3734 _ = tx.send(result);
3735 }
3736 } else {
3737 task.result = Some(result);
3738 Waitable::Guest(guest_task).set_event(state, Some(Event::Subtask { status }))?;
3739 }
3740
3741 Ok(())
3742 }
3743
3744 pub(crate) fn waitable_set_new(
3746 self,
3747 store: &mut StoreOpaque,
3748 caller_instance: RuntimeComponentInstanceIndex,
3749 ) -> Result<u32> {
3750 let set = store.concurrent_state_mut()?.push(WaitableSet::default())?;
3751 let handle = store
3752 .instance_state(self.runtime_instance(caller_instance))
3753 .handle_table()
3754 .waitable_set_insert(set.rep())?;
3755 log::trace!("new waitable set {set:?} (handle {handle})");
3756 Ok(handle)
3757 }
3758
3759 pub(crate) fn waitable_set_drop(
3761 self,
3762 store: &mut StoreOpaque,
3763 caller_instance: RuntimeComponentInstanceIndex,
3764 set: u32,
3765 ) -> Result<()> {
3766 let rep = store
3767 .instance_state(self.runtime_instance(caller_instance))
3768 .handle_table()
3769 .waitable_set_remove(set)?;
3770
3771 log::trace!("drop waitable set {rep} (handle {set})");
3772
3773 if !store
3777 .concurrent_state_mut()?
3778 .get_mut(TableId::<WaitableSet>::new(rep))?
3779 .waiting
3780 .is_empty()
3781 {
3782 bail!(Trap::WaitableSetDropHasWaiters);
3783 }
3784
3785 store
3786 .concurrent_state_mut()?
3787 .delete(TableId::<WaitableSet>::new(rep))?;
3788
3789 Ok(())
3790 }
3791
3792 pub(crate) fn waitable_join(
3794 self,
3795 store: &mut StoreOpaque,
3796 caller_instance: RuntimeComponentInstanceIndex,
3797 waitable_handle: u32,
3798 set_handle: u32,
3799 ) -> Result<()> {
3800 let mut instance = self.id().get_mut(store);
3801 let waitable =
3802 Waitable::from_instance(instance.as_mut(), caller_instance, waitable_handle)?;
3803
3804 let set = if set_handle == 0 {
3805 None
3806 } else {
3807 let set = instance.instance_states().0[caller_instance]
3808 .handle_table()
3809 .waitable_set_rep(set_handle)?;
3810
3811 let state = store.concurrent_state_mut()?;
3812 if let Some(old) = waitable.common(state)?.set
3813 && state.get_mut(old)?.is_sync_call_set
3814 {
3815 bail!(Trap::WaitableSyncAndAsync);
3816 }
3817
3818 Some(TableId::<WaitableSet>::new(set))
3819 };
3820
3821 log::trace!(
3822 "waitable {waitable:?} (handle {waitable_handle}) join set {set:?} (handle {set_handle})",
3823 );
3824
3825 waitable.join(store.concurrent_state_mut()?, set)
3826 }
3827
3828 pub(crate) fn subtask_drop(
3830 self,
3831 store: &mut StoreOpaque,
3832 caller_instance: RuntimeComponentInstanceIndex,
3833 task_id: u32,
3834 ) -> Result<()> {
3835 self.waitable_join(store, caller_instance, task_id, 0)?;
3836
3837 let (rep, is_host) = store
3838 .instance_state(self.runtime_instance(caller_instance))
3839 .handle_table()
3840 .subtask_remove(task_id)?;
3841
3842 let concurrent_state = store.concurrent_state_mut()?;
3843 let (waitable, delete) = if is_host {
3844 let id = TableId::<HostTask>::new(rep);
3845 let task = concurrent_state.get_mut(id)?;
3846 match &task.state {
3847 HostTaskState::CalleeRunning(_) => bail!(Trap::SubtaskDropNotResolved),
3848 HostTaskState::CalleeDone { .. } => {}
3849 HostTaskState::CalleeStarted | HostTaskState::CalleeFinished(_) => {
3850 bail_bug!("invalid state for callee in `subtask.drop`")
3851 }
3852 }
3853
3854 (Waitable::Host(id), true)
3855 } else {
3856 let id = TableId::<GuestTask>::new(rep);
3857 let task = concurrent_state.get_mut(id)?;
3858 if task.lift_result.is_some() {
3859 bail!(Trap::SubtaskDropNotResolved);
3860 }
3861 (
3862 Waitable::Guest(id),
3863 concurrent_state.get_mut(id)?.ready_to_delete(),
3864 )
3865 };
3866
3867 waitable.common(concurrent_state)?.handle = None;
3868
3869 if waitable.take_event(concurrent_state)?.is_some() {
3872 bail!(Trap::SubtaskDropNotResolved);
3873 }
3874
3875 if delete {
3876 waitable.delete_from(store)?;
3877 }
3878
3879 log::trace!("subtask_drop {waitable:?} (handle {task_id})");
3880 Ok(())
3881 }
3882
3883 pub(crate) fn waitable_set_wait(
3885 self,
3886 store: &mut StoreOpaque,
3887 options: OptionsIndex,
3888 set: u32,
3889 payload: u32,
3890 ) -> Result<u32> {
3891 let &CanonicalOptions {
3892 instance: caller_instance,
3893 ..
3894 } = &self.id().get(store).component().env_component().options[options];
3895 let caller = self.runtime_instance(caller_instance);
3896 let rep = store
3897 .instance_state(self.runtime_instance(caller_instance))
3898 .handle_table()
3899 .waitable_set_rep(set)?;
3900
3901 self.waitable_check(
3902 store,
3903 caller,
3904 WaitableCheck::Wait,
3905 WaitableCheckParams {
3906 set: TableId::new(rep),
3907 options,
3908 payload,
3909 },
3910 )
3911 }
3912
3913 pub(crate) fn waitable_set_poll(
3915 self,
3916 store: &mut StoreOpaque,
3917 options: OptionsIndex,
3918 set: u32,
3919 payload: u32,
3920 ) -> Result<u32> {
3921 let &CanonicalOptions {
3922 instance: caller_instance,
3923 ..
3924 } = &self.id().get(store).component().env_component().options[options];
3925 let caller = self.runtime_instance(caller_instance);
3926 let rep = store
3927 .instance_state(caller)
3928 .handle_table()
3929 .waitable_set_rep(set)?;
3930
3931 self.waitable_check(
3932 store,
3933 caller,
3934 WaitableCheck::Poll,
3935 WaitableCheckParams {
3936 set: TableId::new(rep),
3937 options,
3938 payload,
3939 },
3940 )
3941 }
3942
3943 pub(crate) fn thread_index(&self, store: &mut dyn VMStore) -> Result<u32> {
3945 let thread_id = store.current_guest_thread()?.thread;
3946 match store
3947 .concurrent_state_mut()?
3948 .get_mut(thread_id)?
3949 .instance_rep
3950 {
3951 Some(r) => Ok(r),
3952 None => bail_bug!("thread should have instance_rep by now"),
3953 }
3954 }
3955
3956 pub(crate) fn thread_new_indirect<T: 'static>(
3958 self,
3959 mut store: StoreContextMut<T>,
3960 runtime_instance: RuntimeComponentInstanceIndex,
3961 _func_ty_idx: TypeFuncIndex, start_func_table_idx: RuntimeTableIndex,
3963 start_func_idx: u32,
3964 context: i32,
3965 ) -> Result<u32> {
3966 log::trace!("creating new thread");
3967
3968 let start_func_ty = FuncType::new(store.engine(), [ValType::I32], []);
3969 let (instance, registry) = self.id().get_mut_and_registry(store.0);
3970 let callee = instance
3971 .index_runtime_func_table(registry, start_func_table_idx, start_func_idx as u64)?
3972 .ok_or_else(|| Trap::ThreadNewIndirectUninitialized)?;
3973 if callee.type_index(store.0) != start_func_ty.type_index() {
3974 bail!(Trap::ThreadNewIndirectInvalidType);
3975 }
3976
3977 let token = StoreToken::new(store.as_context_mut());
3978 let start_func = Box::new(
3979 move |store: &mut dyn VMStore, guest_thread: QualifiedThreadId| -> Result<()> {
3980 let old_thread = store.set_thread(guest_thread)?;
3981 log::trace!(
3982 "thread start: replaced {old_thread:?} with {guest_thread:?} as current thread"
3983 );
3984
3985 let mut store = token.as_context_mut(store);
3986 let mut params = [ValRaw::i32(context)];
3987 unsafe { callee.call_unchecked(store.as_context_mut(), &mut params)? };
3990
3991 store.0.set_thread(old_thread)?;
3992
3993 let runtime_instance = self.runtime_instance(runtime_instance);
3994
3995 store
3998 .0
3999 .switch_or_trap_if_may_not_suspend(runtime_instance)?;
4000
4001 store
4002 .0
4003 .cleanup_thread(guest_thread, runtime_instance, CleanupTask::Yes)?;
4004
4005 log::trace!("explicit thread {guest_thread:?} completed");
4006 let state = store.0.concurrent_state_mut()?;
4007 if let Some(t) = old_thread.guest() {
4008 state.get_mut(t.thread)?.state = GuestThreadState::Running;
4009 }
4010 log::trace!("thread start: restored {old_thread:?} as current thread");
4011
4012 Ok(())
4013 },
4014 );
4015
4016 let current_thread = store.0.current_guest_thread()?;
4017 let state = store.0.concurrent_state_mut()?;
4018 let parent_task = current_thread.task;
4019
4020 let new_thread = GuestThread::new_explicit(state, parent_task, start_func)?;
4021 let thread_id = state.push(new_thread)?;
4022 state.get_mut(parent_task)?.threads.insert(thread_id);
4023
4024 log::trace!("new thread with id {thread_id:?} created");
4025
4026 self.add_guest_thread_to_instance_table(thread_id, store.0, runtime_instance)
4027 }
4028
4029 pub(crate) fn resume_thread(
4030 self,
4031 store: &mut StoreOpaque,
4032 runtime_instance: RuntimeComponentInstanceIndex,
4033 thread_idx: u32,
4034 how: ResumeThread,
4035 ) -> Result<bool> {
4036 let thread_id =
4037 GuestThread::from_instance(self.id().get_mut(store), runtime_instance, thread_idx)?;
4038 let state = store.concurrent_state_mut()?;
4039 let guest_thread = QualifiedThreadId::qualify(state, thread_id)?;
4040
4041 if store.current_guest_thread()? == guest_thread {
4042 bail!(Trap::CannotResumeThread);
4043 }
4044
4045 let state = store.concurrent_state_mut()?;
4046 let thread = state.get_mut(guest_thread.thread)?;
4047 let priority = match how {
4048 ResumeThread::Promote | ResumeThread::Resume => Priority::Switch,
4049 ResumeThread::ResumeLater => Priority::Low,
4050 };
4051
4052 match (&how, &thread.state) {
4053 (ResumeThread::Promote, GuestThreadState::Ready { .. }) => {}
4055 (ResumeThread::Promote, _) => return Ok(false),
4056
4057 (
4060 ResumeThread::Resume | ResumeThread::ResumeLater,
4061 GuestThreadState::NotStartedExplicit(_) | GuestThreadState::Suspended(_),
4062 ) => {}
4063 (ResumeThread::Resume | ResumeThread::ResumeLater, _) => {
4064 bail!(Trap::CannotResumeThread)
4065 }
4066 }
4067
4068 match mem::replace(&mut thread.state, GuestThreadState::Running) {
4069 GuestThreadState::NotStartedExplicit(start_func) => {
4070 log::trace!("starting thread {guest_thread:?}");
4071 let guest_call = WorkItem::GuestCall {
4072 instance: self.runtime_instance(runtime_instance),
4073 call: GuestCall {
4074 thread: guest_thread,
4075 kind: GuestCallKind::StartExplicit(Box::new(move |store| {
4076 start_func(store, guest_thread)
4077 })),
4078 },
4079 };
4080 store
4081 .concurrent_state_mut()?
4082 .push_work_item(guest_call, priority)?;
4083 }
4084 GuestThreadState::Suspended(fiber) => {
4085 log::trace!("resuming thread {thread_id:?} that was suspended");
4086 store.concurrent_state_mut()?.push_work_item(
4087 WorkItem::ResumeFiber {
4088 instance: self.runtime_instance(runtime_instance),
4089 thread: guest_thread,
4090 fiber,
4091 },
4092 priority,
4093 )?;
4094 }
4095 GuestThreadState::Ready { fiber } => {
4096 log::trace!("resuming thread {thread_id:?} that was ready");
4097 thread.state = GuestThreadState::Ready { fiber };
4098 store
4099 .concurrent_state_mut()?
4100 .promote_thread_work_item(guest_thread)?;
4101 }
4102 other @ (GuestThreadState::NotStartedImplicit
4103 | GuestThreadState::Running
4104 | GuestThreadState::Completed) => {
4105 thread.state = other;
4106 }
4107 }
4108 Ok(true)
4109 }
4110
4111 fn add_guest_thread_to_instance_table(
4112 self,
4113 thread_id: TableId<GuestThread>,
4114 store: &mut StoreOpaque,
4115 runtime_instance: RuntimeComponentInstanceIndex,
4116 ) -> Result<u32> {
4117 let guest_id = store
4118 .instance_state(self.runtime_instance(runtime_instance))
4119 .thread_handle_table()
4120 .guest_thread_insert(thread_id.rep())?;
4121 store
4122 .concurrent_state_mut()?
4123 .get_mut(thread_id)?
4124 .instance_rep = Some(guest_id);
4125 Ok(guest_id)
4126 }
4127
4128 pub(crate) fn suspension_intrinsic(
4132 self,
4133 store: &mut StoreOpaque,
4134 caller: RuntimeComponentInstanceIndex,
4135 yielding: bool,
4136 to_thread: SuspensionTarget,
4137 ) -> Result<WaitResult> {
4138 let check_suspend = match to_thread {
4139 SuspensionTarget::Promote(thread) => {
4140 !self.resume_thread(store, caller, thread, ResumeThread::Promote)?
4141 }
4142 SuspensionTarget::Resume(thread) => {
4143 if !self.resume_thread(store, caller, thread, ResumeThread::Resume)? {
4144 bail_bug!(
4145 "`resume_thread` should only ever return false \
4146 when `ResumeThread::Promote` is passed to it"
4147 );
4148 }
4149 false
4150 }
4151 SuspensionTarget::None => true,
4152 };
4153
4154 if check_suspend && !store.switch_if_may_not_suspend(self.runtime_instance(caller))? {
4155 return if yielding {
4156 Ok(WaitResult::Completed)
4157 } else {
4158 Err(Trap::CannotBlockSyncTask.into())
4159 };
4160 }
4161
4162 let guest_thread = store.current_guest_thread()?;
4163
4164 let reason = if yielding {
4165 SuspendReason::Yielding {
4166 thread: guest_thread,
4167 }
4168 } else {
4169 SuspendReason::ExplicitlySuspending {
4170 thread: guest_thread,
4171 }
4172 };
4173
4174 store.suspend(reason)?;
4175
4176 Ok(WaitResult::Completed)
4177 }
4178
4179 fn waitable_check(
4181 self,
4182 store: &mut StoreOpaque,
4183 caller: RuntimeInstance,
4184 check: WaitableCheck,
4185 params: WaitableCheckParams,
4186 ) -> Result<u32> {
4187 let guest_thread = store.current_guest_thread()?;
4188
4189 log::trace!("waitable check for {guest_thread:?}; set {:?}", params.set);
4190
4191 let state = store.concurrent_state_mut()?;
4192 let task = state.get_mut(guest_thread.task)?;
4193
4194 match &check {
4197 WaitableCheck::Wait => {
4198 let set = params.set;
4199
4200 if (task.event.is_none() || matches!(task.event, Some(Event::Cancelled)))
4201 && state.get_mut(set)?.ready.is_empty()
4202 {
4203 store.switch_or_trap_if_may_not_suspend(caller)?;
4204
4205 store.suspend(SuspendReason::Waiting {
4206 set,
4207 thread: guest_thread,
4208 })?;
4209 }
4210 }
4211 WaitableCheck::Poll => {}
4212 }
4213
4214 log::trace!(
4215 "waitable check for {guest_thread:?}; set {:?}, part two",
4216 params.set
4217 );
4218
4219 let event = self.get_event(store, guest_thread.task, Some(params.set), false)?;
4221
4222 let (ordinal, handle, result) = match &check {
4223 WaitableCheck::Wait => {
4224 let (event, waitable) = match event {
4225 Some(p) => p,
4226 None => bail_bug!("event expected to be present"),
4227 };
4228 let handle = waitable.map(|(_, v)| v).unwrap_or(0);
4229 let (ordinal, result) = event.parts();
4230 (ordinal, handle, result)
4231 }
4232 WaitableCheck::Poll => {
4233 if let Some((event, waitable)) = event {
4234 let handle = waitable.map(|(_, v)| v).unwrap_or(0);
4235 let (ordinal, result) = event.parts();
4236 (ordinal, handle, result)
4237 } else {
4238 log::trace!(
4239 "no events ready to deliver via waitable-set.poll to {:?}; set {:?}",
4240 guest_thread.task,
4241 params.set
4242 );
4243 let (ordinal, result) = Event::None.parts();
4244 (ordinal, 0, result)
4245 }
4246 }
4247 };
4248 let memory = self.options_memory_mut(store, params.options);
4249 let ptr = crate::component::func::validate_inbounds_dynamic(
4250 &CanonicalAbiInfo::POINTER_PAIR,
4251 memory,
4252 &ValRaw::u32(params.payload),
4253 )?;
4254 memory[ptr + 0..][..4].copy_from_slice(&handle.to_le_bytes());
4255 memory[ptr + 4..][..4].copy_from_slice(&result.to_le_bytes());
4256 Ok(ordinal)
4257 }
4258
4259 pub(crate) fn subtask_cancel(
4261 self,
4262 store: &mut StoreOpaque,
4263 caller_instance: RuntimeComponentInstanceIndex,
4264 async_: bool,
4265 task_id: u32,
4266 ) -> Result<u32> {
4267 let (rep, is_host) = store
4268 .instance_state(self.runtime_instance(caller_instance))
4269 .handle_table()
4270 .subtask_rep(task_id)?;
4271 let waitable = if is_host {
4272 Waitable::Host(TableId::<HostTask>::new(rep))
4273 } else {
4274 Waitable::Guest(TableId::<GuestTask>::new(rep))
4275 };
4276 let concurrent_state = store.concurrent_state_mut()?;
4277
4278 log::trace!("subtask_cancel {waitable:?} (handle {task_id}; async {async_})");
4279
4280 waitable.trap_if_in_waitable_set(concurrent_state)?;
4281
4282 let needs_block;
4283 if let Waitable::Host(host_task) = waitable {
4284 let state = &mut concurrent_state.get_mut(host_task)?.state;
4285 match mem::replace(state, HostTaskState::CalleeDone { cancelled: true }) {
4286 HostTaskState::CalleeRunning(handle) => {
4293 handle.abort();
4294 needs_block = true;
4295 }
4296
4297 HostTaskState::CalleeDone { cancelled } => {
4300 if cancelled {
4301 bail!(Trap::SubtaskCancelAfterTerminal);
4302 } else {
4303 needs_block = false;
4306 }
4307 }
4308
4309 HostTaskState::CalleeStarted | HostTaskState::CalleeFinished(_) => {
4312 bail_bug!("invalid states for host callee")
4313 }
4314 }
4315 } else {
4316 let guest_task = TableId::<GuestTask>::new(rep);
4317 let task = concurrent_state.get_mut(guest_task)?;
4318 if !task.already_lowered_parameters() {
4319 store.cancel_guest_subtask_without_lowered_parameters(
4320 self.runtime_instance(caller_instance),
4321 guest_task,
4322 )?;
4323 return Ok(Status::StartCancelled as u32);
4324 } else if !task.returned_or_cancelled() {
4325 task.event = Some(Event::Cancelled);
4333 let runtime_instance = task.instance;
4334 for thread in task.threads.clone() {
4335 let thread = QualifiedThreadId {
4336 task: guest_task,
4337 thread,
4338 };
4339 let thread_mut = concurrent_state.get_mut(thread.thread)?;
4340
4341 let yield_ = |store: &mut StoreOpaque| {
4342 let state = store.instance_state(runtime_instance).concurrent_state();
4347 let old_do_not_suspend = state.do_not_suspend;
4348 state.do_not_suspend = false;
4349
4350 let caller = store.current_guest_thread()?;
4351
4352 let state = store.concurrent_state_mut()?;
4357 let set = state.get_mut(caller.thread)?.sync_call_set;
4358 waitable.join(state, Some(set))?;
4359
4360 store.suspend(SuspendReason::YieldingToSubtask { thread: caller })?;
4361
4362 let state = store.concurrent_state_mut()?;
4363 waitable.join(state, None)?;
4364
4365 store
4366 .instance_state(runtime_instance)
4367 .concurrent_state()
4368 .do_not_suspend = old_do_not_suspend;
4369
4370 Ok::<(), crate::Error>(())
4371 };
4372
4373 match thread_mut.wake_on_cancel.take() {
4374 WakeOnCancel::Waiting(set) => {
4375 let item = match concurrent_state.get_mut(set)?.waiting.remove(&thread)
4377 {
4378 Some(WaitMode::Callback(instance)) => WorkItem::GuestCall {
4379 instance: runtime_instance,
4380 call: GuestCall {
4381 thread,
4382 kind: GuestCallKind::DeliverEvent {
4383 instance,
4384 set: None,
4385 },
4386 },
4387 },
4388 other => bail_bug!(
4389 "expected `Some(WaitMode::Callback(_))`; got `{other:?}`"
4390 ),
4391 };
4392 concurrent_state.set_switch_item(item)?;
4393
4394 yield_(store)?;
4395
4396 break;
4397 }
4398 WakeOnCancel::Yielding => {
4399 if concurrent_state.promote_thread_work_item(thread)? {
4400 yield_(store)?;
4401 break;
4402 } else {
4403 bail_bug!("thread with `WakeOnCancel::Yielding` not promotable");
4404 }
4405 }
4406 WakeOnCancel::None => {}
4407 }
4408 }
4409
4410 needs_block = !store
4413 .concurrent_state_mut()?
4414 .get_mut(guest_task)?
4415 .returned_or_cancelled()
4416 } else {
4417 needs_block = false;
4418 }
4419 };
4420
4421 if needs_block {
4425 if async_ {
4426 return Ok(BLOCKED);
4427 }
4428
4429 let old_next_switch_item = {
4432 let state = store.concurrent_state_mut()?;
4433 let item = state.next_switch_item.take();
4434 state.push(item)?
4438 };
4439
4440 store.wait_for_event(self.runtime_instance(caller_instance), waitable)?;
4443
4444 let state = store.concurrent_state_mut()?;
4445 state.next_switch_item = state.delete(old_next_switch_item)?;
4446
4447 }
4449
4450 let event = waitable.take_event(store.concurrent_state_mut()?)?;
4451 if let Some(Event::Subtask {
4452 status: status @ (Status::Returned | Status::ReturnCancelled),
4453 }) = event
4454 {
4455 Ok(status as u32)
4456 } else {
4457 bail!(Trap::SubtaskCancelAfterTerminal);
4458 }
4459 }
4460}
4461
4462pub trait VMComponentAsyncStore {
4470 unsafe fn prepare_call(
4476 &mut self,
4477 instance: Instance,
4478 memory: *mut VMMemoryDefinition,
4479 start: NonNull<VMFuncRef>,
4480 return_: NonNull<VMFuncRef>,
4481 caller_instance: RuntimeComponentInstanceIndex,
4482 callee_instance: RuntimeComponentInstanceIndex,
4483 task_return_type: TypeTupleIndex,
4484 callee_async: bool,
4485 string_encoding: StringEncoding,
4486 result_count: u32,
4487 storage: *mut ValRaw,
4488 storage_len: usize,
4489 ) -> Result<()>;
4490
4491 unsafe fn sync_start(
4494 &mut self,
4495 instance: Instance,
4496 callback: *mut VMFuncRef,
4497 callee: NonNull<VMFuncRef>,
4498 param_count: u32,
4499 storage: *mut MaybeUninit<ValRaw>,
4500 storage_len: usize,
4501 ) -> Result<()>;
4502
4503 unsafe fn async_start(
4506 &mut self,
4507 instance: Instance,
4508 callback: *mut VMFuncRef,
4509 post_return: *mut VMFuncRef,
4510 callee: NonNull<VMFuncRef>,
4511 param_count: u32,
4512 result_count: u32,
4513 flags: u32,
4514 ) -> Result<u32>;
4515
4516 fn future_write(
4518 &mut self,
4519 instance: Instance,
4520 caller: RuntimeComponentInstanceIndex,
4521 ty: TypeFutureTableIndex,
4522 options: OptionsIndex,
4523 future: u32,
4524 address: u32,
4525 ) -> Result<u32>;
4526
4527 fn future_read(
4529 &mut self,
4530 instance: Instance,
4531 caller: RuntimeComponentInstanceIndex,
4532 ty: TypeFutureTableIndex,
4533 options: OptionsIndex,
4534 future: u32,
4535 address: u32,
4536 ) -> Result<u32>;
4537
4538 fn future_drop_writable(
4540 &mut self,
4541 instance: Instance,
4542 ty: TypeFutureTableIndex,
4543 writer: u32,
4544 ) -> Result<()>;
4545
4546 fn stream_write(
4548 &mut self,
4549 instance: Instance,
4550 caller: RuntimeComponentInstanceIndex,
4551 ty: TypeStreamTableIndex,
4552 options: OptionsIndex,
4553 stream: u32,
4554 address: u32,
4555 count: u32,
4556 ) -> Result<u32>;
4557
4558 fn stream_read(
4560 &mut self,
4561 instance: Instance,
4562 caller: RuntimeComponentInstanceIndex,
4563 ty: TypeStreamTableIndex,
4564 options: OptionsIndex,
4565 stream: u32,
4566 address: u32,
4567 count: u32,
4568 ) -> Result<u32>;
4569
4570 fn flat_stream_write(
4573 &mut self,
4574 instance: Instance,
4575 caller: RuntimeComponentInstanceIndex,
4576 ty: TypeStreamTableIndex,
4577 options: OptionsIndex,
4578 payload_size: u32,
4579 payload_align: u32,
4580 stream: u32,
4581 address: u32,
4582 count: u32,
4583 ) -> Result<u32>;
4584
4585 fn flat_stream_read(
4588 &mut self,
4589 instance: Instance,
4590 caller: RuntimeComponentInstanceIndex,
4591 ty: TypeStreamTableIndex,
4592 options: OptionsIndex,
4593 payload_size: u32,
4594 payload_align: u32,
4595 stream: u32,
4596 address: u32,
4597 count: u32,
4598 ) -> Result<u32>;
4599
4600 fn stream_drop_writable(
4602 &mut self,
4603 instance: Instance,
4604 ty: TypeStreamTableIndex,
4605 writer: u32,
4606 ) -> Result<()>;
4607
4608 fn error_context_debug_message(
4610 &mut self,
4611 instance: Instance,
4612 ty: TypeComponentLocalErrorContextTableIndex,
4613 options: OptionsIndex,
4614 err_ctx_handle: u32,
4615 debug_msg_address: u32,
4616 ) -> Result<()>;
4617
4618 fn thread_new_indirect(
4620 &mut self,
4621 instance: Instance,
4622 caller: RuntimeComponentInstanceIndex,
4623 func_ty_idx: TypeFuncIndex,
4624 start_func_table_idx: RuntimeTableIndex,
4625 start_func_idx: u32,
4626 context: i32,
4627 ) -> Result<u32>;
4628}
4629
4630impl<T: 'static> VMComponentAsyncStore for StoreInner<T> {
4632 unsafe fn prepare_call(
4633 &mut self,
4634 instance: Instance,
4635 memory: *mut VMMemoryDefinition,
4636 start: NonNull<VMFuncRef>,
4637 return_: NonNull<VMFuncRef>,
4638 caller_instance: RuntimeComponentInstanceIndex,
4639 callee_instance: RuntimeComponentInstanceIndex,
4640 task_return_type: TypeTupleIndex,
4641 callee_async: bool,
4642 string_encoding: StringEncoding,
4643 result_count_or_max_if_async: u32,
4644 storage: *mut ValRaw,
4645 storage_len: usize,
4646 ) -> Result<()> {
4647 let params = unsafe { core::slice::from_raw_parts(storage, storage_len) }.to_vec();
4651
4652 unsafe {
4653 instance.prepare_call(
4654 StoreContextMut(self),
4655 start,
4656 return_,
4657 caller_instance,
4658 callee_instance,
4659 task_return_type,
4660 callee_async,
4661 memory,
4662 string_encoding,
4663 match result_count_or_max_if_async {
4664 PREPARE_ASYNC_NO_RESULT => CallerInfo::Async {
4665 params,
4666 has_result: false,
4667 },
4668 PREPARE_ASYNC_WITH_RESULT => CallerInfo::Async {
4669 params,
4670 has_result: true,
4671 },
4672 result_count => CallerInfo::Sync {
4673 params,
4674 result_count,
4675 },
4676 },
4677 )
4678 }
4679 }
4680
4681 unsafe fn sync_start(
4682 &mut self,
4683 instance: Instance,
4684 callback: *mut VMFuncRef,
4685 callee: NonNull<VMFuncRef>,
4686 param_count: u32,
4687 storage: *mut MaybeUninit<ValRaw>,
4688 storage_len: usize,
4689 ) -> Result<()> {
4690 unsafe {
4691 instance
4692 .start_call(
4693 StoreContextMut(self),
4694 callback,
4695 ptr::null_mut(),
4696 callee,
4697 param_count,
4698 1,
4699 START_FLAG_ASYNC_CALLEE,
4700 Some(core::slice::from_raw_parts_mut(storage, storage_len)),
4704 )
4705 .map(drop)
4706 }
4707 }
4708
4709 unsafe fn async_start(
4710 &mut self,
4711 instance: Instance,
4712 callback: *mut VMFuncRef,
4713 post_return: *mut VMFuncRef,
4714 callee: NonNull<VMFuncRef>,
4715 param_count: u32,
4716 result_count: u32,
4717 flags: u32,
4718 ) -> Result<u32> {
4719 unsafe {
4720 instance.start_call(
4721 StoreContextMut(self),
4722 callback,
4723 post_return,
4724 callee,
4725 param_count,
4726 result_count,
4727 flags,
4728 None,
4729 )
4730 }
4731 }
4732
4733 fn future_write(
4734 &mut self,
4735 instance: Instance,
4736 caller: RuntimeComponentInstanceIndex,
4737 ty: TypeFutureTableIndex,
4738 options: OptionsIndex,
4739 future: u32,
4740 address: u32,
4741 ) -> Result<u32> {
4742 instance
4743 .guest_write(
4744 StoreContextMut(self),
4745 caller,
4746 TransmitIndex::Future(ty),
4747 options,
4748 None,
4749 future,
4750 address,
4751 1,
4752 )
4753 .map(|result| result.encode())
4754 }
4755
4756 fn future_read(
4757 &mut self,
4758 instance: Instance,
4759 caller: RuntimeComponentInstanceIndex,
4760 ty: TypeFutureTableIndex,
4761 options: OptionsIndex,
4762 future: u32,
4763 address: u32,
4764 ) -> Result<u32> {
4765 instance
4766 .guest_read(
4767 StoreContextMut(self),
4768 caller,
4769 TransmitIndex::Future(ty),
4770 options,
4771 None,
4772 future,
4773 address,
4774 1,
4775 )
4776 .map(|result| result.encode())
4777 }
4778
4779 fn stream_write(
4780 &mut self,
4781 instance: Instance,
4782 caller: RuntimeComponentInstanceIndex,
4783 ty: TypeStreamTableIndex,
4784 options: OptionsIndex,
4785 stream: u32,
4786 address: u32,
4787 count: u32,
4788 ) -> Result<u32> {
4789 instance
4790 .guest_write(
4791 StoreContextMut(self),
4792 caller,
4793 TransmitIndex::Stream(ty),
4794 options,
4795 None,
4796 stream,
4797 address,
4798 count,
4799 )
4800 .map(|result| result.encode())
4801 }
4802
4803 fn stream_read(
4804 &mut self,
4805 instance: Instance,
4806 caller: RuntimeComponentInstanceIndex,
4807 ty: TypeStreamTableIndex,
4808 options: OptionsIndex,
4809 stream: u32,
4810 address: u32,
4811 count: u32,
4812 ) -> Result<u32> {
4813 instance
4814 .guest_read(
4815 StoreContextMut(self),
4816 caller,
4817 TransmitIndex::Stream(ty),
4818 options,
4819 None,
4820 stream,
4821 address,
4822 count,
4823 )
4824 .map(|result| result.encode())
4825 }
4826
4827 fn future_drop_writable(
4828 &mut self,
4829 instance: Instance,
4830 ty: TypeFutureTableIndex,
4831 writer: u32,
4832 ) -> Result<()> {
4833 instance.guest_drop_writable(self, TransmitIndex::Future(ty), writer)
4834 }
4835
4836 fn flat_stream_write(
4837 &mut self,
4838 instance: Instance,
4839 caller: RuntimeComponentInstanceIndex,
4840 ty: TypeStreamTableIndex,
4841 options: OptionsIndex,
4842 payload_size: u32,
4843 payload_align: u32,
4844 stream: u32,
4845 address: u32,
4846 count: u32,
4847 ) -> Result<u32> {
4848 instance
4849 .guest_write(
4850 StoreContextMut(self),
4851 caller,
4852 TransmitIndex::Stream(ty),
4853 options,
4854 Some(FlatAbi {
4855 size: payload_size,
4856 align: payload_align,
4857 }),
4858 stream,
4859 address,
4860 count,
4861 )
4862 .map(|result| result.encode())
4863 }
4864
4865 fn flat_stream_read(
4866 &mut self,
4867 instance: Instance,
4868 caller: RuntimeComponentInstanceIndex,
4869 ty: TypeStreamTableIndex,
4870 options: OptionsIndex,
4871 payload_size: u32,
4872 payload_align: u32,
4873 stream: u32,
4874 address: u32,
4875 count: u32,
4876 ) -> Result<u32> {
4877 instance
4878 .guest_read(
4879 StoreContextMut(self),
4880 caller,
4881 TransmitIndex::Stream(ty),
4882 options,
4883 Some(FlatAbi {
4884 size: payload_size,
4885 align: payload_align,
4886 }),
4887 stream,
4888 address,
4889 count,
4890 )
4891 .map(|result| result.encode())
4892 }
4893
4894 fn stream_drop_writable(
4895 &mut self,
4896 instance: Instance,
4897 ty: TypeStreamTableIndex,
4898 writer: u32,
4899 ) -> Result<()> {
4900 instance.guest_drop_writable(self, TransmitIndex::Stream(ty), writer)
4901 }
4902
4903 fn error_context_debug_message(
4904 &mut self,
4905 instance: Instance,
4906 ty: TypeComponentLocalErrorContextTableIndex,
4907 options: OptionsIndex,
4908 err_ctx_handle: u32,
4909 debug_msg_address: u32,
4910 ) -> Result<()> {
4911 instance.error_context_debug_message(
4912 StoreContextMut(self),
4913 ty,
4914 options,
4915 err_ctx_handle,
4916 debug_msg_address,
4917 )
4918 }
4919
4920 fn thread_new_indirect(
4921 &mut self,
4922 instance: Instance,
4923 caller: RuntimeComponentInstanceIndex,
4924 func_ty_idx: TypeFuncIndex,
4925 start_func_table_idx: RuntimeTableIndex,
4926 start_func_idx: u32,
4927 context: i32,
4928 ) -> Result<u32> {
4929 instance.thread_new_indirect(
4930 StoreContextMut(self),
4931 caller,
4932 func_ty_idx,
4933 start_func_table_idx,
4934 start_func_idx,
4935 context,
4936 )
4937 }
4938}
4939
4940type HostTaskFuture = Pin<Box<dyn Future<Output = Result<()>> + Send + 'static>>;
4941
4942async fn run_with_host_task_set<F>(task: TableId<HostTask>, future: F) -> Result<F::Output>
4945where
4946 F: Future,
4947{
4948 let mut future = pin!(future);
4949 future::poll_fn(|cx| {
4950 let old_thread = match tls::get(|store| store.set_thread(task)) {
4951 Ok(thread) => thread,
4952 Err(error) => return Poll::Ready(Err(error)),
4953 };
4954 let result = future.as_mut().poll(cx);
4955 match tls::get(|store| store.set_thread(old_thread)) {
4956 Ok(_) => result.map(Ok),
4957 Err(error) => Poll::Ready(Err(error)),
4958 }
4959 })
4960 .await
4961}
4962
4963pub(crate) struct HostTask {
4967 common: WaitableCommon,
4968
4969 call_context: CallContext,
4972
4973 state: HostTaskState,
4974
4975 group: TaskGroupId,
4976}
4977
4978enum HostTaskState {
4979 CalleeStarted,
4984
4985 CalleeRunning(JoinHandle),
4990
4991 CalleeFinished(LiftedResult),
4995
4996 CalleeDone { cancelled: bool },
4999}
5000
5001impl HostTask {
5002 fn new(
5003 concurrent_state: &mut ConcurrentState,
5004 state: HostTaskState,
5005 caller: QualifiedThreadId,
5006 ) -> Result<Self> {
5007 let group = concurrent_state.get_mut(caller.task)?.group;
5008 concurrent_state.increment_group_ref_count(group)?;
5009
5010 Ok(Self {
5011 common: WaitableCommon::default(),
5012 call_context: CallContext::default(),
5013 state,
5014 group,
5015 })
5016 }
5017}
5018
5019impl TableDebug for HostTask {
5020 fn type_name() -> &'static str {
5021 "HostTask"
5022 }
5023}
5024
5025type CallbackFn = Box<dyn Fn(&mut dyn VMStore, Event, u32) -> Result<u32> + Send + Sync + 'static>;
5026
5027enum Caller {
5029 Host {
5031 tx: Option<oneshot::Sender<LiftedResult>>,
5033 host_future_present: bool,
5036 caller: Option<TableId<HostTask>>,
5040 },
5041 Guest {
5043 thread: QualifiedThreadId,
5045 },
5046}
5047
5048struct LiftResult {
5051 lift: RawLift,
5052 ty: TypeTupleIndex,
5053 memory: Option<SendSyncPtr<VMMemoryDefinition>>,
5054 string_encoding: StringEncoding,
5055}
5056
5057#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
5062pub(crate) struct QualifiedThreadId {
5063 task: TableId<GuestTask>,
5064 thread: TableId<GuestThread>,
5065}
5066
5067impl QualifiedThreadId {
5068 fn qualify(
5069 state: &mut ConcurrentState,
5070 thread: TableId<GuestThread>,
5071 ) -> Result<QualifiedThreadId> {
5072 Ok(QualifiedThreadId {
5073 task: state.get_mut(thread)?.parent_task,
5074 thread,
5075 })
5076 }
5077}
5078
5079impl fmt::Debug for QualifiedThreadId {
5080 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
5081 f.debug_tuple("QualifiedThreadId")
5082 .field(&self.task.rep())
5083 .field(&self.thread.rep())
5084 .finish()
5085 }
5086}
5087
5088enum GuestThreadState {
5089 NotStartedImplicit,
5090 NotStartedExplicit(
5091 Box<dyn FnOnce(&mut dyn VMStore, QualifiedThreadId) -> Result<()> + Send + Sync>,
5092 ),
5093 Running,
5094 Suspended(StoreFiber<'static>),
5095 Ready {
5096 fiber: StoreFiber<'static>,
5097 },
5098 Completed,
5099}
5100
5101impl fmt::Debug for GuestThreadState {
5102 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
5103 match self {
5104 Self::NotStartedImplicit => f.debug_tuple("NotStartedImplicit").finish(),
5105 Self::NotStartedExplicit(_) => f.debug_tuple("NotStartedExplicit").finish(),
5106 Self::Running => f.debug_tuple("Running").finish(),
5107 Self::Suspended(_) => f.debug_tuple("Suspended").finish(),
5108 Self::Ready { .. } => f.debug_struct("Ready").finish(),
5109 Self::Completed => f.debug_tuple("Completed").finish(),
5110 }
5111 }
5112}
5113
5114#[derive(Copy, Clone, PartialEq, Eq, Debug)]
5115enum WakeOnCancel {
5116 None,
5117 Waiting(TableId<WaitableSet>),
5118 Yielding,
5119}
5120
5121impl WakeOnCancel {
5122 fn is_none(self) -> bool {
5123 matches!(self, WakeOnCancel::None)
5124 }
5125
5126 fn replace(&mut self, other: WakeOnCancel) -> Self {
5127 let old = *self;
5128 *self = other;
5129 old
5130 }
5131
5132 fn take(&mut self) -> Self {
5133 self.replace(WakeOnCancel::None)
5134 }
5135}
5136
5137pub struct GuestThread {
5138 context: [u32; NUM_COMPONENT_CONTEXT_SLOTS],
5141 parent_task: TableId<GuestTask>,
5143 wake_on_cancel: WakeOnCancel,
5146 state: GuestThreadState,
5148 instance_rep: Option<u32>,
5151 sync_call_set: TableId<WaitableSet>,
5153 old_do_not_suspend: Option<bool>,
5156}
5157
5158impl GuestThread {
5159 fn from_instance(
5162 state: Pin<&mut ComponentInstance>,
5163 caller_instance: RuntimeComponentInstanceIndex,
5164 guest_thread: u32,
5165 ) -> Result<TableId<Self>> {
5166 let rep = state.instance_states().0[caller_instance]
5167 .thread_handle_table()
5168 .guest_thread_rep(guest_thread)?;
5169 Ok(TableId::new(rep))
5170 }
5171
5172 fn new_implicit(state: &mut ConcurrentState, parent_task: TableId<GuestTask>) -> Result<Self> {
5173 let sync_call_set = state.push(WaitableSet {
5174 is_sync_call_set: true,
5175 ..WaitableSet::default()
5176 })?;
5177 Ok(Self {
5178 context: [0; NUM_COMPONENT_CONTEXT_SLOTS],
5179 parent_task,
5180 wake_on_cancel: WakeOnCancel::None,
5181 state: GuestThreadState::NotStartedImplicit,
5182 instance_rep: None,
5183 sync_call_set,
5184 old_do_not_suspend: None,
5185 })
5186 }
5187
5188 fn new_explicit(
5189 state: &mut ConcurrentState,
5190 parent_task: TableId<GuestTask>,
5191 start_func: Box<
5192 dyn FnOnce(&mut dyn VMStore, QualifiedThreadId) -> Result<()> + Send + Sync,
5193 >,
5194 ) -> Result<Self> {
5195 let sync_call_set = state.push(WaitableSet {
5196 is_sync_call_set: true,
5197 ..WaitableSet::default()
5198 })?;
5199 Ok(Self {
5200 context: [0; NUM_COMPONENT_CONTEXT_SLOTS],
5201 parent_task,
5202 wake_on_cancel: WakeOnCancel::None,
5203 state: GuestThreadState::NotStartedExplicit(start_func),
5204 instance_rep: None,
5205 sync_call_set,
5206 old_do_not_suspend: None,
5207 })
5208 }
5209}
5210
5211impl TableDebug for GuestThread {
5212 fn type_name() -> &'static str {
5213 "GuestThread"
5214 }
5215}
5216
5217enum SyncResult {
5218 NotProduced,
5219 Produced(Option<ValRaw>),
5220 Taken,
5221}
5222
5223impl SyncResult {
5224 fn take(&mut self) -> Result<Option<Option<ValRaw>>> {
5225 Ok(match mem::replace(self, SyncResult::Taken) {
5226 SyncResult::NotProduced => None,
5227 SyncResult::Produced(val) => Some(val),
5228 SyncResult::Taken => {
5229 bail_bug!("attempted to take a synchronous result that was already taken")
5230 }
5231 })
5232 }
5233}
5234
5235#[derive(Debug)]
5236enum HostFutureState {
5237 NotApplicable,
5238 Live,
5239 Dropped,
5240}
5241
5242pub(crate) struct GuestTask {
5244 common: WaitableCommon,
5246 lower_params: Option<RawLower>,
5248 lift_result: Option<LiftResult>,
5250 result: Option<LiftedResult>,
5253 callback: Option<CallbackFn>,
5256 caller: Caller,
5258 call_context: CallContext,
5263 sync_result: SyncResult,
5266 cancel_request_delivered: bool,
5270 starting_sent: bool,
5273 instance: RuntimeInstance,
5280 event: Option<Event>,
5283 exited: bool,
5285 threads: HashSet<TableId<GuestThread>>,
5287 host_future_state: HostFutureState,
5290 async_typed: bool,
5293 async_lifted: bool,
5296
5297 decremented_interesting_task_count: bool,
5298
5299 group: TaskGroupId,
5300}
5301
5302impl GuestTask {
5303 fn already_lowered_parameters(&self) -> bool {
5304 self.lower_params.is_none()
5306 }
5307
5308 fn returned_or_cancelled(&self) -> bool {
5309 self.lift_result.is_none()
5311 }
5312
5313 fn ready_to_delete(&self) -> bool {
5314 let threads_completed = self.threads.is_empty();
5315 let has_sync_result = matches!(self.sync_result, SyncResult::Produced(_));
5316 let pending_completion_event = matches!(
5317 self.common.event,
5318 Some(Event::Subtask {
5319 status: Status::Returned | Status::ReturnCancelled
5320 })
5321 );
5322 let ready = threads_completed
5323 && !has_sync_result
5324 && !pending_completion_event
5325 && !matches!(self.host_future_state, HostFutureState::Live);
5326 log::trace!(
5327 "ready to delete? {ready} (threads_completed: {}, has_sync_result: {}, pending_completion_event: {}, host_future_state: {:?})",
5328 threads_completed,
5329 has_sync_result,
5330 pending_completion_event,
5331 self.host_future_state
5332 );
5333 ready
5334 }
5335
5336 fn new(
5337 state: &mut ConcurrentState,
5338 lower_params: RawLower,
5339 lift_result: LiftResult,
5340 caller: Caller,
5341 callback: Option<CallbackFn>,
5342 instance: RuntimeInstance,
5343 async_typed: bool,
5344 async_lifted: bool,
5345 ) -> Result<QualifiedThreadId> {
5346 let host_future_state = match &caller {
5347 Caller::Guest { .. } => HostFutureState::NotApplicable,
5348 Caller::Host {
5349 host_future_present,
5350 ..
5351 } => {
5352 if *host_future_present {
5353 HostFutureState::Live
5354 } else {
5355 HostFutureState::NotApplicable
5356 }
5357 }
5358 };
5359
5360 let group = match caller {
5361 Caller::Guest { thread } => {
5362 let group = state.get_mut(thread.task)?.group;
5363 state.increment_group_ref_count(group)?;
5364 group
5365 }
5366 Caller::Host { .. } => state.make_task_group()?,
5367 };
5368
5369 let task = state.push(Self {
5370 common: WaitableCommon::default(),
5371 lower_params: Some(lower_params),
5372 lift_result: Some(lift_result),
5373 result: None,
5374 callback,
5375 caller,
5376 call_context: CallContext::default(),
5377 sync_result: SyncResult::NotProduced,
5378 cancel_request_delivered: false,
5379 starting_sent: false,
5380 instance,
5381 event: None,
5382 exited: false,
5383 threads: HashSet::new(),
5384 host_future_state,
5385 async_typed,
5386 async_lifted,
5387 decremented_interesting_task_count: false,
5388 group,
5389 })?;
5390 let new_thread = GuestThread::new_implicit(state, task)?;
5391 let thread = state.push(new_thread)?;
5392 state.get_mut(task)?.threads.insert(thread);
5393 state.interesting_tasks += 1;
5394 let thread = QualifiedThreadId { task, thread };
5395 log::trace!("new implicit thread {thread:?} for instance {instance:?}");
5396 Ok(thread)
5397 }
5398}
5399
5400impl TableDebug for GuestTask {
5401 fn type_name() -> &'static str {
5402 "GuestTask"
5403 }
5404}
5405
5406#[derive(Default)]
5408struct WaitableCommon {
5409 event: Option<Event>,
5411 set: Option<TableId<WaitableSet>>,
5413 handle: Option<u32>,
5415}
5416
5417#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
5419enum Waitable {
5420 Host(TableId<HostTask>),
5422 Guest(TableId<GuestTask>),
5424 Transmit(TableId<TransmitHandle>),
5426}
5427
5428impl Waitable {
5429 fn from_instance(
5432 state: Pin<&mut ComponentInstance>,
5433 caller_instance: RuntimeComponentInstanceIndex,
5434 waitable: u32,
5435 ) -> Result<Self> {
5436 use crate::runtime::vm::component::Waitable;
5437
5438 let (waitable, kind) = state.instance_states().0[caller_instance]
5439 .handle_table()
5440 .waitable_rep(waitable)?;
5441
5442 Ok(match kind {
5443 Waitable::Subtask { is_host: true } => Self::Host(TableId::new(waitable)),
5444 Waitable::Subtask { is_host: false } => Self::Guest(TableId::new(waitable)),
5445 Waitable::Stream | Waitable::Future => Self::Transmit(TableId::new(waitable)),
5446 })
5447 }
5448
5449 fn rep(&self) -> u32 {
5451 match self {
5452 Self::Host(id) => id.rep(),
5453 Self::Guest(id) => id.rep(),
5454 Self::Transmit(id) => id.rep(),
5455 }
5456 }
5457
5458 fn join(&self, state: &mut ConcurrentState, set: Option<TableId<WaitableSet>>) -> Result<()> {
5462 log::trace!("waitable {self:?} join set {set:?}");
5463
5464 let old = mem::replace(&mut self.common(state)?.set, set);
5465
5466 if let Some(old) = old {
5467 match *self {
5468 Waitable::Host(id) => state.remove_child(id, old),
5469 Waitable::Guest(id) => state.remove_child(id, old),
5470 Waitable::Transmit(id) => state.remove_child(id, old),
5471 }?;
5472
5473 state.get_mut(old)?.ready.remove(self);
5474 }
5475
5476 if let Some(set) = set {
5477 match *self {
5478 Waitable::Host(id) => state.add_child(id, set),
5479 Waitable::Guest(id) => state.add_child(id, set),
5480 Waitable::Transmit(id) => state.add_child(id, set),
5481 }?;
5482
5483 if self.common(state)?.event.is_some() {
5484 self.mark_ready(state)?;
5485 }
5486 }
5487
5488 Ok(())
5489 }
5490
5491 fn common<'a>(&self, state: &'a mut ConcurrentState) -> Result<&'a mut WaitableCommon> {
5493 Ok(match self {
5494 Self::Host(id) => &mut state.get_mut(*id)?.common,
5495 Self::Guest(id) => &mut state.get_mut(*id)?.common,
5496 Self::Transmit(id) => &mut state.get_mut(*id)?.common,
5497 })
5498 }
5499
5500 fn trap_if_in_waitable_set(&self, state: &mut ConcurrentState) -> Result<()> {
5506 if self.common(state)?.set.is_some() {
5507 bail!(Trap::WaitableSyncAndAsync);
5508 }
5509 Ok(())
5510 }
5511
5512 fn set_event(&self, state: &mut ConcurrentState, event: Option<Event>) -> Result<()> {
5516 log::trace!("set event for {self:?}: {event:?}");
5517 self.common(state)?.event = event;
5518 self.mark_ready(state)
5519 }
5520
5521 fn take_event(&self, state: &mut ConcurrentState) -> Result<Option<Event>> {
5523 let common = self.common(state)?;
5524 let event = common.event.take();
5525 if let Some(set) = self.common(state)?.set {
5526 state.get_mut(set)?.ready.remove(self);
5527 }
5528
5529 Ok(event)
5530 }
5531
5532 fn mark_ready(&self, state: &mut ConcurrentState) -> Result<()> {
5536 if let Some(set) = self.common(state)?.set {
5537 let set_state = state.get_mut(set)?;
5538 set_state.ready.insert(*self);
5539
5540 if let Some((thread, mode)) = set_state.waiting.pop_first() {
5541 let wake_on_cancel = state.get_mut(thread.thread)?.wake_on_cancel.take();
5542 assert!(wake_on_cancel.is_none() || wake_on_cancel == WakeOnCancel::Waiting(set));
5543
5544 let item = match mode {
5545 WaitMode::Fiber(fiber) => Some(WorkItem::ResumeFiber {
5546 instance: state.get_mut(thread.task)?.instance,
5547 thread,
5548 fiber,
5549 }),
5550 WaitMode::Callback(instance) => Some(WorkItem::GuestCall {
5551 instance: state.get_mut(thread.task)?.instance,
5552 call: GuestCall {
5553 thread,
5554 kind: GuestCallKind::DeliverEvent {
5555 instance,
5556 set: Some(set),
5557 },
5558 },
5559 }),
5560 };
5561
5562 if let Some(item) = item {
5563 state.push_high_priority(item);
5564 }
5565 }
5566 }
5567 Ok(())
5568 }
5569
5570 fn delete_from(&self, store: &mut StoreOpaque) -> Result<()> {
5572 match self {
5573 Self::Host(task) => {
5574 log::trace!("delete host task {task:?}");
5575 let state = store.concurrent_state_mut()?;
5576 let task = state.delete(*task)?;
5577
5578 state.decrement_group_ref_count(task.group)?;
5579 }
5580 Self::Guest(task) => {
5581 log::trace!("delete guest task {task:?}");
5582 let state = store.concurrent_state_mut()?;
5583 let task = state.delete(*task)?;
5584
5585 state.decrement_group_ref_count(task.group)?;
5586
5587 debug_assert!(task.decremented_interesting_task_count);
5594 }
5595 Self::Transmit(task) => {
5596 store.concurrent_state_mut()?.delete(*task)?;
5597 }
5598 }
5599
5600 Ok(())
5601 }
5602}
5603
5604impl fmt::Debug for Waitable {
5605 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
5606 match self {
5607 Self::Host(id) => write!(f, "{id:?}"),
5608 Self::Guest(id) => write!(f, "{id:?}"),
5609 Self::Transmit(id) => write!(f, "{id:?}"),
5610 }
5611 }
5612}
5613
5614#[derive(Default)]
5616struct WaitableSet {
5617 ready: BTreeSet<Waitable>,
5619 waiting: BTreeMap<QualifiedThreadId, WaitMode>,
5621 is_sync_call_set: bool,
5624}
5625
5626impl TableDebug for WaitableSet {
5627 fn type_name() -> &'static str {
5628 "WaitableSet"
5629 }
5630}
5631
5632type RawLower =
5634 Box<dyn FnOnce(&mut dyn VMStore, &mut [MaybeUninit<ValRaw>]) -> Result<()> + Send + Sync>;
5635
5636type RawLift = Box<
5638 dyn FnOnce(&mut dyn VMStore, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>> + Send + Sync,
5639>;
5640
5641type LiftedResult = Box<dyn Any + Send + Sync>;
5645
5646struct DummyResult;
5649
5650#[derive(Default)]
5652pub struct ConcurrentInstanceState {
5653 backpressure: u16,
5655 do_not_enter: bool,
5657 do_not_suspend: bool,
5660 pending: BTreeMap<QualifiedThreadId, GuestCallKind>,
5663}
5664
5665impl ConcurrentInstanceState {
5666 pub fn pending_is_empty(&self) -> bool {
5667 self.pending.is_empty()
5668 }
5669}
5670
5671#[derive(Debug, Copy, Clone)]
5672pub(crate) enum CurrentThread {
5673 Guest(QualifiedThreadId),
5676 Host(TableId<HostTask>),
5678 DeferredHost(QualifiedThreadId),
5681 None,
5684}
5685
5686impl CurrentThread {
5687 fn guest(&self) -> Option<&QualifiedThreadId> {
5688 match self {
5689 Self::Guest(id) => Some(id),
5690 _ => None,
5691 }
5692 }
5693
5694 fn guest_task(&self) -> Option<TableId<GuestTask>> {
5695 match self {
5696 Self::Guest(id) => Some(id.task),
5697 _ => None,
5698 }
5699 }
5700
5701 fn is_none(&self) -> bool {
5702 matches!(self, Self::None)
5703 }
5704}
5705
5706impl From<QualifiedThreadId> for CurrentThread {
5707 fn from(id: QualifiedThreadId) -> Self {
5708 Self::Guest(id)
5709 }
5710}
5711
5712impl From<TableId<HostTask>> for CurrentThread {
5713 fn from(id: TableId<HostTask>) -> Self {
5714 Self::Host(id)
5715 }
5716}
5717
5718enum Priority {
5719 Switch,
5720 High,
5721 Low,
5722}
5723
5724pub struct ConcurrentState {
5726 unforced_current_thread: CurrentThread,
5732
5733 deferred_host_call_context: Option<CallContext>,
5739
5740 futures: AlwaysMut<Option<FuturesUnordered<HostTaskFuture>>>,
5745 table: AlwaysMut<ResourceTable>,
5747 switch_item: Option<WorkItem>,
5755 next_switch_item: Option<WorkItem>,
5761 high_priority: VecDeque<WorkItem>,
5763 low_priority: VecDeque<WorkItem>,
5765 suspend_reason: Option<SuspendReason>,
5769 worker: Option<StoreFiber<'static>>,
5773 worker_item: Option<WorkerItem>,
5775
5776 global_error_context_ref_counts:
5789 BTreeMap<TypeComponentGlobalErrorContextTableIndex, GlobalErrorContextRefCount>,
5790
5791 interesting_tasks: usize,
5804
5805 interesting_tasks_empty_waker: Option<Waker>,
5809
5810 ready_for_concurrent_call_waker: Option<Waker>,
5815
5816 event_loop_running: bool,
5818
5819 #[cfg(feature = "task-group-hook")]
5821 task_group_hook: Option<Box<dyn TaskGroupHook>>,
5822}
5823
5824impl Default for ConcurrentState {
5825 fn default() -> Self {
5826 Self {
5827 unforced_current_thread: CurrentThread::None,
5828 deferred_host_call_context: None,
5829 table: AlwaysMut::new(ResourceTable::new()),
5830 futures: AlwaysMut::new(Some(FuturesUnordered::new())),
5831 switch_item: None,
5832 next_switch_item: None,
5833 high_priority: VecDeque::new(),
5834 low_priority: VecDeque::new(),
5835 suspend_reason: None,
5836 worker: None,
5837 worker_item: None,
5838 global_error_context_ref_counts: BTreeMap::new(),
5839 interesting_tasks: 0,
5840 interesting_tasks_empty_waker: None,
5841 ready_for_concurrent_call_waker: None,
5842 event_loop_running: false,
5843 #[cfg(feature = "task-group-hook")]
5844 task_group_hook: None,
5845 }
5846 }
5847}
5848
5849impl ConcurrentState {
5850 pub(crate) fn take_fibers_and_futures(
5867 &mut self,
5868 fibers: &mut Vec<StoreFiber<'static>>,
5869 futures: &mut Vec<FuturesUnordered<HostTaskFuture>>,
5870 ) {
5871 let mut items = Vec::new();
5872 for (_, entry) in self.table.get_mut().iter_mut() {
5873 if let Some(set) = entry.downcast_mut::<WaitableSet>() {
5874 for mode in mem::take(&mut set.waiting).into_values() {
5875 match mode {
5876 WaitMode::Fiber(fiber) => {
5877 fibers.push(fiber);
5878 }
5879 WaitMode::Callback(_) => {}
5880 }
5881 }
5882 } else if let Some(thread) = entry.downcast_mut::<GuestThread>() {
5883 if let GuestThreadState::Suspended(fiber) | GuestThreadState::Ready { fiber, .. } =
5884 mem::replace(&mut thread.state, GuestThreadState::Completed)
5885 {
5886 fibers.push(fiber);
5887 }
5888 } else if let Some(item) = entry.downcast_mut::<Option<WorkItem>>() {
5889 if let Some(item) = item.take() {
5890 items.push(item);
5891 }
5892 }
5893 }
5894
5895 if let Some(fiber) = self.worker.take() {
5896 fibers.push(fiber);
5897 }
5898
5899 let mut handle_item = |item| match item {
5900 WorkItem::ResumeFiber { fiber, .. } => {
5901 fibers.push(fiber);
5902 }
5903 WorkItem::PushFuture(future) => {
5904 self.futures
5905 .get_mut()
5906 .as_mut()
5907 .unwrap()
5908 .push(future.into_inner());
5909 }
5910 WorkItem::ResumeThread { .. }
5911 | WorkItem::GuestCall { .. }
5912 | WorkItem::WorkerFunction(_) => {}
5913 };
5914
5915 for item in items {
5916 handle_item(item);
5917 }
5918 if let Some(item) = self.switch_item.take() {
5919 handle_item(item);
5920 }
5921 if let Some(item) = self.next_switch_item.take() {
5922 handle_item(item);
5923 }
5924 for item in mem::take(&mut self.high_priority) {
5925 handle_item(item);
5926 }
5927 for item in mem::take(&mut self.low_priority) {
5928 handle_item(item);
5929 }
5930
5931 if let Some(them) = self.futures.get_mut().take() {
5932 futures.push(them);
5933 }
5934 }
5935
5936 #[cfg(feature = "gc")]
5937 pub(crate) fn trace_fiber_roots(
5938 &mut self,
5939 modules: &ModuleRegistry,
5940 unwind: &dyn Unwind,
5941 gc_roots_list: &mut GcRootsList,
5942 ) {
5943 let ConcurrentState {
5944 table,
5945 worker,
5946 switch_item,
5947 next_switch_item,
5948 high_priority,
5949 low_priority,
5950
5951 futures: _,
5955
5956 worker_item: _,
5958 unforced_current_thread: _,
5959 deferred_host_call_context: _,
5960 suspend_reason: _,
5961 global_error_context_ref_counts: _,
5962 interesting_tasks: _,
5963 interesting_tasks_empty_waker: _,
5964 ready_for_concurrent_call_waker: _,
5965 event_loop_running: _,
5966 #[cfg(feature = "task-group-hook")]
5967 task_group_hook: _,
5968 } = self;
5969
5970 for (_, entry) in table.get_mut().iter_mut() {
5971 if let Some(set) = entry.downcast_mut::<WaitableSet>() {
5972 for mode in set.waiting.values_mut() {
5973 match mode {
5974 WaitMode::Fiber(fiber) => {
5975 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
5976 }
5977 WaitMode::Callback(_) => {}
5978 }
5979 }
5980 } else if let Some(thread) = entry.downcast_mut::<GuestThread>() {
5981 if let GuestThreadState::Suspended(fiber) | GuestThreadState::Ready { fiber, .. } =
5982 &mut thread.state
5983 {
5984 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
5985 }
5986 } else if let Some(Some(WorkItem::ResumeFiber { fiber, .. })) =
5987 entry.downcast_mut::<Option<WorkItem>>()
5988 {
5989 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
5990 }
5991 }
5992
5993 if let Some(fiber) = worker {
5994 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
5995 }
5996
5997 let mut handle_item = |item: &mut WorkItem| match item {
5998 WorkItem::ResumeFiber { fiber, .. } => {
5999 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6000 }
6001 WorkItem::PushFuture(_future) => {
6002 }
6005 WorkItem::ResumeThread { .. }
6006 | WorkItem::GuestCall { .. }
6007 | WorkItem::WorkerFunction(_) => {}
6008 };
6009
6010 if let Some(item) = switch_item {
6011 handle_item(item);
6012 }
6013 if let Some(item) = next_switch_item {
6014 handle_item(item);
6015 }
6016 for item in high_priority {
6017 handle_item(item);
6018 }
6019 for item in low_priority {
6020 handle_item(item);
6021 }
6022 }
6023
6024 fn push<V: Send + Sync + 'static>(
6025 &mut self,
6026 value: V,
6027 ) -> Result<TableId<V>, ResourceTableError> {
6028 self.table.get_mut().push(value).map(TableId::from)
6029 }
6030
6031 fn get_mut<V: 'static>(&mut self, id: TableId<V>) -> Result<&mut V, ResourceTableError> {
6032 self.table.get_mut().get_mut(&Resource::from(id))
6033 }
6034
6035 pub fn add_child<T: 'static, U: 'static>(
6036 &mut self,
6037 child: TableId<T>,
6038 parent: TableId<U>,
6039 ) -> Result<(), ResourceTableError> {
6040 self.table
6041 .get_mut()
6042 .add_child(Resource::from(child), Resource::from(parent))
6043 }
6044
6045 pub fn remove_child<T: 'static, U: 'static>(
6046 &mut self,
6047 child: TableId<T>,
6048 parent: TableId<U>,
6049 ) -> Result<(), ResourceTableError> {
6050 self.table
6051 .get_mut()
6052 .remove_child(Resource::from(child), Resource::from(parent))
6053 }
6054
6055 fn delete<V: 'static>(&mut self, id: TableId<V>) -> Result<V, ResourceTableError> {
6056 self.table.get_mut().delete(Resource::from(id))
6057 }
6058
6059 fn push_future(&mut self, future: HostTaskFuture) {
6060 self.push_high_priority(WorkItem::PushFuture(AlwaysMut::new(future)));
6067 }
6068
6069 fn set_switch_item(&mut self, item: WorkItem) -> Result<()> {
6070 log::trace!("set switch item: {item:?}");
6071
6072 if self.switch_item.is_some() {
6073 bail_bug!("switch item already set");
6074 }
6075
6076 self.switch_item = Some(item);
6077
6078 Ok(())
6079 }
6080
6081 fn take_next_switch_item(&mut self) -> Result<()> {
6082 if let Some(item) = self.next_switch_item.take() {
6083 self.set_switch_item(item)?;
6084 }
6085 Ok(())
6086 }
6087
6088 fn push_high_priority(&mut self, item: WorkItem) {
6089 log::trace!("push high priority: {item:?}");
6090 self.high_priority.push_front(item);
6091 }
6092
6093 fn push_low_priority(&mut self, item: WorkItem) {
6094 log::trace!("push low priority: {item:?}");
6095 self.low_priority.push_front(item);
6096 }
6097
6098 fn push_work_item(&mut self, item: WorkItem, priority: Priority) -> Result<()> {
6099 match priority {
6100 Priority::Switch => self.set_switch_item(item)?,
6101 Priority::High => self.push_high_priority(item),
6102 Priority::Low => self.push_low_priority(item),
6103 }
6104
6105 Ok(())
6106 }
6107
6108 fn promote_instance_local_thread_work_item(
6109 &mut self,
6110 current_instance: RuntimeInstance,
6111 ) -> Result<bool> {
6112 log::trace!("promote thread work items for {current_instance:?}");
6113
6114 self.promote_work_item_matching(|item: &WorkItem| {
6115 let result = match item {
6116 WorkItem::ResumeThread { instance, .. }
6117 | WorkItem::ResumeFiber { instance, .. }
6118 | WorkItem::GuestCall { instance, .. } => *instance == current_instance,
6119 _ => false,
6120 };
6121
6122 log::trace!("candidate {item:?}: {result}");
6123 result
6124 })
6125 }
6126
6127 fn promote_thread_work_item(&mut self, thread: QualifiedThreadId) -> Result<bool> {
6128 self.promote_work_item_matching(|item: &WorkItem| match item {
6129 WorkItem::ResumeThread {
6130 thread: item_thread,
6131 ..
6132 }
6133 | WorkItem::GuestCall {
6134 call:
6135 GuestCall {
6136 thread: item_thread,
6137 ..
6138 },
6139 ..
6140 } => *item_thread == thread,
6141 _ => false,
6142 })
6143 }
6144
6145 fn promote_work_item_matching<F>(&mut self, mut predicate: F) -> Result<bool>
6146 where
6147 F: FnMut(&WorkItem) -> bool,
6148 {
6149 for item in mem::take(&mut self.high_priority).into_iter().rev() {
6154 if self.switch_item.is_none() && predicate(&item) {
6155 self.set_switch_item(item)?;
6156 } else {
6157 self.push_high_priority(item);
6158 }
6159 }
6160
6161 if self.switch_item.is_none() {
6162 for item in mem::take(&mut self.low_priority).into_iter().rev() {
6163 if self.switch_item.is_none() && predicate(&item) {
6164 self.set_switch_item(item)?;
6165 } else {
6166 self.push_low_priority(item);
6167 }
6168 }
6169 }
6170
6171 Ok(self.switch_item.is_some())
6172 }
6173
6174 pub fn call_context(&mut self, task: Scope) -> Result<&mut CallContext> {
6177 match task {
6178 Scope::HostId(task) => {
6179 let task: TableId<HostTask> = TableId::new(task);
6180 Ok(&mut self.get_mut(task)?.call_context)
6181 }
6182 Scope::Id(task) => {
6183 let task: TableId<GuestTask> = TableId::new(task);
6184 Ok(&mut self.get_mut(task)?.call_context)
6185 }
6186 }
6187 }
6188
6189 pub(crate) fn deferred_host_call_context(&mut self) -> Option<&mut CallContext> {
6190 self.deferred_host_call_context.as_mut()
6191 }
6192
6193 fn futures_mut(&mut self) -> Result<&mut FuturesUnordered<HostTaskFuture>> {
6194 match self.futures.get_mut().as_mut() {
6195 Some(f) => Ok(f),
6196 None => bail_bug!("futures field of concurrent state is currently taken"),
6197 }
6198 }
6199
6200 pub(crate) fn table(&mut self) -> &mut ResourceTable {
6201 self.table.get_mut()
6202 }
6203
6204 fn debug_assert_deferred_host_invariant(&self) {
6205 debug_assert_eq!(
6206 self.deferred_host_call_context.is_some(),
6207 matches!(self.unforced_current_thread, CurrentThread::DeferredHost(_)),
6208 "a deferred host thread and call context must exist together",
6209 );
6210 }
6211
6212 fn materialize_host_task(&mut self) -> Result<CurrentThread> {
6213 self.debug_assert_deferred_host_invariant();
6214 let caller = match self.unforced_current_thread {
6215 CurrentThread::DeferredHost(caller) => caller,
6216 thread => return Ok(thread),
6217 };
6218
6219 let task = HostTask::new(self, HostTaskState::CalleeStarted, caller)?;
6221 let task = self.push(task)?;
6222 let call_context = self
6223 .deferred_host_call_context
6224 .take()
6225 .expect("deferred host call context should be present");
6226 self.get_mut(task)
6227 .expect("newly inserted host task should be present")
6228 .call_context = call_context;
6229 self.unforced_current_thread = CurrentThread::Host(task);
6230 self.debug_assert_deferred_host_invariant();
6231 log::trace!("new host task materialized {task:?}");
6232 Ok(CurrentThread::Host(task))
6233 }
6234
6235 fn materialize_current_host_task_id(&mut self) -> Result<Option<TableId<HostTask>>> {
6236 match self.materialize_host_task()? {
6237 CurrentThread::Host(id) => Ok(Some(id)),
6238 CurrentThread::None => Ok(None),
6239 CurrentThread::Guest(_) => {
6240 bail_bug!("tried to materialize a host task id from a guest thread")
6241 }
6242 CurrentThread::DeferredHost(_) => {
6243 bail_bug!(
6244 "current thread is a deferred host thread which should have been materialized"
6245 )
6246 }
6247 }
6248 }
6249
6250 pub(crate) fn materialize_current_scope(&mut self) -> Result<Scope> {
6251 match self.materialize_host_task()? {
6252 CurrentThread::Host(id) => Ok(Scope::HostId(id.rep())),
6253 _ => bail_bug!("current scope is not a deferred host scope"),
6254 }
6255 }
6256}
6257
6258fn for_any_lower<
6261 F: FnOnce(&mut dyn VMStore, &mut [MaybeUninit<ValRaw>]) -> Result<()> + Send + Sync,
6262>(
6263 fun: F,
6264) -> F {
6265 fun
6266}
6267
6268fn for_any_lift<
6270 F: FnOnce(&mut dyn VMStore, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>> + Send + Sync,
6271>(
6272 fun: F,
6273) -> F {
6274 fun
6275}
6276
6277fn check_ambient_store(id: StoreId) {
6278 let message = "\
6279 `Future`s which depend on asynchronous component tasks, streams, or \
6280 futures to complete may only be polled from the event loop of the \
6281 store to which they belong. Please use \
6282 `StoreContextMut::{run_concurrent,spawn}` to poll or await them.\
6283 ";
6284 tls::try_get(|store| {
6285 let matched = match store {
6286 tls::TryGet::Some(store) => store.id() == id,
6287 tls::TryGet::Taken | tls::TryGet::None => false,
6288 };
6289
6290 if !matched {
6291 panic!("{message}")
6292 }
6293 });
6294}
6295
6296fn unpack_callback_code(code: u32) -> (u32, u32) {
6297 (code & 0xF, code >> 4)
6298}
6299
6300struct WaitableCheckParams {
6304 set: TableId<WaitableSet>,
6305 options: OptionsIndex,
6306 payload: u32,
6307}
6308
6309enum WaitableCheck {
6312 Wait,
6313 Poll,
6314}
6315
6316pub(crate) struct PreparedCall<R> {
6318 handle: Func,
6320 thread: QualifiedThreadId,
6322 param_count: usize,
6324 rx: oneshot::Receiver<LiftedResult>,
6327 runtime_instance: RuntimeInstance,
6329 _phantom: PhantomData<R>,
6330}
6331
6332impl<R> PreparedCall<R> {
6333 pub(crate) fn task_id(&self) -> TaskId {
6335 TaskId {
6336 task: self.thread.task,
6337 runtime_instance: self.runtime_instance,
6338 }
6339 }
6340}
6341
6342pub(crate) struct TaskId {
6344 task: TableId<GuestTask>,
6345 runtime_instance: RuntimeInstance,
6346}
6347
6348impl TaskId {
6349 pub(crate) fn host_future_dropped(&self, store: &mut StoreOpaque) -> Result<()> {
6355 let task = store.concurrent_state_mut()?.get_mut(self.task)?;
6356 let delete = if !task.already_lowered_parameters() {
6357 store.cancel_guest_subtask_without_lowered_parameters(
6358 self.runtime_instance,
6359 self.task,
6360 )?;
6361 true
6362 } else {
6363 task.host_future_state = HostFutureState::Dropped;
6364 task.ready_to_delete()
6365 };
6366 if delete {
6367 Waitable::Guest(self.task).delete_from(store)?
6368 }
6369 Ok(())
6370 }
6371}
6372
6373pub(crate) fn prepare_call<T, R>(
6379 mut store: StoreContextMut<T>,
6380 handle: Func,
6381 param_count: usize,
6382 host_future_present: bool,
6383 lower_params: impl FnOnce(StoreContextMut<T>, &mut [MaybeUninit<ValRaw>]) -> Result<()>
6384 + Send
6385 + Sync
6386 + 'static,
6387 lift_result: impl FnOnce(&mut StoreOpaque, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>>
6388 + Send
6389 + Sync
6390 + 'static,
6391) -> Result<PreparedCall<R>> {
6392 if !store.0.may_enter() {
6393 bail!(Trap::CannotEnterComponent);
6394 }
6395
6396 let (options, _flags, ty, raw_options) = handle.abi_info(store.0);
6397
6398 let instance = handle.instance().id().get(store.0);
6399 let options = &instance.component().env_component().options[options];
6400 let ty = &instance.component().types()[ty];
6401 let async_typed = ty.async_;
6402 let async_lifted = raw_options.async_;
6403 let task_return_type = ty.results;
6404 let component_instance = raw_options.instance;
6405 let callback = options.callback.map(|i| instance.runtime_callback(i));
6406 let memory = options
6407 .memory()
6408 .map(|i| instance.runtime_memory(i))
6409 .map(SendSyncPtr::new);
6410 let string_encoding = options.string_encoding;
6411 let token = StoreToken::new(store.as_context_mut());
6412 let caller = store.0.materialize_host_task_id()?;
6413 let state = store.0.concurrent_state_mut()?;
6414
6415 let (tx, rx) = oneshot::channel();
6416
6417 let instance = handle.instance().runtime_instance(component_instance);
6418 let thread = GuestTask::new(
6419 state,
6420 Box::new(for_any_lower(move |store, params| {
6421 lower_params(token.as_context_mut(store), params)
6422 })),
6423 LiftResult {
6424 lift: Box::new(for_any_lift(move |store, result| {
6425 lift_result(store, result)
6426 })),
6427 ty: task_return_type,
6428 memory,
6429 string_encoding,
6430 },
6431 Caller::Host {
6432 tx: Some(tx),
6433 host_future_present,
6434 caller,
6435 },
6436 callback.map(|callback| {
6437 let callback = SendSyncPtr::new(callback);
6438 let instance = handle.instance();
6439 Box::new(move |store: &mut dyn VMStore, event, handle| {
6440 let store = token.as_context_mut(store);
6441 unsafe { instance.call_callback(store, callback, event, handle) }
6444 }) as CallbackFn
6445 }),
6446 instance,
6447 async_typed,
6448 async_lifted,
6449 )?;
6450
6451 Ok(PreparedCall {
6452 handle,
6453 thread,
6454 param_count,
6455 runtime_instance: instance,
6456 rx,
6457 _phantom: PhantomData,
6458 })
6459}
6460
6461pub(crate) struct StagedCall<R> {
6462 store: StoreId,
6463 rx: oneshot::Receiver<LiftedResult>,
6464 _marker: PhantomData<fn() -> R>,
6465 group: TaskGroupId,
6466}
6467
6468impl<R> StagedCall<R> {
6469 pub(crate) fn new<T: 'static>(
6476 mut store: StoreContextMut<T>,
6477 prepared: PreparedCall<R>,
6478 ) -> Result<StagedCall<R>> {
6479 let PreparedCall {
6480 handle,
6481 thread,
6482 param_count,
6483 rx,
6484 ..
6485 } = prepared;
6486
6487 stage_call0(store.as_context_mut(), handle, thread, param_count)?;
6488
6489 Ok(StagedCall {
6490 store: store.0.id(),
6491 rx,
6492 _marker: PhantomData,
6493 group: store.0.concurrent_state_mut()?.get_mut(thread.task)?.group,
6494 })
6495 }
6496}
6497
6498impl<R> Future for StagedCall<R>
6499where
6500 R: 'static,
6501{
6502 type Output = Result<R>;
6503
6504 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
6505 check_ambient_store(self.store);
6506 Pin::new(&mut self.rx).poll(cx).map(|result| match result {
6507 Ok(r) => match r.downcast() {
6508 Ok(r) => Ok(*r),
6509 Err(_) => bail_bug!("wrong type of value produced"),
6510 },
6511 Err(oneshot::Canceled) => bail_bug!("channel erroneously dropped"),
6512 })
6513 }
6514}
6515
6516fn stage_call0<T: 'static>(
6519 store: StoreContextMut<T>,
6520 handle: Func,
6521 guest_thread: QualifiedThreadId,
6522 param_count: usize,
6523) -> Result<()> {
6524 let (_options, _, _ty, raw_options) = handle.abi_info(store.0);
6525 let is_concurrent = raw_options.async_;
6526 let callback = raw_options.callback;
6527 let instance = handle.instance();
6528 let callee = handle.lifted_core_func(store.0);
6529 let post_return = raw_options
6530 .post_return
6531 .map(|i| instance.id().get(store.0).runtime_post_return(i));
6532 let callback = callback.map(|i| {
6533 let instance = instance.id().get(store.0);
6534 SendSyncPtr::new(instance.runtime_callback(i))
6535 });
6536
6537 log::trace!("queueing call {guest_thread:?}");
6538
6539 unsafe {
6543 instance.stage_call(
6544 store,
6545 guest_thread,
6546 SendSyncPtr::new(callee),
6547 param_count,
6548 1,
6549 is_concurrent,
6550 callback,
6551 post_return.map(SendSyncPtr::new),
6552 true,
6553 )
6554 }
6555}