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 loop {
1351 let futures = self.0.concurrent_state_mut()?.futures.get_mut().take();
1355 let mut reset = Reset {
1356 store: self.as_context_mut(),
1357 futures,
1358 };
1359 let mut next = match reset.futures.as_mut() {
1360 Some(f) => pin!(f.next()),
1361 None => bail_bug!("concurrent state missing futures field"),
1362 };
1363
1364 enum PollResult<R> {
1365 Complete(R),
1366 ProcessWork {
1367 ready: Option<WorkItem>,
1368 low_priority: bool,
1369 },
1370 }
1371
1372 let result = future::poll_fn(|cx| {
1373 if let Poll::Ready(value) = tls::set(reset.store.0, || future.as_mut().poll(cx)) {
1376 return Poll::Ready(Ok(PollResult::Complete(value)));
1377 }
1378
1379 let next = match tls::set(reset.store.0, || next.as_mut().poll(cx)) {
1383 Poll::Ready(Some(output)) => {
1384 match output {
1385 Err(e) => return Poll::Ready(Err(e)),
1386 Ok(()) => {}
1387 }
1388 Poll::Ready(true)
1389 }
1390 Poll::Ready(None) => Poll::Ready(false),
1391 Poll::Pending => Poll::Pending,
1392 };
1393
1394 let state = reset.store.0.concurrent_state_mut()?;
1409 let mut ready = state.switch_item.take();
1410 let mut low_priority = false;
1411 if ready.is_none() {
1412 ready = state.high_priority.pop_back();
1413 if ready.is_none() {
1414 ready = state.low_priority.pop_back();
1415 low_priority = true;
1416 }
1417 }
1418 if ready.is_some() {
1419 return Poll::Ready(Ok(PollResult::ProcessWork {
1420 ready,
1421 low_priority,
1422 }));
1423 }
1424
1425 return match next {
1429 Poll::Ready(true) => {
1430 Poll::Ready(Ok(PollResult::ProcessWork {
1436 ready: None,
1437 low_priority: false,
1438 }))
1439 }
1440 Poll::Ready(false) => {
1441 if let Poll::Ready(value) =
1445 tls::set(reset.store.0, || future.as_mut().poll(cx))
1446 {
1447 Poll::Ready(Ok(PollResult::Complete(value)))
1448 } else {
1449 if trap_on_idle {
1455 Poll::Ready(Err(if reset.store.0.any_may_not_suspend()? {
1462 Trap::CannotBlockSyncTask.into()
1463 } else {
1464 Trap::AsyncDeadlock.into()
1466 }))
1467 } else {
1468 Poll::Pending
1472 }
1473 }
1474 }
1475 Poll::Pending => Poll::Pending,
1480 };
1481 })
1482 .await;
1483
1484 drop(reset);
1488
1489 match result? {
1490 PollResult::Complete(value) => break Ok(value),
1493 PollResult::ProcessWork {
1496 ready,
1497 low_priority,
1498 } => {
1499 struct Dispose<'a, T: 'static> {
1500 store: StoreContextMut<'a, T>,
1501 ready: Option<WorkItem>,
1502 }
1503
1504 impl<'a, T> Drop for Dispose<'a, T> {
1505 fn drop(&mut self) {
1506 if let Some(item) = self.ready.take() {
1507 match item {
1508 WorkItem::ResumeFiber { mut fiber, .. } => {
1509 fiber.dispose(self.store.0)
1510 }
1511 WorkItem::PushFuture(future) => {
1512 tls::set(self.store.0, move || drop(future))
1513 }
1514 _ => {}
1515 }
1516 }
1517 }
1518 }
1519
1520 let mut dispose = Dispose {
1521 store: self.as_context_mut(),
1522 ready,
1523 };
1524
1525 if low_priority {
1547 dispose.store.0.yield_now().await
1548 }
1549
1550 if let Some(item) = dispose.ready.take() {
1551 dispose
1552 .store
1553 .as_context_mut()
1554 .handle_work_item(item)
1555 .await?;
1556 }
1557 }
1558 }
1559 }
1560 }
1561
1562 async fn handle_work_item(self, item: WorkItem) -> Result<()> {
1564 log::trace!("handle work item {item:?}");
1565 match item {
1566 WorkItem::PushFuture(future) => {
1567 self.0
1568 .concurrent_state_mut()?
1569 .futures_mut()?
1570 .push(future.into_inner());
1571 }
1572 WorkItem::ResumeFiber { fiber, .. } => {
1573 self.0.resume_fiber(fiber).await?;
1574 }
1575 WorkItem::ResumeThread { thread, .. } => {
1576 if let GuestThreadState::Ready { fiber, .. } = mem::replace(
1577 &mut self.0.concurrent_state_mut()?.get_mut(thread.thread)?.state,
1578 GuestThreadState::Running,
1579 ) {
1580 self.0.resume_fiber(fiber).await?;
1581 } else {
1582 bail_bug!("cannot resume non-pending thread {thread:?}");
1583 }
1584 }
1585 WorkItem::GuestCall { call, .. } => {
1586 if call.is_ready(self.0)? {
1587 self.0
1588 .concurrent_state_mut()?
1589 .get_mut(call.thread.thread)?
1590 .wake_on_cancel = WakeOnCancel::None;
1591 self.run_on_worker(WorkerItem::GuestCall(call)).await?;
1592 } else {
1593 let state = self.0.concurrent_state_mut()?;
1594 let task = state.get_mut(call.thread.task)?;
1595 if !task.starting_sent {
1596 task.starting_sent = true;
1597 if let GuestCallKind::StartImplicit(_) = &call.kind {
1598 Waitable::Guest(call.thread.task).set_event(
1599 state,
1600 Some(Event::Subtask {
1601 status: Status::Starting,
1602 }),
1603 )?;
1604 }
1605 }
1606
1607 let instance = state.get_mut(call.thread.task)?.instance;
1608 self.0
1609 .instance_state(instance)
1610 .concurrent_state()
1611 .pending
1612 .insert(call.thread, call.kind);
1613
1614 self.0.concurrent_state_mut()?.take_next_switch_item()?;
1618 }
1619 }
1620 WorkItem::WorkerFunction(fun) => {
1621 self.run_on_worker(WorkerItem::Function(fun)).await?;
1622 }
1623 }
1624
1625 Ok(())
1626 }
1627
1628 async fn run_on_worker(self, item: WorkerItem) -> Result<()> {
1630 let worker = if let Some(fiber) = self.0.concurrent_state_mut()?.worker.take() {
1631 fiber
1632 } else {
1633 unsafe {
1652 fiber::make_fiber_unchecked(self.0, move |store| {
1653 loop {
1654 let Some(item) = store.concurrent_state_mut()?.worker_item.take() else {
1655 bail_bug!("worker_item not present when resuming fiber")
1656 };
1657 match item {
1658 WorkerItem::GuestCall(call) => handle_guest_call(store, call)?,
1659 WorkerItem::Function(fun) => fun.into_inner()(store)?,
1660 }
1661
1662 store.suspend(SuspendReason::NeedWork)?;
1663 }
1664 })?
1665 }
1666 };
1667
1668 let worker_item = &mut self.0.concurrent_state_mut()?.worker_item;
1669 assert!(worker_item.is_none());
1670 *worker_item = Some(item);
1671
1672 self.0.resume_fiber(worker).await
1673 }
1674
1675 pub(crate) fn wrap_call<F, R>(self, closure: F) -> impl Future<Output = Result<R>> + 'static
1680 where
1681 T: 'static,
1682 F: FnOnce(&Accessor<T>) -> Pin<Box<dyn Future<Output = Result<R>> + Send + '_>>
1683 + Send
1684 + Sync
1685 + 'static,
1686 R: Send + Sync + 'static,
1687 {
1688 let token = StoreToken::new(self);
1689 async move {
1690 let mut accessor = Accessor::new(token);
1691 closure(&mut accessor).await
1692 }
1693 }
1694
1695 pub(crate) async fn start_instance(
1696 &mut self,
1697 instance: ModuleInstance,
1698 ) -> Result<ModuleInstance> {
1699 let (tx, rx) = oneshot::channel();
1700 let token = StoreToken::new(self.as_context_mut());
1701 self.0.queue_task(move |store| {
1702 _ = tx.send(
1703 instance
1704 .start_raw(&mut token.as_context_mut(store))
1705 .map(|()| instance),
1706 );
1707 Ok(())
1708 })?;
1709 self.as_context_mut()
1710 .run_concurrent_trap_on_idle(async |_| {
1711 rx.await
1712 .map_err(|_| format_err!("oneshot channel canceled"))
1713 })
1714 .await??
1715 }
1716}
1717
1718pub type EnteredHostTask = Option<QualifiedThreadId>;
1725
1726impl StoreOpaque {
1727 #[inline]
1731 pub(crate) fn current_thread(&mut self) -> Result<CurrentThread> {
1732 if !self.concurrency_support() {
1734 return Ok(CurrentThread::None);
1735 }
1736
1737 if !self
1740 .vm_store_context_mut()
1741 .current_thread_mut()
1742 .is_deferred()
1743 {
1744 return Ok(self
1745 .concurrent_state_mut_already_forced_current_thread()
1746 .unforced_current_thread);
1747 }
1748
1749 self.force_deferred_current_thread()
1750 }
1751
1752 #[cold]
1755 fn force_deferred_current_thread(&mut self) -> Result<CurrentThread> {
1756 let state = self.concurrent_state_mut_without_forcing_current_thread();
1765 let id = match state.unforced_current_thread.guest_task() {
1766 Some(task) => state.get_mut(task)?.instance.instance,
1767 None => bail_bug!("deferred component-model thread with non-guest base"),
1768 };
1769
1770 let mut frames = Vec::new();
1773 let mut cur = *self.vm_store_context_mut().current_thread_mut();
1774 while let Some(ptr) = cur.as_deferred() {
1775 let deferred = unsafe { ptr.as_non_null().as_ref() };
1780 frames.push((
1781 deferred.callee_async != 0,
1782 deferred.callee_instance,
1783 deferred.saved_context,
1784 ));
1785 cur = deferred.parent;
1786 }
1787
1788 *self.vm_store_context_mut().current_thread_mut() = VMLazyThread::forced();
1792
1793 let current_context = *self.vm_store_context_mut().component_context_mut();
1796
1797 for (callee_async, callee_instance, saved_context) in frames.into_iter().rev() {
1801 *self.vm_store_context_mut().component_context_mut() = saved_context;
1805 let callee = RuntimeInstance {
1806 instance: id,
1807 index: RuntimeComponentInstanceIndex::from_u32(callee_instance),
1808 };
1809 self.enter_guest_sync_call(callee_async, callee)?;
1810 }
1811
1812 *self.vm_store_context_mut().component_context_mut() = current_context;
1814
1815 Ok(self
1816 .concurrent_state_mut_without_forcing_current_thread()
1817 .unforced_current_thread)
1818 }
1819
1820 fn current_guest_thread(&mut self) -> Result<QualifiedThreadId> {
1821 match self.current_thread()?.guest() {
1822 Some(id) => Ok(*id),
1823 None => bail_bug!("current thread is not a guest thread"),
1824 }
1825 }
1826
1827 pub(crate) fn current_materialized_host_task(&mut self) -> Result<Option<TableId<HostTask>>> {
1831 match self.current_thread()? {
1832 CurrentThread::Host(id) => Ok(Some(id)),
1833 CurrentThread::DeferredHost(_) | CurrentThread::None => Ok(None),
1834 _ => bail_bug!("current thread is not a host thread"),
1835 }
1836 }
1837
1838 fn materialize_host_task_id(&mut self) -> Result<Option<TableId<HostTask>>> {
1841 Ok(self
1842 .concurrent_state_mut()?
1843 .materialize_current_host_task_id()?)
1844 }
1845
1846 fn enter_sync_call(&mut self, callee: RuntimeInstance) -> Result<()> {
1847 log::trace!("enter sync-typed call {callee:?}");
1848 let state = self.instance_state(callee).concurrent_state();
1849 let old_do_not_suspend = state.do_not_suspend;
1850 state.do_not_suspend = true;
1851
1852 let thread = self.current_guest_thread()?;
1853 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
1854 if thread.old_do_not_suspend.is_some() {
1855 bail_bug!("current thread already has `old_do_not_suspend` value");
1856 }
1857
1858 thread.old_do_not_suspend = Some(old_do_not_suspend);
1859
1860 Ok(())
1861 }
1862
1863 fn exit_sync_call(&mut self, callee: RuntimeInstance) -> Result<()> {
1864 log::trace!("exit sync-typed call {callee:?}");
1865 let thread = self.current_guest_thread()?;
1866 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
1867 let Some(old_do_not_suspend) = thread.old_do_not_suspend.take() else {
1868 bail_bug!("current thread missing `old_do_not_suspend` value");
1869 };
1870 let state = self.instance_state(callee).concurrent_state();
1871 state.do_not_suspend = old_do_not_suspend;
1872 Ok(())
1873 }
1874
1875 pub(crate) fn enter_guest_sync_call(
1887 &mut self,
1888 callee_async_typed: bool,
1889 callee: RuntimeInstance,
1890 ) -> Result<()> {
1891 log::trace!("enter sync-lifted call {callee:?}");
1892 if !self.concurrency_support() {
1893 return self.enter_call_not_concurrent();
1894 }
1895
1896 let thread = self.current_thread()?;
1897 let caller = if let Some(thread) = thread.guest() {
1898 Caller::Guest { thread: *thread }
1899 } else {
1900 Caller::Host {
1901 tx: None,
1902 host_future_present: false,
1903 caller: self.materialize_host_task_id()?,
1904 }
1905 };
1906
1907 let state = self.concurrent_state_mut()?;
1908 let guest_thread = GuestTask::new(
1909 state,
1910 Box::new(move |_, _| bail_bug!("cannot lower params in sync call")),
1911 LiftResult {
1912 lift: Box::new(move |_, _| bail_bug!("cannot lift result in sync call")),
1913 ty: TypeTupleIndex::reserved_value(),
1914 memory: None,
1915 string_encoding: StringEncoding::Utf8,
1916 },
1917 caller,
1918 None,
1919 callee,
1920 callee_async_typed,
1921 true,
1922 )?;
1923
1924 Instance::from_wasmtime(self, callee.instance).add_guest_thread_to_instance_table(
1925 guest_thread.thread,
1926 self,
1927 callee.index,
1928 )?;
1929 self.set_thread(guest_thread)?;
1930
1931 if !callee_async_typed {
1932 self.enter_sync_call(callee)?;
1933 }
1934
1935 Ok(())
1936 }
1937
1938 pub(crate) fn exit_guest_sync_call(&mut self) -> Result<()> {
1946 if !self.concurrency_support() {
1947 return Ok(self.exit_call_not_concurrent());
1948 }
1949
1950 let thread = match self.current_thread()?.guest() {
1951 Some(t) => *t,
1952 None => bail_bug!("expected task when exiting"),
1953 };
1954 let task = self.concurrent_state_mut()?.get_mut(thread.task)?;
1955 let instance = task.instance;
1956
1957 let caller = match &task.caller {
1958 &Caller::Guest { thread } => thread.into(),
1959 &Caller::Host { caller, .. } => caller
1960 .map(CurrentThread::Host)
1961 .unwrap_or(CurrentThread::None),
1962 };
1963 task.lift_result = None;
1964 task.exited = true;
1965 let async_typed = task.async_typed;
1966
1967 if !async_typed {
1968 self.exit_sync_call(instance)?;
1969 }
1970
1971 self.set_thread(caller)?;
1972
1973 log::trace!("exit sync-lifted call {instance:?}");
1974
1975 if async_typed {
1976 self.switch_or_trap_if_may_not_suspend(instance)?;
1981 }
1982
1983 self.cleanup_thread(thread, instance, CleanupTask::Yes)?;
1984
1985 Ok(())
1986 }
1987
1988 pub(crate) fn host_task_create(&mut self) -> Result<EnteredHostTask> {
1995 if !self.concurrency_support() {
1996 self.enter_call_not_concurrent()?;
1997 return Ok(None);
1998 }
1999 let caller = self.current_guest_thread()?;
2000 log::trace!("new deferred host task with caller {caller:?}");
2001
2002 self.set_thread(CurrentThread::DeferredHost(caller))?;
2003 let state = self.concurrent_state_mut()?;
2004 debug_assert!(state.deferred_host_call_context.is_none());
2005 state.deferred_host_call_context = Some(CallContext::default());
2006 state.debug_assert_deferred_host_invariant();
2007 Ok(Some(caller))
2008 }
2009
2010 pub(crate) fn host_task_delete(
2017 &mut self,
2018 original_task: EnteredHostTask,
2019 materialized_task: Option<TableId<HostTask>>,
2020 ) -> Result<()> {
2021 match original_task {
2022 Some(caller) => {
2023 self.set_thread(caller)?;
2024 if materialized_task.is_none() {
2025 let state = self.concurrent_state_mut()?;
2026 let context = state
2027 .deferred_host_call_context
2028 .take()
2029 .expect("deferred host call context should be present");
2030 debug_assert!(context.is_empty());
2031 state.debug_assert_deferred_host_invariant();
2032 }
2033 log::trace!(
2034 "delete host task with caller {original_task:?} and materialized as {materialized_task:?}"
2035 );
2036 if let Some(task) = materialized_task {
2037 Waitable::Host(task).delete_from(self)?;
2038 }
2039 }
2040 None => {
2041 debug_assert!(materialized_task.is_none());
2042 self.exit_call_not_concurrent();
2043 }
2044 }
2045 Ok(())
2046 }
2047
2048 fn instance_state(&mut self, instance: RuntimeInstance) -> &mut InstanceState {
2051 self.component_instance_mut(instance.instance)
2052 .instance_state(instance.index)
2053 }
2054
2055 pub(crate) fn set_thread(&mut self, thread: impl Into<CurrentThread>) -> Result<CurrentThread> {
2061 let thread = thread.into();
2062 let state = self.concurrent_state_mut()?;
2063 state.debug_assert_deferred_host_invariant();
2064 let old_thread = mem::replace(&mut state.unforced_current_thread, thread);
2065
2066 state.handle_thread_switch(old_thread, thread)?;
2067
2068 if let Some(old_thread) = old_thread.guest() {
2076 let old_context = *self.vm_store_context_mut().component_context_mut();
2077 self.concurrent_state_mut()?
2078 .get_mut(old_thread.thread)?
2079 .context = old_context;
2080 }
2081 if cfg!(debug_assertions) {
2082 *self.vm_store_context_mut().component_context_mut() =
2083 [u32::MAX; NUM_COMPONENT_CONTEXT_SLOTS];
2084 }
2085 if let Some(thread) = thread.guest() {
2086 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
2087 let context = thread.context;
2088 if cfg!(debug_assertions) {
2089 thread.context = [u32::MAX; NUM_COMPONENT_CONTEXT_SLOTS];
2090 }
2091 *self.vm_store_context_mut().component_context_mut() = context;
2092 }
2093
2094 *self.vm_store_context_mut().current_thread_mut() = if thread.is_none() {
2096 VMLazyThread::none()
2097 } else {
2098 VMLazyThread::forced()
2099 };
2100
2101 Ok(old_thread)
2102 }
2103
2104 fn switch_or_trap_if_may_not_suspend(&mut self, instance: RuntimeInstance) -> Result<()> {
2106 if self.switch_if_may_not_suspend(instance)? {
2107 Ok(())
2108 } else {
2109 Err(Trap::CannotBlockSyncTask.into())
2110 }
2111 }
2112
2113 fn switch_if_may_not_suspend(&mut self, instance: RuntimeInstance) -> Result<bool> {
2117 self.concurrent_state_mut()?;
2121
2122 Ok(!self.concurrency_support()
2123 || !self
2124 .instance_state(instance)
2125 .concurrent_state()
2126 .do_not_suspend
2127 || self
2128 .concurrent_state_mut()?
2129 .promote_instance_local_thread_work_item(instance)?)
2130 }
2131
2132 fn enter_instance(&mut self, instance: RuntimeInstance) {
2136 log::trace!("enter {instance:?}");
2137 self.instance_state(instance)
2138 .concurrent_state()
2139 .do_not_enter = true;
2140 }
2141
2142 fn exit_instance(&mut self, instance: RuntimeInstance) -> Result<()> {
2146 log::trace!("exit {instance:?}");
2147 self.instance_state(instance)
2148 .concurrent_state()
2149 .do_not_enter = false;
2150 self.partition_pending(instance)
2151 }
2152
2153 fn partition_pending(&mut self, instance: RuntimeInstance) -> Result<()> {
2161 for (thread, kind) in
2162 mem::take(&mut self.instance_state(instance).concurrent_state().pending).into_iter()
2163 {
2164 let call = GuestCall { thread, kind };
2165 if call.is_ready(self)? {
2166 self.concurrent_state_mut()?
2167 .push_high_priority(WorkItem::GuestCall { instance, call });
2168 } else {
2169 self.instance_state(instance)
2170 .concurrent_state()
2171 .pending
2172 .insert(call.thread, call.kind);
2173 }
2174 }
2175
2176 if let Some(waker) = self
2177 .concurrent_state_mut()?
2178 .ready_for_concurrent_call_waker
2179 .take()
2180 {
2181 waker.wake();
2182 }
2183
2184 Ok(())
2185 }
2186
2187 pub(crate) fn backpressure_modify(
2189 &mut self,
2190 caller_instance: RuntimeInstance,
2191 modify: impl FnOnce(u16) -> Option<u16>,
2192 ) -> Result<()> {
2193 let state = self.instance_state(caller_instance).concurrent_state();
2194 let old = state.backpressure;
2195 let new = modify(old).ok_or_else(|| Trap::BackpressureOverflow)?;
2196 state.backpressure = new;
2197
2198 if old > 0 && new == 0 {
2199 self.partition_pending(caller_instance)?;
2202 }
2203
2204 Ok(())
2205 }
2206
2207 async fn resume_fiber(&mut self, fiber: StoreFiber<'static>) -> Result<()> {
2210 let old_thread = self.current_thread()?;
2211 log::trace!("resume_fiber: save current thread {old_thread:?}");
2212
2213 let fiber = fiber::resolve_or_release(self, fiber).await?;
2214
2215 self.set_thread(old_thread)?;
2216
2217 let state = self.concurrent_state_mut()?;
2218
2219 if let Some(ot) = old_thread.guest() {
2220 state.get_mut(ot.thread)?.state = GuestThreadState::Running;
2221 }
2222 log::trace!("resume_fiber: restore current thread {old_thread:?}");
2223
2224 if let Some(mut fiber) = fiber {
2225 log::trace!("resume_fiber: suspend reason {:?}", &state.suspend_reason);
2226 let reason = match state.suspend_reason.take() {
2228 Some(r) => r,
2229 None => bail_bug!("suspend reason missing when resuming fiber"),
2230 };
2231 match reason {
2232 SuspendReason::NeedWork => {
2233 if state.worker.is_none() {
2234 state.worker = Some(fiber);
2235 } else {
2236 fiber.dispose(self);
2237 }
2238 }
2239 SuspendReason::Yielding { thread } => {
2240 state.get_mut(thread.thread)?.state = GuestThreadState::Ready { fiber };
2241 let instance = state.get_mut(thread.task)?.instance;
2242 state.push_low_priority(WorkItem::ResumeThread { instance, thread });
2243 }
2244 SuspendReason::ExplicitlySuspending { thread } => {
2245 state.get_mut(thread.thread)?.state = GuestThreadState::Suspended(fiber);
2246 }
2247 SuspendReason::Waiting { set, thread } => {
2248 let old = state
2249 .get_mut(set)?
2250 .waiting
2251 .insert(thread, WaitMode::Fiber(fiber));
2252 assert!(old.is_none());
2253 }
2254 SuspendReason::YieldingToSubtask { thread } => {
2255 let item = WorkItem::ResumeFiber {
2264 instance: state.get_mut(thread.task)?.instance,
2265 thread,
2266 fiber,
2267 };
2268
2269 if state.next_switch_item.replace(item).is_some() {
2270 bail_bug!(
2273 "`ConcurrentState::next_switch_item` was already `Some(_)` when \
2274 a thread wanted to wait on a subtask"
2275 );
2276 }
2277 }
2278 };
2279 } else {
2280 log::trace!("resume_fiber: fiber has exited");
2281 }
2282
2283 Ok(())
2284 }
2285
2286 fn suspend(&mut self, reason: SuspendReason) -> Result<()> {
2292 log::trace!("suspend fiber: {reason:?}");
2293
2294 let state = self.concurrent_state_mut()?;
2295
2296 let (save_and_restore_thread, save_and_restore_next_switch_item) = match &reason {
2303 SuspendReason::Yielding { .. }
2304 | SuspendReason::Waiting { .. }
2305 | SuspendReason::ExplicitlySuspending { .. } => {
2306 if state.switch_item.is_none() {
2309 state.take_next_switch_item()?;
2310 }
2311
2312 (true, false)
2313 }
2314 SuspendReason::YieldingToSubtask { .. } => (true, true),
2315 SuspendReason::NeedWork => (false, false),
2316 };
2317
2318 let old_next_switch_item = if save_and_restore_next_switch_item {
2319 let item = state.next_switch_item.take();
2320 Some(state.push(item)?)
2324 } else {
2325 None
2326 };
2327
2328 let old_guest_thread = if save_and_restore_thread {
2329 self.current_thread()?
2330 } else {
2331 CurrentThread::None
2332 };
2333
2334 let suspend_reason = &mut self.concurrent_state_mut()?.suspend_reason;
2335 assert!(suspend_reason.is_none());
2336 *suspend_reason = Some(reason);
2337
2338 if !self.fiber_async_state_mut().can_block() {
2341 return Err(format_err!("future dropped"));
2342 }
2343
2344 self.with_blocking(|_, cx| cx.suspend(StoreFiberYield::ReleaseStore))?;
2345
2346 if save_and_restore_thread {
2347 self.set_thread(old_guest_thread)?;
2348 }
2349
2350 if let Some(item) = old_next_switch_item {
2351 let state = self.concurrent_state_mut()?;
2352 state.next_switch_item = state.delete(item)?;
2353 }
2354
2355 Ok(())
2356 }
2357
2358 fn wait_for_event(
2359 &mut self,
2360 caller_instance: RuntimeInstance,
2361 waitable: Waitable,
2362 ) -> Result<()> {
2363 let caller = self.current_guest_thread()?;
2364 let state = self.concurrent_state_mut()?;
2365
2366 waitable.trap_if_in_waitable_set(state)?;
2367
2368 let set = state.get_mut(caller.thread)?.sync_call_set;
2369 waitable.join(state, Some(set))?;
2370
2371 self.switch_or_trap_if_may_not_suspend(caller_instance)?;
2372
2373 self.suspend(SuspendReason::Waiting {
2374 set,
2375 thread: caller,
2376 })?;
2377 let state = self.concurrent_state_mut()?;
2378
2379 waitable.join(state, None)
2380 }
2381
2382 fn cleanup_thread(
2404 &mut self,
2405 guest_thread: QualifiedThreadId,
2406 runtime_instance: RuntimeInstance,
2407 cleanup_task: CleanupTask,
2408 ) -> Result<()> {
2409 let state = self.concurrent_state_mut()?;
2410 state.take_next_switch_item()?;
2413 let thread_data = state.get_mut(guest_thread.thread)?;
2414 let sync_call_set = thread_data.sync_call_set;
2415 if let Some(guest_id) = thread_data.instance_rep {
2416 self.instance_state(runtime_instance)
2417 .thread_handle_table()
2418 .guest_thread_remove(guest_id)?;
2419 }
2420 let state = self.concurrent_state_mut()?;
2421
2422 for waitable in mem::take(&mut state.get_mut(sync_call_set)?.ready) {
2424 if let Some(Event::Subtask {
2425 status: Status::Returned | Status::ReturnCancelled,
2426 }) = waitable.common(self.concurrent_state_mut()?)?.event
2427 {
2428 waitable.delete_from(self)?;
2429 }
2430 }
2431
2432 let state = self.concurrent_state_mut()?;
2433 state.delete(guest_thread.thread)?;
2434 state.delete(sync_call_set)?;
2435 let task = state.get_mut(guest_thread.task)?;
2436 task.threads.remove(&guest_thread.thread);
2437
2438 if task.threads.is_empty() && !task.returned_or_cancelled() {
2439 bail!(Trap::NoAsyncResult);
2440 }
2441 let ready_to_delete = task.ready_to_delete();
2442
2443 if !task.decremented_interesting_task_count && task.exited && task.returned_or_cancelled() {
2444 task.decremented_interesting_task_count = true;
2445
2446 debug_assert!(state.interesting_tasks > 0);
2447 state.interesting_tasks -= 1;
2448 if state.interesting_tasks == 0
2449 && let Some(waker) = state.interesting_tasks_empty_waker.take()
2450 {
2451 waker.wake();
2452 }
2453 }
2454
2455 match cleanup_task {
2456 CleanupTask::Yes => {
2457 if ready_to_delete {
2458 Waitable::Guest(guest_thread.task).delete_from(self)?;
2459 }
2460 }
2461 CleanupTask::No => {}
2462 }
2463
2464 Ok(())
2465 }
2466
2467 fn cancel_guest_subtask_without_lowered_parameters(
2480 &mut self,
2481 caller_instance: RuntimeInstance,
2482 guest_task: TableId<GuestTask>,
2483 ) -> Result<()> {
2484 let concurrent_state = self.concurrent_state_mut()?;
2485 let task = concurrent_state.get_mut(guest_task)?;
2486 assert!(!task.already_lowered_parameters());
2487 task.lower_params = None;
2491 task.lift_result = None;
2492 task.exited = true;
2493 let instance = task.instance;
2494
2495 assert_eq!(1, task.threads.len());
2498 let thread = *task.threads.iter().next().unwrap();
2499 self.cleanup_thread(
2500 QualifiedThreadId {
2501 task: guest_task,
2502 thread,
2503 },
2504 caller_instance,
2505 CleanupTask::No,
2506 )?;
2507
2508 let pending = &mut self.instance_state(instance).concurrent_state().pending;
2510 let pending_count = pending.len();
2511 pending.retain(|thread, _| thread.task != guest_task);
2512 if pending.len() == pending_count {
2514 bail!(Trap::SubtaskCancelAfterTerminal);
2515 }
2516 Ok(())
2517 }
2518
2519 pub(crate) fn current_scope(&mut self) -> Result<Option<CurrentScope>> {
2522 if !self.concurrency_support() {
2523 return Ok(self
2524 .current_scope_id_not_concurrent()?
2525 .map(|id| CurrentScope::Id(Scope::Id(id))));
2526 }
2527
2528 Ok(match self.current_thread()? {
2529 CurrentThread::Guest(id) => Some(CurrentScope::Id(Scope::Id(id.task.rep()))),
2530 CurrentThread::Host(id) => Some(CurrentScope::Id(Scope::HostId(id.rep()))),
2531 CurrentThread::DeferredHost(_) => Some(CurrentScope::DeferredHost),
2532 CurrentThread::None => return Ok(None),
2533 })
2534 }
2535
2536 pub(crate) fn queue_task(
2537 &mut self,
2538 task: impl FnOnce(&mut dyn VMStore) -> Result<()> + Send + 'static,
2539 ) -> Result<()> {
2540 self.concurrent_state_mut()?
2541 .push_high_priority(WorkItem::WorkerFunction(AlwaysMut::new(Box::new(task))));
2542 Ok(())
2543 }
2544
2545 fn any_may_not_suspend(&mut self) -> Result<bool> {
2554 Ok(self
2562 .concurrent_state_mut()?
2563 .table
2564 .get_mut()
2565 .iter_mut()
2566 .filter_map(|(_, entry)| {
2567 if let Some(task) = entry.downcast_ref::<GuestTask>() {
2568 Some(task.instance)
2569 } else {
2570 None
2571 }
2572 })
2573 .collect::<Vec<_>>()
2574 .into_iter()
2575 .any(|instance| {
2576 self.instance_state(instance)
2577 .concurrent_state()
2578 .do_not_suspend
2579 }))
2580 }
2581}
2582
2583enum CleanupTask {
2584 Yes,
2585 No,
2586}
2587
2588impl Instance {
2589 fn get_event(
2592 self,
2593 store: &mut StoreOpaque,
2594 guest_task: TableId<GuestTask>,
2595 set: Option<TableId<WaitableSet>>,
2596 cancellable: bool,
2597 ) -> Result<Option<(Event, Option<(Waitable, u32)>)>> {
2598 let state = store.concurrent_state_mut()?;
2599
2600 let task = state.get_mut(guest_task)?;
2601 let event = &mut task.event;
2602 if let Some(ev) = event
2603 && (cancellable || !matches!(ev, Event::Cancelled))
2604 {
2605 log::trace!("deliver event {ev:?} to {guest_task:?}");
2606
2607 if matches!(ev, Event::Cancelled) {
2608 task.cancel_request_delivered = true;
2609 }
2610
2611 let ev = *ev;
2612 *event = None;
2613 return Ok(Some((ev, None)));
2614 }
2615
2616 let set = match set {
2617 Some(set) => set,
2618 None => return Ok(None),
2619 };
2620 let waitable = match state.get_mut(set)?.ready.pop_first() {
2621 Some(v) => v,
2622 None => return Ok(None),
2623 };
2624
2625 let common = waitable.common(state)?;
2626 let handle = match common.handle {
2627 Some(h) => h,
2628 None => bail_bug!("handle not set when delivering event"),
2629 };
2630 let event = match common.event.take() {
2631 Some(e) => e,
2632 None => bail_bug!("event not set when delivering event"),
2633 };
2634
2635 log::trace!(
2636 "deliver event {event:?} to {guest_task:?} for {waitable:?} (handle {handle}); set {set:?}"
2637 );
2638
2639 waitable.on_delivery(store, self, event)?;
2640
2641 Ok(Some((event, Some((waitable, handle)))))
2642 }
2643
2644 fn handle_callback_code(
2650 self,
2651 store: &mut StoreOpaque,
2652 guest_thread: QualifiedThreadId,
2653 runtime_instance: RuntimeComponentInstanceIndex,
2654 code: u32,
2655 ) -> Result<()> {
2656 let (code, set) = unpack_callback_code(code);
2657
2658 log::trace!("received callback code from {guest_thread:?}: {code} (set: {set})");
2659
2660 let state = store.concurrent_state_mut()?;
2661
2662 state.take_next_switch_item()?;
2663
2664 let get_set = |store: &mut StoreOpaque, handle| -> Result<_> {
2665 let set = store
2666 .instance_state(self.runtime_instance(runtime_instance))
2667 .handle_table()
2668 .waitable_set_rep(handle)?;
2669
2670 Ok(TableId::<WaitableSet>::new(set))
2671 };
2672
2673 match code {
2674 callback_code::EXIT => {
2675 log::trace!("implicit thread {guest_thread:?} completed");
2676 let task = store.concurrent_state_mut()?.get_mut(guest_thread.task)?;
2677 task.exited = true;
2678 task.callback = None;
2679
2680 let runtime_instance = self.runtime_instance(runtime_instance);
2681
2682 store.switch_or_trap_if_may_not_suspend(runtime_instance)?;
2687
2688 store.cleanup_thread(guest_thread, runtime_instance, CleanupTask::Yes)?;
2689 }
2690 callback_code::YIELD => {
2691 let old = state
2694 .get_mut(guest_thread.thread)?
2695 .wake_on_cancel
2696 .replace(WakeOnCancel::Yielding);
2697 if !old.is_none() {
2698 bail_bug!("thread unexpectedly had wake_on_cancel set");
2699 }
2700
2701 let task = state.get_mut(guest_thread.task)?;
2702 if let Some(event) = task.event {
2707 assert!(matches!(event, Event::None | Event::Cancelled));
2708 } else {
2709 task.event = Some(Event::None);
2710 }
2711 let call = GuestCall {
2712 thread: guest_thread,
2713 kind: GuestCallKind::DeliverEvent {
2714 instance: self,
2715 set: None,
2716 },
2717 };
2718 state.push_low_priority(WorkItem::GuestCall {
2721 instance: self.runtime_instance(runtime_instance),
2722 call,
2723 });
2724 }
2725 callback_code::WAIT => {
2726 let set = get_set(store, set)?;
2727 let state = store.concurrent_state_mut()?;
2728
2729 if state.get_mut(guest_thread.task)?.event.is_some()
2730 || !state.get_mut(set)?.ready.is_empty()
2731 {
2732 state.push_high_priority(WorkItem::GuestCall {
2734 instance: self.runtime_instance(runtime_instance),
2735 call: GuestCall {
2736 thread: guest_thread,
2737 kind: GuestCallKind::DeliverEvent {
2738 instance: self,
2739 set: Some(set),
2740 },
2741 },
2742 });
2743 } else {
2744 let old = state
2752 .get_mut(guest_thread.thread)?
2753 .wake_on_cancel
2754 .replace(WakeOnCancel::Waiting(set));
2755 if !old.is_none() {
2756 bail_bug!("thread unexpectedly had wake_on_cancel set");
2757 }
2758 let old = state
2759 .get_mut(set)?
2760 .waiting
2761 .insert(guest_thread, WaitMode::Callback(self));
2762 if !old.is_none() {
2763 bail_bug!("set's waiting set already had this thread registered");
2764 }
2765 }
2766 }
2767 _ => bail!(Trap::UnsupportedCallbackCode),
2768 }
2769
2770 Ok(())
2771 }
2772
2773 unsafe fn stage_call<T: 'static>(
2780 self,
2781 mut store: StoreContextMut<T>,
2782 guest_thread: QualifiedThreadId,
2783 callee: SendSyncPtr<VMFuncRef>,
2784 param_count: usize,
2785 result_count: usize,
2786 async_: bool,
2787 callback: Option<SendSyncPtr<VMFuncRef>>,
2788 post_return: Option<SendSyncPtr<VMFuncRef>>,
2789 host_caller: bool,
2790 ) -> Result<()> {
2791 unsafe fn make_call<T: 'static>(
2806 store: StoreContextMut<T>,
2807 guest_thread: QualifiedThreadId,
2808 callee: SendSyncPtr<VMFuncRef>,
2809 param_count: usize,
2810 result_count: usize,
2811 ) -> impl FnOnce(&mut dyn VMStore) -> Result<[MaybeUninit<ValRaw>; MAX_FLAT_PARAMS]>
2812 + Send
2813 + Sync
2814 + 'static
2815 + use<T> {
2816 let token = StoreToken::new(store);
2817 move |store: &mut dyn VMStore| {
2818 let mut storage = [MaybeUninit::uninit(); MAX_FLAT_PARAMS];
2819
2820 store
2821 .concurrent_state_mut()?
2822 .get_mut(guest_thread.thread)?
2823 .state = GuestThreadState::Running;
2824 let task = store.concurrent_state_mut()?.get_mut(guest_thread.task)?;
2825 let lower = match task.lower_params.take() {
2826 Some(l) => l,
2827 None => bail_bug!("lower_params missing"),
2828 };
2829
2830 lower(store, &mut storage[..param_count])?;
2831
2832 let mut store = token.as_context_mut(store);
2833
2834 unsafe {
2837 crate::Func::call_unchecked_raw(
2838 &mut store,
2839 callee.as_non_null(),
2840 NonNull::new(
2841 &mut storage[..param_count.max(result_count)]
2842 as *mut [MaybeUninit<ValRaw>] as _,
2843 )
2844 .unwrap(),
2845 )?;
2846 }
2847
2848 Ok(storage)
2849 }
2850 }
2851
2852 let call = unsafe {
2856 make_call(
2857 store.as_context_mut(),
2858 guest_thread,
2859 callee,
2860 param_count,
2861 result_count,
2862 )
2863 };
2864
2865 let callee_instance = store
2866 .0
2867 .concurrent_state_mut()?
2868 .get_mut(guest_thread.task)?
2869 .instance;
2870
2871 let fun = if callback.is_some() {
2872 assert!(async_);
2873
2874 Box::new(move |store: &mut dyn VMStore| {
2875 self.add_guest_thread_to_instance_table(
2876 guest_thread.thread,
2877 store,
2878 callee_instance.index,
2879 )?;
2880 let old_thread = store.set_thread(guest_thread)?;
2881 log::trace!(
2882 "stackless call: replaced {old_thread:?} with {guest_thread:?} as current thread"
2883 );
2884
2885 store.enter_instance(callee_instance);
2886
2887 let storage = call(store)?;
2894
2895 store.exit_instance(callee_instance)?;
2896
2897 store.set_thread(old_thread)?;
2898 let state = store.concurrent_state_mut()?;
2899 if let Some(t) = old_thread.guest() {
2900 state.get_mut(t.thread)?.state = GuestThreadState::Running;
2901 }
2902 log::trace!("stackless call: restored {old_thread:?} as current thread");
2903
2904 let code = unsafe { storage[0].assume_init() }.get_i32() as u32;
2907
2908 self.handle_callback_code(store, guest_thread, callee_instance.index, code)
2909 }) as Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send + Sync>
2910 } else {
2911 let token = StoreToken::new(store.as_context_mut());
2912 Box::new(move |store: &mut dyn VMStore| {
2913 self.add_guest_thread_to_instance_table(
2914 guest_thread.thread,
2915 store,
2916 callee_instance.index,
2917 )?;
2918 let old_thread = store.set_thread(guest_thread)?;
2919 log::trace!(
2920 "sync/async-stackful call: replaced {old_thread:?} with {guest_thread:?} as current thread",
2921 );
2922 let flags = self.id().get(store).instance_flags(callee_instance.index);
2923
2924 let callee_async_typed = store
2925 .concurrent_state_mut()?
2926 .get_mut(guest_thread.task)?
2927 .async_typed;
2928
2929 if !async_ && callee_async_typed {
2933 store.enter_instance(callee_instance);
2934 }
2935
2936 if !callee_async_typed {
2937 store.enter_sync_call(callee_instance)?;
2938 }
2939
2940 let storage = call(store)?;
2947
2948 if !callee_async_typed {
2949 store.exit_sync_call(callee_instance)?;
2950 }
2951
2952 if !async_ {
2953 if callee_async_typed {
2959 store.exit_instance(callee_instance)?;
2960 }
2961
2962 let lift = {
2963 let state = store.concurrent_state_mut()?;
2964 if !state.get_mut(guest_thread.task)?.result.is_none() {
2965 bail_bug!("task has already produced a result");
2966 }
2967
2968 match state.get_mut(guest_thread.task)?.lift_result.take() {
2969 Some(lift) => lift,
2970 None => bail_bug!("lift_result field is missing"),
2971 }
2972 };
2973
2974 let result = (lift.lift)(store, unsafe {
2977 mem::transmute::<&[MaybeUninit<ValRaw>], &[ValRaw]>(
2978 &storage[..result_count],
2979 )
2980 })?;
2981
2982 let post_return_arg = match result_count {
2983 0 => ValRaw::i32(0),
2984 1 => unsafe { storage[0].assume_init() },
2987 _ => unreachable!(),
2988 };
2989
2990 unsafe {
2991 call_post_return(
2992 token.as_context_mut(store),
2993 post_return.map(|v| v.as_non_null()),
2994 post_return_arg,
2995 flags,
2996 )?;
2997 }
2998
2999 self.task_complete(store, guest_thread.task, result, Status::Returned)?;
3000 }
3001
3002 store.set_thread(old_thread)?;
3003
3004 store
3005 .concurrent_state_mut()?
3006 .get_mut(guest_thread.task)?
3007 .exited = true;
3008
3009 log::trace!(
3010 "clean up thread; async lifted? {async_} async typed? {callee_async_typed}"
3011 );
3012
3013 if callee_async_typed {
3014 store.switch_or_trap_if_may_not_suspend(callee_instance)?;
3019 }
3020
3021 store.cleanup_thread(guest_thread, callee_instance, CleanupTask::Yes)?;
3023 Ok(())
3024 })
3025 };
3026
3027 store.0.concurrent_state_mut()?.push_work_item(
3028 WorkItem::GuestCall {
3029 instance: callee_instance,
3030 call: GuestCall {
3031 thread: guest_thread,
3032 kind: GuestCallKind::StartImplicit(fun),
3033 },
3034 },
3035 if host_caller {
3036 Priority::High
3037 } else {
3038 Priority::Switch
3039 },
3040 )?;
3041
3042 Ok(())
3043 }
3044
3045 unsafe fn prepare_call<T: 'static>(
3058 self,
3059 mut store: StoreContextMut<T>,
3060 start: NonNull<VMFuncRef>,
3061 return_: NonNull<VMFuncRef>,
3062 caller_instance: RuntimeComponentInstanceIndex,
3063 callee_instance: RuntimeComponentInstanceIndex,
3064 task_return_type: TypeTupleIndex,
3065 callee_async_typed: bool,
3066 memory: *mut VMMemoryDefinition,
3067 string_encoding: StringEncoding,
3068 caller_info: CallerInfo,
3069 ) -> Result<()> {
3070 enum ResultInfo {
3071 Heap { results: u32 },
3072 Stack { result_count: u32 },
3073 }
3074
3075 let result_info = match &caller_info {
3076 CallerInfo::Async {
3077 has_result: true,
3078 params,
3079 } => ResultInfo::Heap {
3080 results: match params.last() {
3081 Some(r) => r.get_u32(),
3082 None => bail_bug!("retptr missing"),
3083 },
3084 },
3085 CallerInfo::Async {
3086 has_result: false, ..
3087 } => ResultInfo::Stack { result_count: 0 },
3088 CallerInfo::Sync {
3089 result_count,
3090 params,
3091 } if *result_count > u32::try_from(MAX_FLAT_RESULTS)? => ResultInfo::Heap {
3092 results: match params.last() {
3093 Some(r) => r.get_u32(),
3094 None => bail_bug!("arg ptr missing"),
3095 },
3096 },
3097 CallerInfo::Sync { result_count, .. } => ResultInfo::Stack {
3098 result_count: *result_count,
3099 },
3100 };
3101
3102 let sync_caller = matches!(caller_info, CallerInfo::Sync { .. });
3103
3104 let start = SendSyncPtr::new(start);
3108 let return_ = SendSyncPtr::new(return_);
3109 let token = StoreToken::new(store.as_context_mut());
3110 let old_thread = store.0.current_guest_thread()?;
3111
3112 let state = store.0.concurrent_state_mut()?;
3113
3114 debug_assert_eq!(
3115 state.get_mut(old_thread.task)?.instance,
3116 self.runtime_instance(caller_instance)
3117 );
3118
3119 let guest_thread = GuestTask::new(
3120 state,
3121 Box::new(move |store, dst| {
3122 let mut store = token.as_context_mut(store);
3123 assert!(dst.len() <= MAX_FLAT_PARAMS);
3124 let mut src = [MaybeUninit::uninit(); MAX_FLAT_PARAMS + 1];
3126 let count = match caller_info {
3127 CallerInfo::Async { params, has_result } => {
3131 let params = ¶ms[..params.len() - usize::from(has_result)];
3132 for (param, src) in params.iter().zip(&mut src) {
3133 src.write(*param);
3134 }
3135 params.len()
3136 }
3137
3138 CallerInfo::Sync { params, .. } => {
3140 for (param, src) in params.iter().zip(&mut src) {
3141 src.write(*param);
3142 }
3143 params.len()
3144 }
3145 };
3146 unsafe {
3153 crate::Func::call_unchecked_raw(
3154 &mut store,
3155 start.as_non_null(),
3156 NonNull::new(
3157 &mut src[..count.max(dst.len())] as *mut [MaybeUninit<ValRaw>] as _,
3158 )
3159 .unwrap(),
3160 )?;
3161 }
3162 dst.copy_from_slice(&src[..dst.len()]);
3163 let task = store.0.current_guest_thread()?.task;
3164 let state = store.0.concurrent_state_mut()?;
3165 Waitable::Guest(task).set_event(
3166 state,
3167 Some(Event::Subtask {
3168 status: Status::Started,
3169 }),
3170 )?;
3171 Ok(())
3172 }),
3173 LiftResult {
3174 lift: Box::new(move |store, src| {
3175 let mut store = token.as_context_mut(store);
3178 let mut my_src = src.to_owned(); if let ResultInfo::Heap { results } = &result_info {
3180 my_src.push(ValRaw::u32(*results));
3181 }
3182
3183 unsafe {
3190 crate::Func::call_unchecked_raw(
3191 &mut store,
3192 return_.as_non_null(),
3193 my_src.as_mut_slice().into(),
3194 )?;
3195 }
3196
3197 let thread = store.0.current_guest_thread()?;
3198 let state = store.0.concurrent_state_mut()?;
3199 if sync_caller {
3200 state.get_mut(thread.task)?.sync_result = SyncResult::Produced(
3201 if let ResultInfo::Stack { result_count } = &result_info {
3202 match result_count {
3203 0 => None,
3204 1 => Some(my_src[0]),
3205 _ => unreachable!(),
3206 }
3207 } else {
3208 None
3209 },
3210 );
3211 }
3212 Ok(Box::new(DummyResult) as Box<dyn Any + Send + Sync>)
3213 }),
3214 ty: task_return_type,
3215 memory: NonNull::new(memory).map(SendSyncPtr::new),
3216 string_encoding,
3217 },
3218 Caller::Guest { thread: old_thread },
3219 None,
3220 self.runtime_instance(callee_instance),
3221 callee_async_typed,
3222 false,
3225 )?;
3226
3227 store.0.set_thread(guest_thread)?;
3230 log::trace!("pushed {guest_thread:?} as current thread; old thread was {old_thread:?}");
3231
3232 Ok(())
3233 }
3234
3235 unsafe fn call_callback<T>(
3240 self,
3241 mut store: StoreContextMut<T>,
3242 function: SendSyncPtr<VMFuncRef>,
3243 event: Event,
3244 handle: u32,
3245 ) -> Result<u32> {
3246 let (ordinal, result) = event.parts();
3247 let params = &mut [
3248 ValRaw::u32(ordinal),
3249 ValRaw::u32(handle),
3250 ValRaw::u32(result),
3251 ];
3252 unsafe {
3257 crate::Func::call_unchecked_raw(
3258 &mut store,
3259 function.as_non_null(),
3260 params.as_mut_slice().into(),
3261 )?;
3262 }
3263 Ok(params[0].get_u32())
3264 }
3265
3266 unsafe fn start_call<T: 'static>(
3279 self,
3280 mut store: StoreContextMut<T>,
3281 callback: *mut VMFuncRef,
3282 post_return: *mut VMFuncRef,
3283 callee: NonNull<VMFuncRef>,
3284 param_count: u32,
3285 result_count: u32,
3286 flags: u32,
3287 storage: Option<&mut [MaybeUninit<ValRaw>]>,
3288 ) -> Result<u32> {
3289 let token = StoreToken::new(store.as_context_mut());
3290 let async_caller = storage.is_none();
3291 let guest_thread = store.0.current_guest_thread()?;
3292 let state = store.0.concurrent_state_mut()?;
3293
3294 if !state.event_loop_running {
3295 bail_bug!("Instance::start_call called without a running event loop");
3296 }
3297
3298 let callee = SendSyncPtr::new(callee);
3299 let param_count = usize::try_from(param_count)?;
3300 assert!(param_count <= MAX_FLAT_PARAMS);
3301 let result_count = usize::try_from(result_count)?;
3302 assert!(result_count <= MAX_FLAT_RESULTS);
3303
3304 let task = state.get_mut(guest_thread.task)?;
3305 let callee_async_typed = task.async_typed;
3306 let callee_instance = task.instance;
3307
3308 task.async_lifted = (flags & START_FLAG_ASYNC_CALLEE) != 0;
3309
3310 if let Some(callback) = NonNull::new(callback) {
3311 let callback = SendSyncPtr::new(callback);
3315 task.callback = Some(Box::new(move |store, event, handle| {
3316 let store = token.as_context_mut(store);
3317 unsafe { self.call_callback::<T>(store, callback, event, handle) }
3318 }));
3319 }
3320
3321 let Caller::Guest { thread: caller } = &task.caller else {
3322 bail_bug!("start_call unexpectedly invoked for host->guest call");
3325 };
3326 let caller = *caller;
3327 let caller_instance = state.get_mut(caller.task)?.instance;
3328
3329 unsafe {
3331 self.stage_call(
3332 store.as_context_mut(),
3333 guest_thread,
3334 callee,
3335 param_count,
3336 result_count,
3337 (flags & START_FLAG_ASYNC_CALLEE) != 0,
3338 NonNull::new(callback).map(SendSyncPtr::new),
3339 NonNull::new(post_return).map(SendSyncPtr::new),
3340 false,
3341 )?;
3342 }
3343
3344 let old_do_not_suspend = if callee_async_typed {
3345 let state = store.0.instance_state(callee_instance).concurrent_state();
3352 let old_do_not_suspend = state.do_not_suspend;
3353 state.do_not_suspend = false;
3354 Some(old_do_not_suspend)
3355 } else {
3356 None
3357 };
3358
3359 let state = store.0.concurrent_state_mut()?;
3360
3361 let guest_waitable = Waitable::Guest(guest_thread.task);
3364 let old_set = guest_waitable.common(state)?.set;
3365 let set = state.get_mut(caller.thread)?.sync_call_set;
3366 guest_waitable.join(state, Some(set))?;
3367
3368 store.0.set_thread(CurrentThread::None)?;
3369
3370 let mut yielded = false;
3386 let (status, waitable) = loop {
3387 store.0.suspend(if yielded {
3388 SuspendReason::Waiting {
3389 set,
3390 thread: caller,
3391 }
3392 } else {
3393 yielded = true;
3394 SuspendReason::YieldingToSubtask { thread: caller }
3395 })?;
3396
3397 if let Some(old_do_not_suspend) = old_do_not_suspend {
3398 store
3399 .0
3400 .instance_state(callee_instance)
3401 .concurrent_state()
3402 .do_not_suspend = old_do_not_suspend;
3403 }
3404
3405 let state = store.0.concurrent_state_mut()?;
3406
3407 log::trace!("taking event for {:?}", guest_thread.task);
3408 let event = guest_waitable.take_event(state)?;
3409 let Some(Event::Subtask { status }) = event else {
3410 bail_bug!("subtasks should only get subtask events, got {event:?}")
3411 };
3412
3413 log::trace!("status {status:?} for {:?}", guest_thread.task);
3414
3415 if status == Status::Returned {
3416 break (status, None);
3418 } else if async_caller {
3419 let handle = store
3423 .0
3424 .instance_state(caller_instance)
3425 .handle_table()
3426 .subtask_insert_guest(guest_thread.task.rep())?;
3427 store
3428 .0
3429 .concurrent_state_mut()?
3430 .get_mut(guest_thread.task)?
3431 .common
3432 .handle = Some(handle);
3433 break (status, Some(handle));
3434 } else {
3435 store.0.switch_or_trap_if_may_not_suspend(caller_instance)?;
3439 }
3440 };
3441
3442 guest_waitable.join(store.0.concurrent_state_mut()?, old_set)?;
3443
3444 store.0.set_thread(caller)?;
3446 store
3447 .0
3448 .concurrent_state_mut()?
3449 .get_mut(caller.thread)?
3450 .state = GuestThreadState::Running;
3451 log::trace!("popped current thread {guest_thread:?}; new thread is {caller:?}");
3452
3453 if let Some(storage) = storage {
3454 let state = store.0.concurrent_state_mut()?;
3458 let task = state.get_mut(guest_thread.task)?;
3459 if let Some(result) = task.sync_result.take()? {
3460 if let Some(result) = result {
3461 storage[0] = MaybeUninit::new(result);
3462 }
3463
3464 if task.exited && task.ready_to_delete() {
3465 Waitable::Guest(guest_thread.task).delete_from(store.0)?;
3466 }
3467 }
3468 }
3469
3470 Ok(status.pack(waitable))
3471 }
3472
3473 pub(crate) fn first_poll<T: 'static, R: Send + 'static>(
3489 self,
3490 mut store: StoreContextMut<'_, T>,
3491 host_task: EnteredHostTask,
3492 future: impl Future<Output = Result<R>> + Send + 'static,
3493 lower: impl FnOnce(StoreContextMut<T>, Option<R>, bool, Option<TableId<HostTask>>) -> Result<()>
3494 + Send
3495 + 'static,
3496 ) -> Result<u32> {
3497 let token = StoreToken::new(store.as_context_mut());
3498
3499 let (join_handle, future) = JoinHandle::run(future);
3502 let mut future = Box::pin(future);
3503
3504 let poll = tls::set(store.0, || {
3509 future
3510 .as_mut()
3511 .poll(&mut Context::from_waker(&Waker::noop()))
3512 });
3513
3514 match poll {
3515 Poll::Ready(result) => {
3517 let result = result.transpose()?;
3518 let task = store.0.current_materialized_host_task()?;
3521 lower(store.as_context_mut(), result, true, task)?;
3522 return Ok(Status::Returned.pack(None));
3523 }
3524
3525 Poll::Pending => {}
3527 }
3528
3529 let Some(task) = store.0.materialize_host_task_id()? else {
3533 bail_bug!("current thread is not a host thread")
3534 };
3535 {
3536 let state = &mut store.0.concurrent_state_mut()?.get_mut(task)?.state;
3537 assert!(matches!(state, HostTaskState::CalleeStarted));
3538 *state = HostTaskState::CalleeRunning(join_handle);
3539 }
3540
3541 let future = Box::pin(async move {
3549 let result = match run_with_host_task_set(task, future).await? {
3550 Some(result) => Some(result?),
3551 None => None,
3552 };
3553 let on_complete = move |store: &mut dyn VMStore| {
3554 let mut store = token.as_context_mut(store);
3558 let old = store.0.set_thread(task)?;
3559
3560 let status = if result.is_some() {
3561 Status::Returned
3562 } else {
3563 Status::ReturnCancelled
3564 };
3565
3566 lower(store.as_context_mut(), result, false, Some(task))?;
3567 let state = store.0.concurrent_state_mut()?;
3568 match &mut state.get_mut(task)?.state {
3569 HostTaskState::CalleeDone { .. } => {}
3572
3573 other => *other = HostTaskState::CalleeDone { cancelled: false },
3575 }
3576 Waitable::Host(task).set_event(state, Some(Event::Subtask { status }))?;
3577
3578 store.0.set_thread(old)?;
3579 Ok(())
3580 };
3581
3582 tls::get(move |store| {
3587 store
3588 .concurrent_state_mut()?
3589 .push_high_priority(WorkItem::WorkerFunction(AlwaysMut::new(Box::new(
3590 on_complete,
3591 ))));
3592 Ok(())
3593 })
3594 });
3595
3596 let caller = match host_task {
3599 Some(caller) => caller,
3600 None => bail_bug!("host task wasn't created but should have been"),
3601 };
3602 let state = store.0.concurrent_state_mut()?;
3603 state.push_future(future);
3604 let instance = state.get_mut(caller.task)?.instance;
3605 let handle = store
3606 .0
3607 .instance_state(instance)
3608 .handle_table()
3609 .subtask_insert_host(task.rep())?;
3610 store.0.concurrent_state_mut()?.get_mut(task)?.common.handle = Some(handle);
3611 log::trace!("assign {task:?} handle {handle} for {caller:?} instance {instance:?}");
3612
3613 store.0.set_thread(caller)?;
3617 Ok(Status::Started.pack(Some(handle)))
3618 }
3619
3620 pub(crate) fn task_return(
3623 self,
3624 store: &mut dyn VMStore,
3625 ty: TypeTupleIndex,
3626 options: OptionsIndex,
3627 storage: &[ValRaw],
3628 ) -> Result<()> {
3629 let guest_thread = store.current_guest_thread()?;
3630 let state = store.concurrent_state_mut()?;
3631 let lift = state
3632 .get_mut(guest_thread.task)?
3633 .lift_result
3634 .take()
3635 .ok_or_else(|| Trap::TaskCancelOrReturnTwice)?;
3636 if !state.get_mut(guest_thread.task)?.result.is_none() {
3637 bail_bug!("task result unexpectedly already set");
3638 }
3639
3640 let CanonicalOptions {
3641 string_encoding,
3642 data_model,
3643 ..
3644 } = &self.id().get(store).component().env_component().options[options];
3645
3646 let invalid = ty != lift.ty
3647 || string_encoding != &lift.string_encoding
3648 || match data_model {
3649 CanonicalOptionsDataModel::LinearMemory(opts) => match opts.memory {
3650 Some(memory) => {
3651 let expected = lift.memory.map(|v| v.as_ptr()).unwrap_or(ptr::null_mut());
3652 let actual = self.id().get(store).runtime_memory(memory);
3653 expected != actual.as_ptr()
3654 }
3655 None => false,
3658 },
3659 CanonicalOptionsDataModel::Gc { .. } => true,
3661 };
3662
3663 if invalid {
3664 bail!(Trap::TaskReturnInvalid);
3665 }
3666
3667 log::trace!("task.return for {guest_thread:?}");
3668
3669 let result = (lift.lift)(store, storage)?;
3670 self.task_complete(store, guest_thread.task, result, Status::Returned)
3671 }
3672
3673 pub(crate) fn task_cancel(self, store: &mut StoreOpaque) -> Result<()> {
3675 let guest_thread = store.current_guest_thread()?;
3676 let state = store.concurrent_state_mut()?;
3677 let task = state.get_mut(guest_thread.task)?;
3678 if !task.cancel_request_delivered {
3679 bail!(Trap::TaskCancelNotCancelled);
3680 }
3681 _ = task
3682 .lift_result
3683 .take()
3684 .ok_or_else(|| Trap::TaskCancelOrReturnTwice)?;
3685
3686 if !task.result.is_none() {
3687 bail_bug!("task result should not bet set yet");
3688 }
3689
3690 log::trace!("task.cancel for {guest_thread:?}");
3691
3692 self.task_complete(
3693 store,
3694 guest_thread.task,
3695 Box::new(DummyResult),
3696 Status::ReturnCancelled,
3697 )
3698 }
3699
3700 fn task_complete(
3706 self,
3707 store: &mut StoreOpaque,
3708 guest_task: TableId<GuestTask>,
3709 result: Box<dyn Any + Send + Sync>,
3710 status: Status,
3711 ) -> Result<()> {
3712 store
3713 .component_resource_tables(Some(self))?
3714 .validate_scope_exit()?;
3715
3716 let state = store.concurrent_state_mut()?;
3717 let task = state.get_mut(guest_task)?;
3718
3719 if let Caller::Host { tx, .. } = &mut task.caller {
3720 if let Some(tx) = tx.take() {
3721 _ = tx.send(result);
3722 }
3723 } else {
3724 task.result = Some(result);
3725 Waitable::Guest(guest_task).set_event(state, Some(Event::Subtask { status }))?;
3726 }
3727
3728 Ok(())
3729 }
3730
3731 pub(crate) fn waitable_set_new(
3733 self,
3734 store: &mut StoreOpaque,
3735 caller_instance: RuntimeComponentInstanceIndex,
3736 ) -> Result<u32> {
3737 let set = store.concurrent_state_mut()?.push(WaitableSet::default())?;
3738 let handle = store
3739 .instance_state(self.runtime_instance(caller_instance))
3740 .handle_table()
3741 .waitable_set_insert(set.rep())?;
3742 log::trace!("new waitable set {set:?} (handle {handle})");
3743 Ok(handle)
3744 }
3745
3746 pub(crate) fn waitable_set_drop(
3748 self,
3749 store: &mut StoreOpaque,
3750 caller_instance: RuntimeComponentInstanceIndex,
3751 set: u32,
3752 ) -> Result<()> {
3753 let rep = store
3754 .instance_state(self.runtime_instance(caller_instance))
3755 .handle_table()
3756 .waitable_set_remove(set)?;
3757
3758 log::trace!("drop waitable set {rep} (handle {set})");
3759
3760 if !store
3764 .concurrent_state_mut()?
3765 .get_mut(TableId::<WaitableSet>::new(rep))?
3766 .waiting
3767 .is_empty()
3768 {
3769 bail!(Trap::WaitableSetDropHasWaiters);
3770 }
3771
3772 store
3773 .concurrent_state_mut()?
3774 .delete(TableId::<WaitableSet>::new(rep))?;
3775
3776 Ok(())
3777 }
3778
3779 pub(crate) fn waitable_join(
3781 self,
3782 store: &mut StoreOpaque,
3783 caller_instance: RuntimeComponentInstanceIndex,
3784 waitable_handle: u32,
3785 set_handle: u32,
3786 ) -> Result<()> {
3787 let mut instance = self.id().get_mut(store);
3788 let waitable =
3789 Waitable::from_instance(instance.as_mut(), caller_instance, waitable_handle)?;
3790
3791 let set = if set_handle == 0 {
3792 None
3793 } else {
3794 let set = instance.instance_states().0[caller_instance]
3795 .handle_table()
3796 .waitable_set_rep(set_handle)?;
3797
3798 let state = store.concurrent_state_mut()?;
3799 if let Some(old) = waitable.common(state)?.set
3800 && state.get_mut(old)?.is_sync_call_set
3801 {
3802 bail!(Trap::WaitableSyncAndAsync);
3803 }
3804
3805 Some(TableId::<WaitableSet>::new(set))
3806 };
3807
3808 log::trace!(
3809 "waitable {waitable:?} (handle {waitable_handle}) join set {set:?} (handle {set_handle})",
3810 );
3811
3812 waitable.join(store.concurrent_state_mut()?, set)
3813 }
3814
3815 pub(crate) fn subtask_drop(
3817 self,
3818 store: &mut StoreOpaque,
3819 caller_instance: RuntimeComponentInstanceIndex,
3820 task_id: u32,
3821 ) -> Result<()> {
3822 self.waitable_join(store, caller_instance, task_id, 0)?;
3823
3824 let (rep, is_host) = store
3825 .instance_state(self.runtime_instance(caller_instance))
3826 .handle_table()
3827 .subtask_remove(task_id)?;
3828
3829 let concurrent_state = store.concurrent_state_mut()?;
3830 let (waitable, delete) = if is_host {
3831 let id = TableId::<HostTask>::new(rep);
3832 let task = concurrent_state.get_mut(id)?;
3833 match &task.state {
3834 HostTaskState::CalleeRunning(_) => bail!(Trap::SubtaskDropNotResolved),
3835 HostTaskState::CalleeDone { .. } => {}
3836 HostTaskState::CalleeStarted | HostTaskState::CalleeFinished(_) => {
3837 bail_bug!("invalid state for callee in `subtask.drop`")
3838 }
3839 }
3840
3841 (Waitable::Host(id), true)
3842 } else {
3843 let id = TableId::<GuestTask>::new(rep);
3844 let task = concurrent_state.get_mut(id)?;
3845 if task.lift_result.is_some() {
3846 bail!(Trap::SubtaskDropNotResolved);
3847 }
3848 (
3849 Waitable::Guest(id),
3850 concurrent_state.get_mut(id)?.ready_to_delete(),
3851 )
3852 };
3853
3854 waitable.common(concurrent_state)?.handle = None;
3855
3856 if waitable.take_event(concurrent_state)?.is_some() {
3859 bail!(Trap::SubtaskDropNotResolved);
3860 }
3861
3862 if delete {
3863 waitable.delete_from(store)?;
3864 }
3865
3866 log::trace!("subtask_drop {waitable:?} (handle {task_id})");
3867 Ok(())
3868 }
3869
3870 pub(crate) fn waitable_set_wait(
3872 self,
3873 store: &mut StoreOpaque,
3874 options: OptionsIndex,
3875 set: u32,
3876 payload: u32,
3877 ) -> Result<u32> {
3878 let &CanonicalOptions {
3879 instance: caller_instance,
3880 ..
3881 } = &self.id().get(store).component().env_component().options[options];
3882 let caller = self.runtime_instance(caller_instance);
3883 let rep = store
3884 .instance_state(self.runtime_instance(caller_instance))
3885 .handle_table()
3886 .waitable_set_rep(set)?;
3887
3888 self.waitable_check(
3889 store,
3890 caller,
3891 WaitableCheck::Wait,
3892 WaitableCheckParams {
3893 set: TableId::new(rep),
3894 options,
3895 payload,
3896 },
3897 )
3898 }
3899
3900 pub(crate) fn waitable_set_poll(
3902 self,
3903 store: &mut StoreOpaque,
3904 options: OptionsIndex,
3905 set: u32,
3906 payload: u32,
3907 ) -> Result<u32> {
3908 let &CanonicalOptions {
3909 instance: caller_instance,
3910 ..
3911 } = &self.id().get(store).component().env_component().options[options];
3912 let caller = self.runtime_instance(caller_instance);
3913 let rep = store
3914 .instance_state(caller)
3915 .handle_table()
3916 .waitable_set_rep(set)?;
3917
3918 self.waitable_check(
3919 store,
3920 caller,
3921 WaitableCheck::Poll,
3922 WaitableCheckParams {
3923 set: TableId::new(rep),
3924 options,
3925 payload,
3926 },
3927 )
3928 }
3929
3930 pub(crate) fn thread_index(&self, store: &mut dyn VMStore) -> Result<u32> {
3932 let thread_id = store.current_guest_thread()?.thread;
3933 match store
3934 .concurrent_state_mut()?
3935 .get_mut(thread_id)?
3936 .instance_rep
3937 {
3938 Some(r) => Ok(r),
3939 None => bail_bug!("thread should have instance_rep by now"),
3940 }
3941 }
3942
3943 pub(crate) fn thread_new_indirect<T: 'static>(
3945 self,
3946 mut store: StoreContextMut<T>,
3947 runtime_instance: RuntimeComponentInstanceIndex,
3948 _func_ty_idx: TypeFuncIndex, start_func_table_idx: RuntimeTableIndex,
3950 start_func_idx: u32,
3951 context: i32,
3952 ) -> Result<u32> {
3953 log::trace!("creating new thread");
3954
3955 let start_func_ty = FuncType::new(store.engine(), [ValType::I32], []);
3956 let (instance, registry) = self.id().get_mut_and_registry(store.0);
3957 let callee = instance
3958 .index_runtime_func_table(registry, start_func_table_idx, start_func_idx as u64)?
3959 .ok_or_else(|| Trap::ThreadNewIndirectUninitialized)?;
3960 if callee.type_index(store.0) != start_func_ty.type_index() {
3961 bail!(Trap::ThreadNewIndirectInvalidType);
3962 }
3963
3964 let token = StoreToken::new(store.as_context_mut());
3965 let start_func = Box::new(
3966 move |store: &mut dyn VMStore, guest_thread: QualifiedThreadId| -> Result<()> {
3967 let old_thread = store.set_thread(guest_thread)?;
3968 log::trace!(
3969 "thread start: replaced {old_thread:?} with {guest_thread:?} as current thread"
3970 );
3971
3972 let mut store = token.as_context_mut(store);
3973 let mut params = [ValRaw::i32(context)];
3974 unsafe { callee.call_unchecked(store.as_context_mut(), &mut params)? };
3977
3978 store.0.set_thread(old_thread)?;
3979
3980 let runtime_instance = self.runtime_instance(runtime_instance);
3981
3982 store
3985 .0
3986 .switch_or_trap_if_may_not_suspend(runtime_instance)?;
3987
3988 store
3989 .0
3990 .cleanup_thread(guest_thread, runtime_instance, CleanupTask::Yes)?;
3991
3992 log::trace!("explicit thread {guest_thread:?} completed");
3993 let state = store.0.concurrent_state_mut()?;
3994 if let Some(t) = old_thread.guest() {
3995 state.get_mut(t.thread)?.state = GuestThreadState::Running;
3996 }
3997 log::trace!("thread start: restored {old_thread:?} as current thread");
3998
3999 Ok(())
4000 },
4001 );
4002
4003 let current_thread = store.0.current_guest_thread()?;
4004 let state = store.0.concurrent_state_mut()?;
4005 let parent_task = current_thread.task;
4006
4007 let new_thread = GuestThread::new_explicit(state, parent_task, start_func)?;
4008 let thread_id = state.push(new_thread)?;
4009 state.get_mut(parent_task)?.threads.insert(thread_id);
4010
4011 log::trace!("new thread with id {thread_id:?} created");
4012
4013 self.add_guest_thread_to_instance_table(thread_id, store.0, runtime_instance)
4014 }
4015
4016 pub(crate) fn resume_thread(
4017 self,
4018 store: &mut StoreOpaque,
4019 runtime_instance: RuntimeComponentInstanceIndex,
4020 thread_idx: u32,
4021 how: ResumeThread,
4022 ) -> Result<bool> {
4023 let thread_id =
4024 GuestThread::from_instance(self.id().get_mut(store), runtime_instance, thread_idx)?;
4025 let state = store.concurrent_state_mut()?;
4026 let guest_thread = QualifiedThreadId::qualify(state, thread_id)?;
4027
4028 if store.current_guest_thread()? == guest_thread {
4029 bail!(Trap::CannotResumeThread);
4030 }
4031
4032 let state = store.concurrent_state_mut()?;
4033 let thread = state.get_mut(guest_thread.thread)?;
4034 let priority = match how {
4035 ResumeThread::Promote | ResumeThread::Resume => Priority::Switch,
4036 ResumeThread::ResumeLater => Priority::Low,
4037 };
4038
4039 match (&how, &thread.state) {
4040 (ResumeThread::Promote, GuestThreadState::Ready { .. }) => {}
4042 (ResumeThread::Promote, _) => return Ok(false),
4043
4044 (
4047 ResumeThread::Resume | ResumeThread::ResumeLater,
4048 GuestThreadState::NotStartedExplicit(_) | GuestThreadState::Suspended(_),
4049 ) => {}
4050 (ResumeThread::Resume | ResumeThread::ResumeLater, _) => {
4051 bail!(Trap::CannotResumeThread)
4052 }
4053 }
4054
4055 match mem::replace(&mut thread.state, GuestThreadState::Running) {
4056 GuestThreadState::NotStartedExplicit(start_func) => {
4057 log::trace!("starting thread {guest_thread:?}");
4058 let guest_call = WorkItem::GuestCall {
4059 instance: self.runtime_instance(runtime_instance),
4060 call: GuestCall {
4061 thread: guest_thread,
4062 kind: GuestCallKind::StartExplicit(Box::new(move |store| {
4063 start_func(store, guest_thread)
4064 })),
4065 },
4066 };
4067 store
4068 .concurrent_state_mut()?
4069 .push_work_item(guest_call, priority)?;
4070 }
4071 GuestThreadState::Suspended(fiber) => {
4072 log::trace!("resuming thread {thread_id:?} that was suspended");
4073 store.concurrent_state_mut()?.push_work_item(
4074 WorkItem::ResumeFiber {
4075 instance: self.runtime_instance(runtime_instance),
4076 thread: guest_thread,
4077 fiber,
4078 },
4079 priority,
4080 )?;
4081 }
4082 GuestThreadState::Ready { fiber } => {
4083 log::trace!("resuming thread {thread_id:?} that was ready");
4084 thread.state = GuestThreadState::Ready { fiber };
4085 store
4086 .concurrent_state_mut()?
4087 .promote_thread_work_item(guest_thread)?;
4088 }
4089 other @ (GuestThreadState::NotStartedImplicit
4090 | GuestThreadState::Running
4091 | GuestThreadState::Completed) => {
4092 thread.state = other;
4093 }
4094 }
4095 Ok(true)
4096 }
4097
4098 fn add_guest_thread_to_instance_table(
4099 self,
4100 thread_id: TableId<GuestThread>,
4101 store: &mut StoreOpaque,
4102 runtime_instance: RuntimeComponentInstanceIndex,
4103 ) -> Result<u32> {
4104 let guest_id = store
4105 .instance_state(self.runtime_instance(runtime_instance))
4106 .thread_handle_table()
4107 .guest_thread_insert(thread_id.rep())?;
4108 store
4109 .concurrent_state_mut()?
4110 .get_mut(thread_id)?
4111 .instance_rep = Some(guest_id);
4112 Ok(guest_id)
4113 }
4114
4115 pub(crate) fn suspension_intrinsic(
4119 self,
4120 store: &mut StoreOpaque,
4121 caller: RuntimeComponentInstanceIndex,
4122 yielding: bool,
4123 to_thread: SuspensionTarget,
4124 ) -> Result<WaitResult> {
4125 let check_suspend = match to_thread {
4126 SuspensionTarget::Promote(thread) => {
4127 !self.resume_thread(store, caller, thread, ResumeThread::Promote)?
4128 }
4129 SuspensionTarget::Resume(thread) => {
4130 if !self.resume_thread(store, caller, thread, ResumeThread::Resume)? {
4131 bail_bug!(
4132 "`resume_thread` should only ever return false \
4133 when `ResumeThread::Promote` is passed to it"
4134 );
4135 }
4136 false
4137 }
4138 SuspensionTarget::None => true,
4139 };
4140
4141 if check_suspend && !store.switch_if_may_not_suspend(self.runtime_instance(caller))? {
4142 return if yielding {
4143 Ok(WaitResult::Completed)
4144 } else {
4145 Err(Trap::CannotBlockSyncTask.into())
4146 };
4147 }
4148
4149 let guest_thread = store.current_guest_thread()?;
4150
4151 let reason = if yielding {
4152 SuspendReason::Yielding {
4153 thread: guest_thread,
4154 }
4155 } else {
4156 SuspendReason::ExplicitlySuspending {
4157 thread: guest_thread,
4158 }
4159 };
4160
4161 store.suspend(reason)?;
4162
4163 Ok(WaitResult::Completed)
4164 }
4165
4166 fn waitable_check(
4168 self,
4169 store: &mut StoreOpaque,
4170 caller: RuntimeInstance,
4171 check: WaitableCheck,
4172 params: WaitableCheckParams,
4173 ) -> Result<u32> {
4174 let guest_thread = store.current_guest_thread()?;
4175
4176 log::trace!("waitable check for {guest_thread:?}; set {:?}", params.set);
4177
4178 let state = store.concurrent_state_mut()?;
4179 let task = state.get_mut(guest_thread.task)?;
4180
4181 match &check {
4184 WaitableCheck::Wait => {
4185 let set = params.set;
4186
4187 if (task.event.is_none() || matches!(task.event, Some(Event::Cancelled)))
4188 && state.get_mut(set)?.ready.is_empty()
4189 {
4190 store.switch_or_trap_if_may_not_suspend(caller)?;
4191
4192 store.suspend(SuspendReason::Waiting {
4193 set,
4194 thread: guest_thread,
4195 })?;
4196 }
4197 }
4198 WaitableCheck::Poll => {}
4199 }
4200
4201 log::trace!(
4202 "waitable check for {guest_thread:?}; set {:?}, part two",
4203 params.set
4204 );
4205
4206 let event = self.get_event(store, guest_thread.task, Some(params.set), false)?;
4208
4209 let (ordinal, handle, result) = match &check {
4210 WaitableCheck::Wait => {
4211 let (event, waitable) = match event {
4212 Some(p) => p,
4213 None => bail_bug!("event expected to be present"),
4214 };
4215 let handle = waitable.map(|(_, v)| v).unwrap_or(0);
4216 let (ordinal, result) = event.parts();
4217 (ordinal, handle, result)
4218 }
4219 WaitableCheck::Poll => {
4220 if let Some((event, waitable)) = event {
4221 let handle = waitable.map(|(_, v)| v).unwrap_or(0);
4222 let (ordinal, result) = event.parts();
4223 (ordinal, handle, result)
4224 } else {
4225 log::trace!(
4226 "no events ready to deliver via waitable-set.poll to {:?}; set {:?}",
4227 guest_thread.task,
4228 params.set
4229 );
4230 let (ordinal, result) = Event::None.parts();
4231 (ordinal, 0, result)
4232 }
4233 }
4234 };
4235 let memory = self.options_memory_mut(store, params.options);
4236 let ptr = crate::component::func::validate_inbounds_dynamic(
4237 &CanonicalAbiInfo::POINTER_PAIR,
4238 memory,
4239 &ValRaw::u32(params.payload),
4240 )?;
4241 memory[ptr + 0..][..4].copy_from_slice(&handle.to_le_bytes());
4242 memory[ptr + 4..][..4].copy_from_slice(&result.to_le_bytes());
4243 Ok(ordinal)
4244 }
4245
4246 pub(crate) fn subtask_cancel(
4248 self,
4249 store: &mut StoreOpaque,
4250 caller_instance: RuntimeComponentInstanceIndex,
4251 async_: bool,
4252 task_id: u32,
4253 ) -> Result<u32> {
4254 let (rep, is_host) = store
4255 .instance_state(self.runtime_instance(caller_instance))
4256 .handle_table()
4257 .subtask_rep(task_id)?;
4258 let waitable = if is_host {
4259 Waitable::Host(TableId::<HostTask>::new(rep))
4260 } else {
4261 Waitable::Guest(TableId::<GuestTask>::new(rep))
4262 };
4263 let concurrent_state = store.concurrent_state_mut()?;
4264
4265 log::trace!("subtask_cancel {waitable:?} (handle {task_id}; async {async_})");
4266
4267 waitable.trap_if_in_waitable_set(concurrent_state)?;
4268
4269 let needs_block;
4270 if let Waitable::Host(host_task) = waitable {
4271 let state = &mut concurrent_state.get_mut(host_task)?.state;
4272 match mem::replace(state, HostTaskState::CalleeDone { cancelled: true }) {
4273 HostTaskState::CalleeRunning(handle) => {
4280 handle.abort();
4281 needs_block = true;
4282 }
4283
4284 HostTaskState::CalleeDone { cancelled } => {
4287 if cancelled {
4288 bail!(Trap::SubtaskCancelAfterTerminal);
4289 } else {
4290 needs_block = false;
4293 }
4294 }
4295
4296 HostTaskState::CalleeStarted | HostTaskState::CalleeFinished(_) => {
4299 bail_bug!("invalid states for host callee")
4300 }
4301 }
4302 } else {
4303 let guest_task = TableId::<GuestTask>::new(rep);
4304 let task = concurrent_state.get_mut(guest_task)?;
4305 if !task.already_lowered_parameters() {
4306 store.cancel_guest_subtask_without_lowered_parameters(
4307 self.runtime_instance(caller_instance),
4308 guest_task,
4309 )?;
4310 return Ok(Status::StartCancelled as u32);
4311 } else if !task.returned_or_cancelled() {
4312 task.event = Some(Event::Cancelled);
4320 let runtime_instance = task.instance;
4321 for thread in task.threads.clone() {
4322 let thread = QualifiedThreadId {
4323 task: guest_task,
4324 thread,
4325 };
4326 let thread_mut = concurrent_state.get_mut(thread.thread)?;
4327
4328 let yield_ = |store: &mut StoreOpaque| {
4329 let state = store.instance_state(runtime_instance).concurrent_state();
4334 let old_do_not_suspend = state.do_not_suspend;
4335 state.do_not_suspend = false;
4336
4337 let caller = store.current_guest_thread()?;
4338
4339 let state = store.concurrent_state_mut()?;
4344 let set = state.get_mut(caller.thread)?.sync_call_set;
4345 waitable.join(state, Some(set))?;
4346
4347 store.suspend(SuspendReason::YieldingToSubtask { thread: caller })?;
4348
4349 let state = store.concurrent_state_mut()?;
4350 waitable.join(state, None)?;
4351
4352 store
4353 .instance_state(runtime_instance)
4354 .concurrent_state()
4355 .do_not_suspend = old_do_not_suspend;
4356
4357 Ok::<(), crate::Error>(())
4358 };
4359
4360 match thread_mut.wake_on_cancel.take() {
4361 WakeOnCancel::Waiting(set) => {
4362 let item = match concurrent_state.get_mut(set)?.waiting.remove(&thread)
4364 {
4365 Some(WaitMode::Callback(instance)) => WorkItem::GuestCall {
4366 instance: runtime_instance,
4367 call: GuestCall {
4368 thread,
4369 kind: GuestCallKind::DeliverEvent {
4370 instance,
4371 set: None,
4372 },
4373 },
4374 },
4375 other => bail_bug!(
4376 "expected `Some(WaitMode::Callback(_))`; got `{other:?}`"
4377 ),
4378 };
4379 concurrent_state.set_switch_item(item)?;
4380
4381 yield_(store)?;
4382
4383 break;
4384 }
4385 WakeOnCancel::Yielding => {
4386 if concurrent_state.promote_thread_work_item(thread)? {
4387 yield_(store)?;
4388 break;
4389 } else {
4390 bail_bug!("thread with `WakeOnCancel::Yielding` not promotable");
4391 }
4392 }
4393 WakeOnCancel::None => {}
4394 }
4395 }
4396
4397 needs_block = !store
4400 .concurrent_state_mut()?
4401 .get_mut(guest_task)?
4402 .returned_or_cancelled()
4403 } else {
4404 needs_block = false;
4405 }
4406 };
4407
4408 if needs_block {
4412 if async_ {
4413 return Ok(BLOCKED);
4414 }
4415
4416 let old_next_switch_item = {
4419 let state = store.concurrent_state_mut()?;
4420 let item = state.next_switch_item.take();
4421 state.push(item)?
4425 };
4426
4427 store.wait_for_event(self.runtime_instance(caller_instance), waitable)?;
4430
4431 let state = store.concurrent_state_mut()?;
4432 state.next_switch_item = state.delete(old_next_switch_item)?;
4433
4434 }
4436
4437 let event = waitable.take_event(store.concurrent_state_mut()?)?;
4438 if let Some(Event::Subtask {
4439 status: status @ (Status::Returned | Status::ReturnCancelled),
4440 }) = event
4441 {
4442 Ok(status as u32)
4443 } else {
4444 bail!(Trap::SubtaskCancelAfterTerminal);
4445 }
4446 }
4447}
4448
4449pub trait VMComponentAsyncStore {
4457 unsafe fn prepare_call(
4463 &mut self,
4464 instance: Instance,
4465 memory: *mut VMMemoryDefinition,
4466 start: NonNull<VMFuncRef>,
4467 return_: NonNull<VMFuncRef>,
4468 caller_instance: RuntimeComponentInstanceIndex,
4469 callee_instance: RuntimeComponentInstanceIndex,
4470 task_return_type: TypeTupleIndex,
4471 callee_async: bool,
4472 string_encoding: StringEncoding,
4473 result_count: u32,
4474 storage: *mut ValRaw,
4475 storage_len: usize,
4476 ) -> Result<()>;
4477
4478 unsafe fn sync_start(
4481 &mut self,
4482 instance: Instance,
4483 callback: *mut VMFuncRef,
4484 callee: NonNull<VMFuncRef>,
4485 param_count: u32,
4486 storage: *mut MaybeUninit<ValRaw>,
4487 storage_len: usize,
4488 ) -> Result<()>;
4489
4490 unsafe fn async_start(
4493 &mut self,
4494 instance: Instance,
4495 callback: *mut VMFuncRef,
4496 post_return: *mut VMFuncRef,
4497 callee: NonNull<VMFuncRef>,
4498 param_count: u32,
4499 result_count: u32,
4500 flags: u32,
4501 ) -> Result<u32>;
4502
4503 fn future_write(
4505 &mut self,
4506 instance: Instance,
4507 caller: RuntimeComponentInstanceIndex,
4508 ty: TypeFutureTableIndex,
4509 options: OptionsIndex,
4510 future: u32,
4511 address: u32,
4512 ) -> Result<u32>;
4513
4514 fn future_read(
4516 &mut self,
4517 instance: Instance,
4518 caller: RuntimeComponentInstanceIndex,
4519 ty: TypeFutureTableIndex,
4520 options: OptionsIndex,
4521 future: u32,
4522 address: u32,
4523 ) -> Result<u32>;
4524
4525 fn future_drop_writable(
4527 &mut self,
4528 instance: Instance,
4529 ty: TypeFutureTableIndex,
4530 writer: u32,
4531 ) -> Result<()>;
4532
4533 fn stream_write(
4535 &mut self,
4536 instance: Instance,
4537 caller: RuntimeComponentInstanceIndex,
4538 ty: TypeStreamTableIndex,
4539 options: OptionsIndex,
4540 stream: u32,
4541 address: u32,
4542 count: u32,
4543 ) -> Result<u32>;
4544
4545 fn stream_read(
4547 &mut self,
4548 instance: Instance,
4549 caller: RuntimeComponentInstanceIndex,
4550 ty: TypeStreamTableIndex,
4551 options: OptionsIndex,
4552 stream: u32,
4553 address: u32,
4554 count: u32,
4555 ) -> Result<u32>;
4556
4557 fn flat_stream_write(
4560 &mut self,
4561 instance: Instance,
4562 caller: RuntimeComponentInstanceIndex,
4563 ty: TypeStreamTableIndex,
4564 options: OptionsIndex,
4565 payload_size: u32,
4566 payload_align: u32,
4567 stream: u32,
4568 address: u32,
4569 count: u32,
4570 ) -> Result<u32>;
4571
4572 fn flat_stream_read(
4575 &mut self,
4576 instance: Instance,
4577 caller: RuntimeComponentInstanceIndex,
4578 ty: TypeStreamTableIndex,
4579 options: OptionsIndex,
4580 payload_size: u32,
4581 payload_align: u32,
4582 stream: u32,
4583 address: u32,
4584 count: u32,
4585 ) -> Result<u32>;
4586
4587 fn stream_drop_writable(
4589 &mut self,
4590 instance: Instance,
4591 ty: TypeStreamTableIndex,
4592 writer: u32,
4593 ) -> Result<()>;
4594
4595 fn error_context_debug_message(
4597 &mut self,
4598 instance: Instance,
4599 ty: TypeComponentLocalErrorContextTableIndex,
4600 options: OptionsIndex,
4601 err_ctx_handle: u32,
4602 debug_msg_address: u32,
4603 ) -> Result<()>;
4604
4605 fn thread_new_indirect(
4607 &mut self,
4608 instance: Instance,
4609 caller: RuntimeComponentInstanceIndex,
4610 func_ty_idx: TypeFuncIndex,
4611 start_func_table_idx: RuntimeTableIndex,
4612 start_func_idx: u32,
4613 context: i32,
4614 ) -> Result<u32>;
4615}
4616
4617impl<T: 'static> VMComponentAsyncStore for StoreInner<T> {
4619 unsafe fn prepare_call(
4620 &mut self,
4621 instance: Instance,
4622 memory: *mut VMMemoryDefinition,
4623 start: NonNull<VMFuncRef>,
4624 return_: NonNull<VMFuncRef>,
4625 caller_instance: RuntimeComponentInstanceIndex,
4626 callee_instance: RuntimeComponentInstanceIndex,
4627 task_return_type: TypeTupleIndex,
4628 callee_async: bool,
4629 string_encoding: StringEncoding,
4630 result_count_or_max_if_async: u32,
4631 storage: *mut ValRaw,
4632 storage_len: usize,
4633 ) -> Result<()> {
4634 let params = unsafe { core::slice::from_raw_parts(storage, storage_len) }.to_vec();
4638
4639 unsafe {
4640 instance.prepare_call(
4641 StoreContextMut(self),
4642 start,
4643 return_,
4644 caller_instance,
4645 callee_instance,
4646 task_return_type,
4647 callee_async,
4648 memory,
4649 string_encoding,
4650 match result_count_or_max_if_async {
4651 PREPARE_ASYNC_NO_RESULT => CallerInfo::Async {
4652 params,
4653 has_result: false,
4654 },
4655 PREPARE_ASYNC_WITH_RESULT => CallerInfo::Async {
4656 params,
4657 has_result: true,
4658 },
4659 result_count => CallerInfo::Sync {
4660 params,
4661 result_count,
4662 },
4663 },
4664 )
4665 }
4666 }
4667
4668 unsafe fn sync_start(
4669 &mut self,
4670 instance: Instance,
4671 callback: *mut VMFuncRef,
4672 callee: NonNull<VMFuncRef>,
4673 param_count: u32,
4674 storage: *mut MaybeUninit<ValRaw>,
4675 storage_len: usize,
4676 ) -> Result<()> {
4677 unsafe {
4678 instance
4679 .start_call(
4680 StoreContextMut(self),
4681 callback,
4682 ptr::null_mut(),
4683 callee,
4684 param_count,
4685 1,
4686 START_FLAG_ASYNC_CALLEE,
4687 Some(core::slice::from_raw_parts_mut(storage, storage_len)),
4691 )
4692 .map(drop)
4693 }
4694 }
4695
4696 unsafe fn async_start(
4697 &mut self,
4698 instance: Instance,
4699 callback: *mut VMFuncRef,
4700 post_return: *mut VMFuncRef,
4701 callee: NonNull<VMFuncRef>,
4702 param_count: u32,
4703 result_count: u32,
4704 flags: u32,
4705 ) -> Result<u32> {
4706 unsafe {
4707 instance.start_call(
4708 StoreContextMut(self),
4709 callback,
4710 post_return,
4711 callee,
4712 param_count,
4713 result_count,
4714 flags,
4715 None,
4716 )
4717 }
4718 }
4719
4720 fn future_write(
4721 &mut self,
4722 instance: Instance,
4723 caller: RuntimeComponentInstanceIndex,
4724 ty: TypeFutureTableIndex,
4725 options: OptionsIndex,
4726 future: u32,
4727 address: u32,
4728 ) -> Result<u32> {
4729 instance
4730 .guest_write(
4731 StoreContextMut(self),
4732 caller,
4733 TransmitIndex::Future(ty),
4734 options,
4735 None,
4736 future,
4737 address,
4738 1,
4739 )
4740 .map(|result| result.encode())
4741 }
4742
4743 fn future_read(
4744 &mut self,
4745 instance: Instance,
4746 caller: RuntimeComponentInstanceIndex,
4747 ty: TypeFutureTableIndex,
4748 options: OptionsIndex,
4749 future: u32,
4750 address: u32,
4751 ) -> Result<u32> {
4752 instance
4753 .guest_read(
4754 StoreContextMut(self),
4755 caller,
4756 TransmitIndex::Future(ty),
4757 options,
4758 None,
4759 future,
4760 address,
4761 1,
4762 )
4763 .map(|result| result.encode())
4764 }
4765
4766 fn stream_write(
4767 &mut self,
4768 instance: Instance,
4769 caller: RuntimeComponentInstanceIndex,
4770 ty: TypeStreamTableIndex,
4771 options: OptionsIndex,
4772 stream: u32,
4773 address: u32,
4774 count: u32,
4775 ) -> Result<u32> {
4776 instance
4777 .guest_write(
4778 StoreContextMut(self),
4779 caller,
4780 TransmitIndex::Stream(ty),
4781 options,
4782 None,
4783 stream,
4784 address,
4785 count,
4786 )
4787 .map(|result| result.encode())
4788 }
4789
4790 fn stream_read(
4791 &mut self,
4792 instance: Instance,
4793 caller: RuntimeComponentInstanceIndex,
4794 ty: TypeStreamTableIndex,
4795 options: OptionsIndex,
4796 stream: u32,
4797 address: u32,
4798 count: u32,
4799 ) -> Result<u32> {
4800 instance
4801 .guest_read(
4802 StoreContextMut(self),
4803 caller,
4804 TransmitIndex::Stream(ty),
4805 options,
4806 None,
4807 stream,
4808 address,
4809 count,
4810 )
4811 .map(|result| result.encode())
4812 }
4813
4814 fn future_drop_writable(
4815 &mut self,
4816 instance: Instance,
4817 ty: TypeFutureTableIndex,
4818 writer: u32,
4819 ) -> Result<()> {
4820 instance.guest_drop_writable(self, TransmitIndex::Future(ty), writer)
4821 }
4822
4823 fn flat_stream_write(
4824 &mut self,
4825 instance: Instance,
4826 caller: RuntimeComponentInstanceIndex,
4827 ty: TypeStreamTableIndex,
4828 options: OptionsIndex,
4829 payload_size: u32,
4830 payload_align: u32,
4831 stream: u32,
4832 address: u32,
4833 count: u32,
4834 ) -> Result<u32> {
4835 instance
4836 .guest_write(
4837 StoreContextMut(self),
4838 caller,
4839 TransmitIndex::Stream(ty),
4840 options,
4841 Some(FlatAbi {
4842 size: payload_size,
4843 align: payload_align,
4844 }),
4845 stream,
4846 address,
4847 count,
4848 )
4849 .map(|result| result.encode())
4850 }
4851
4852 fn flat_stream_read(
4853 &mut self,
4854 instance: Instance,
4855 caller: RuntimeComponentInstanceIndex,
4856 ty: TypeStreamTableIndex,
4857 options: OptionsIndex,
4858 payload_size: u32,
4859 payload_align: u32,
4860 stream: u32,
4861 address: u32,
4862 count: u32,
4863 ) -> Result<u32> {
4864 instance
4865 .guest_read(
4866 StoreContextMut(self),
4867 caller,
4868 TransmitIndex::Stream(ty),
4869 options,
4870 Some(FlatAbi {
4871 size: payload_size,
4872 align: payload_align,
4873 }),
4874 stream,
4875 address,
4876 count,
4877 )
4878 .map(|result| result.encode())
4879 }
4880
4881 fn stream_drop_writable(
4882 &mut self,
4883 instance: Instance,
4884 ty: TypeStreamTableIndex,
4885 writer: u32,
4886 ) -> Result<()> {
4887 instance.guest_drop_writable(self, TransmitIndex::Stream(ty), writer)
4888 }
4889
4890 fn error_context_debug_message(
4891 &mut self,
4892 instance: Instance,
4893 ty: TypeComponentLocalErrorContextTableIndex,
4894 options: OptionsIndex,
4895 err_ctx_handle: u32,
4896 debug_msg_address: u32,
4897 ) -> Result<()> {
4898 instance.error_context_debug_message(
4899 StoreContextMut(self),
4900 ty,
4901 options,
4902 err_ctx_handle,
4903 debug_msg_address,
4904 )
4905 }
4906
4907 fn thread_new_indirect(
4908 &mut self,
4909 instance: Instance,
4910 caller: RuntimeComponentInstanceIndex,
4911 func_ty_idx: TypeFuncIndex,
4912 start_func_table_idx: RuntimeTableIndex,
4913 start_func_idx: u32,
4914 context: i32,
4915 ) -> Result<u32> {
4916 instance.thread_new_indirect(
4917 StoreContextMut(self),
4918 caller,
4919 func_ty_idx,
4920 start_func_table_idx,
4921 start_func_idx,
4922 context,
4923 )
4924 }
4925}
4926
4927type HostTaskFuture = Pin<Box<dyn Future<Output = Result<()>> + Send + 'static>>;
4928
4929async fn run_with_host_task_set<F>(task: TableId<HostTask>, future: F) -> Result<F::Output>
4932where
4933 F: Future,
4934{
4935 let mut future = pin!(future);
4936 future::poll_fn(|cx| {
4937 let old_thread = match tls::get(|store| store.set_thread(task)) {
4938 Ok(thread) => thread,
4939 Err(error) => return Poll::Ready(Err(error)),
4940 };
4941 let result = future.as_mut().poll(cx);
4942 match tls::get(|store| store.set_thread(old_thread)) {
4943 Ok(_) => result.map(Ok),
4944 Err(error) => Poll::Ready(Err(error)),
4945 }
4946 })
4947 .await
4948}
4949
4950pub(crate) struct HostTask {
4954 common: WaitableCommon,
4955
4956 call_context: CallContext,
4959
4960 state: HostTaskState,
4961
4962 group: TaskGroupId,
4963}
4964
4965enum HostTaskState {
4966 CalleeStarted,
4971
4972 CalleeRunning(JoinHandle),
4977
4978 CalleeFinished(LiftedResult),
4982
4983 CalleeDone { cancelled: bool },
4986}
4987
4988impl HostTask {
4989 fn new(
4990 concurrent_state: &mut ConcurrentState,
4991 state: HostTaskState,
4992 caller: QualifiedThreadId,
4993 ) -> Result<Self> {
4994 let group = concurrent_state.get_mut(caller.task)?.group;
4995 concurrent_state.increment_group_ref_count(group)?;
4996
4997 Ok(Self {
4998 common: WaitableCommon::default(),
4999 call_context: CallContext::default(),
5000 state,
5001 group,
5002 })
5003 }
5004}
5005
5006impl TableDebug for HostTask {
5007 fn type_name() -> &'static str {
5008 "HostTask"
5009 }
5010}
5011
5012type CallbackFn = Box<dyn Fn(&mut dyn VMStore, Event, u32) -> Result<u32> + Send + Sync + 'static>;
5013
5014enum Caller {
5016 Host {
5018 tx: Option<oneshot::Sender<LiftedResult>>,
5020 host_future_present: bool,
5023 caller: Option<TableId<HostTask>>,
5027 },
5028 Guest {
5030 thread: QualifiedThreadId,
5032 },
5033}
5034
5035struct LiftResult {
5038 lift: RawLift,
5039 ty: TypeTupleIndex,
5040 memory: Option<SendSyncPtr<VMMemoryDefinition>>,
5041 string_encoding: StringEncoding,
5042}
5043
5044#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
5049pub(crate) struct QualifiedThreadId {
5050 task: TableId<GuestTask>,
5051 thread: TableId<GuestThread>,
5052}
5053
5054impl QualifiedThreadId {
5055 fn qualify(
5056 state: &mut ConcurrentState,
5057 thread: TableId<GuestThread>,
5058 ) -> Result<QualifiedThreadId> {
5059 Ok(QualifiedThreadId {
5060 task: state.get_mut(thread)?.parent_task,
5061 thread,
5062 })
5063 }
5064}
5065
5066impl fmt::Debug for QualifiedThreadId {
5067 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
5068 f.debug_tuple("QualifiedThreadId")
5069 .field(&self.task.rep())
5070 .field(&self.thread.rep())
5071 .finish()
5072 }
5073}
5074
5075enum GuestThreadState {
5076 NotStartedImplicit,
5077 NotStartedExplicit(
5078 Box<dyn FnOnce(&mut dyn VMStore, QualifiedThreadId) -> Result<()> + Send + Sync>,
5079 ),
5080 Running,
5081 Suspended(StoreFiber<'static>),
5082 Ready {
5083 fiber: StoreFiber<'static>,
5084 },
5085 Completed,
5086}
5087
5088impl fmt::Debug for GuestThreadState {
5089 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
5090 match self {
5091 Self::NotStartedImplicit => f.debug_tuple("NotStartedImplicit").finish(),
5092 Self::NotStartedExplicit(_) => f.debug_tuple("NotStartedExplicit").finish(),
5093 Self::Running => f.debug_tuple("Running").finish(),
5094 Self::Suspended(_) => f.debug_tuple("Suspended").finish(),
5095 Self::Ready { .. } => f.debug_struct("Ready").finish(),
5096 Self::Completed => f.debug_tuple("Completed").finish(),
5097 }
5098 }
5099}
5100
5101#[derive(Copy, Clone, PartialEq, Eq, Debug)]
5102enum WakeOnCancel {
5103 None,
5104 Waiting(TableId<WaitableSet>),
5105 Yielding,
5106}
5107
5108impl WakeOnCancel {
5109 fn is_none(self) -> bool {
5110 matches!(self, WakeOnCancel::None)
5111 }
5112
5113 fn replace(&mut self, other: WakeOnCancel) -> Self {
5114 let old = *self;
5115 *self = other;
5116 old
5117 }
5118
5119 fn take(&mut self) -> Self {
5120 self.replace(WakeOnCancel::None)
5121 }
5122}
5123
5124pub struct GuestThread {
5125 context: [u32; NUM_COMPONENT_CONTEXT_SLOTS],
5128 parent_task: TableId<GuestTask>,
5130 wake_on_cancel: WakeOnCancel,
5133 state: GuestThreadState,
5135 instance_rep: Option<u32>,
5138 sync_call_set: TableId<WaitableSet>,
5140 old_do_not_suspend: Option<bool>,
5143}
5144
5145impl GuestThread {
5146 fn from_instance(
5149 state: Pin<&mut ComponentInstance>,
5150 caller_instance: RuntimeComponentInstanceIndex,
5151 guest_thread: u32,
5152 ) -> Result<TableId<Self>> {
5153 let rep = state.instance_states().0[caller_instance]
5154 .thread_handle_table()
5155 .guest_thread_rep(guest_thread)?;
5156 Ok(TableId::new(rep))
5157 }
5158
5159 fn new_implicit(state: &mut ConcurrentState, parent_task: TableId<GuestTask>) -> Result<Self> {
5160 let sync_call_set = state.push(WaitableSet {
5161 is_sync_call_set: true,
5162 ..WaitableSet::default()
5163 })?;
5164 Ok(Self {
5165 context: [0; NUM_COMPONENT_CONTEXT_SLOTS],
5166 parent_task,
5167 wake_on_cancel: WakeOnCancel::None,
5168 state: GuestThreadState::NotStartedImplicit,
5169 instance_rep: None,
5170 sync_call_set,
5171 old_do_not_suspend: None,
5172 })
5173 }
5174
5175 fn new_explicit(
5176 state: &mut ConcurrentState,
5177 parent_task: TableId<GuestTask>,
5178 start_func: Box<
5179 dyn FnOnce(&mut dyn VMStore, QualifiedThreadId) -> Result<()> + Send + Sync,
5180 >,
5181 ) -> Result<Self> {
5182 let sync_call_set = state.push(WaitableSet {
5183 is_sync_call_set: true,
5184 ..WaitableSet::default()
5185 })?;
5186 Ok(Self {
5187 context: [0; NUM_COMPONENT_CONTEXT_SLOTS],
5188 parent_task,
5189 wake_on_cancel: WakeOnCancel::None,
5190 state: GuestThreadState::NotStartedExplicit(start_func),
5191 instance_rep: None,
5192 sync_call_set,
5193 old_do_not_suspend: None,
5194 })
5195 }
5196}
5197
5198impl TableDebug for GuestThread {
5199 fn type_name() -> &'static str {
5200 "GuestThread"
5201 }
5202}
5203
5204enum SyncResult {
5205 NotProduced,
5206 Produced(Option<ValRaw>),
5207 Taken,
5208}
5209
5210impl SyncResult {
5211 fn take(&mut self) -> Result<Option<Option<ValRaw>>> {
5212 Ok(match mem::replace(self, SyncResult::Taken) {
5213 SyncResult::NotProduced => None,
5214 SyncResult::Produced(val) => Some(val),
5215 SyncResult::Taken => {
5216 bail_bug!("attempted to take a synchronous result that was already taken")
5217 }
5218 })
5219 }
5220}
5221
5222#[derive(Debug)]
5223enum HostFutureState {
5224 NotApplicable,
5225 Live,
5226 Dropped,
5227}
5228
5229pub(crate) struct GuestTask {
5231 common: WaitableCommon,
5233 lower_params: Option<RawLower>,
5235 lift_result: Option<LiftResult>,
5237 result: Option<LiftedResult>,
5240 callback: Option<CallbackFn>,
5243 caller: Caller,
5245 call_context: CallContext,
5250 sync_result: SyncResult,
5253 cancel_request_delivered: bool,
5257 starting_sent: bool,
5260 instance: RuntimeInstance,
5267 event: Option<Event>,
5270 exited: bool,
5272 threads: HashSet<TableId<GuestThread>>,
5274 host_future_state: HostFutureState,
5277 async_typed: bool,
5280 async_lifted: bool,
5283
5284 decremented_interesting_task_count: bool,
5285
5286 group: TaskGroupId,
5287}
5288
5289impl GuestTask {
5290 fn already_lowered_parameters(&self) -> bool {
5291 self.lower_params.is_none()
5293 }
5294
5295 fn returned_or_cancelled(&self) -> bool {
5296 self.lift_result.is_none()
5298 }
5299
5300 fn ready_to_delete(&self) -> bool {
5301 let threads_completed = self.threads.is_empty();
5302 let has_sync_result = matches!(self.sync_result, SyncResult::Produced(_));
5303 let pending_completion_event = matches!(
5304 self.common.event,
5305 Some(Event::Subtask {
5306 status: Status::Returned | Status::ReturnCancelled
5307 })
5308 );
5309 let ready = threads_completed
5310 && !has_sync_result
5311 && !pending_completion_event
5312 && !matches!(self.host_future_state, HostFutureState::Live);
5313 log::trace!(
5314 "ready to delete? {ready} (threads_completed: {}, has_sync_result: {}, pending_completion_event: {}, host_future_state: {:?})",
5315 threads_completed,
5316 has_sync_result,
5317 pending_completion_event,
5318 self.host_future_state
5319 );
5320 ready
5321 }
5322
5323 fn new(
5324 state: &mut ConcurrentState,
5325 lower_params: RawLower,
5326 lift_result: LiftResult,
5327 caller: Caller,
5328 callback: Option<CallbackFn>,
5329 instance: RuntimeInstance,
5330 async_typed: bool,
5331 async_lifted: bool,
5332 ) -> Result<QualifiedThreadId> {
5333 let host_future_state = match &caller {
5334 Caller::Guest { .. } => HostFutureState::NotApplicable,
5335 Caller::Host {
5336 host_future_present,
5337 ..
5338 } => {
5339 if *host_future_present {
5340 HostFutureState::Live
5341 } else {
5342 HostFutureState::NotApplicable
5343 }
5344 }
5345 };
5346
5347 let group = match caller {
5348 Caller::Guest { thread } => {
5349 let group = state.get_mut(thread.task)?.group;
5350 state.increment_group_ref_count(group)?;
5351 group
5352 }
5353 Caller::Host { .. } => state.make_task_group()?,
5354 };
5355
5356 let task = state.push(Self {
5357 common: WaitableCommon::default(),
5358 lower_params: Some(lower_params),
5359 lift_result: Some(lift_result),
5360 result: None,
5361 callback,
5362 caller,
5363 call_context: CallContext::default(),
5364 sync_result: SyncResult::NotProduced,
5365 cancel_request_delivered: false,
5366 starting_sent: false,
5367 instance,
5368 event: None,
5369 exited: false,
5370 threads: HashSet::new(),
5371 host_future_state,
5372 async_typed,
5373 async_lifted,
5374 decremented_interesting_task_count: false,
5375 group,
5376 })?;
5377 let new_thread = GuestThread::new_implicit(state, task)?;
5378 let thread = state.push(new_thread)?;
5379 state.get_mut(task)?.threads.insert(thread);
5380 state.interesting_tasks += 1;
5381 let thread = QualifiedThreadId { task, thread };
5382 log::trace!("new implicit thread {thread:?} for instance {instance:?}");
5383 Ok(thread)
5384 }
5385}
5386
5387impl TableDebug for GuestTask {
5388 fn type_name() -> &'static str {
5389 "GuestTask"
5390 }
5391}
5392
5393#[derive(Default)]
5395struct WaitableCommon {
5396 event: Option<Event>,
5398 set: Option<TableId<WaitableSet>>,
5400 handle: Option<u32>,
5402}
5403
5404#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
5406enum Waitable {
5407 Host(TableId<HostTask>),
5409 Guest(TableId<GuestTask>),
5411 Transmit(TableId<TransmitHandle>),
5413}
5414
5415impl Waitable {
5416 fn from_instance(
5419 state: Pin<&mut ComponentInstance>,
5420 caller_instance: RuntimeComponentInstanceIndex,
5421 waitable: u32,
5422 ) -> Result<Self> {
5423 use crate::runtime::vm::component::Waitable;
5424
5425 let (waitable, kind) = state.instance_states().0[caller_instance]
5426 .handle_table()
5427 .waitable_rep(waitable)?;
5428
5429 Ok(match kind {
5430 Waitable::Subtask { is_host: true } => Self::Host(TableId::new(waitable)),
5431 Waitable::Subtask { is_host: false } => Self::Guest(TableId::new(waitable)),
5432 Waitable::Stream | Waitable::Future => Self::Transmit(TableId::new(waitable)),
5433 })
5434 }
5435
5436 fn rep(&self) -> u32 {
5438 match self {
5439 Self::Host(id) => id.rep(),
5440 Self::Guest(id) => id.rep(),
5441 Self::Transmit(id) => id.rep(),
5442 }
5443 }
5444
5445 fn join(&self, state: &mut ConcurrentState, set: Option<TableId<WaitableSet>>) -> Result<()> {
5449 log::trace!("waitable {self:?} join set {set:?}");
5450
5451 let old = mem::replace(&mut self.common(state)?.set, set);
5452
5453 if let Some(old) = old {
5454 match *self {
5455 Waitable::Host(id) => state.remove_child(id, old),
5456 Waitable::Guest(id) => state.remove_child(id, old),
5457 Waitable::Transmit(id) => state.remove_child(id, old),
5458 }?;
5459
5460 state.get_mut(old)?.ready.remove(self);
5461 }
5462
5463 if let Some(set) = set {
5464 match *self {
5465 Waitable::Host(id) => state.add_child(id, set),
5466 Waitable::Guest(id) => state.add_child(id, set),
5467 Waitable::Transmit(id) => state.add_child(id, set),
5468 }?;
5469
5470 if self.common(state)?.event.is_some() {
5471 self.mark_ready(state)?;
5472 }
5473 }
5474
5475 Ok(())
5476 }
5477
5478 fn common<'a>(&self, state: &'a mut ConcurrentState) -> Result<&'a mut WaitableCommon> {
5480 Ok(match self {
5481 Self::Host(id) => &mut state.get_mut(*id)?.common,
5482 Self::Guest(id) => &mut state.get_mut(*id)?.common,
5483 Self::Transmit(id) => &mut state.get_mut(*id)?.common,
5484 })
5485 }
5486
5487 fn trap_if_in_waitable_set(&self, state: &mut ConcurrentState) -> Result<()> {
5493 if self.common(state)?.set.is_some() {
5494 bail!(Trap::WaitableSyncAndAsync);
5495 }
5496 Ok(())
5497 }
5498
5499 fn set_event(&self, state: &mut ConcurrentState, event: Option<Event>) -> Result<()> {
5503 log::trace!("set event for {self:?}: {event:?}");
5504 self.common(state)?.event = event;
5505 self.mark_ready(state)
5506 }
5507
5508 fn take_event(&self, state: &mut ConcurrentState) -> Result<Option<Event>> {
5510 let common = self.common(state)?;
5511 let event = common.event.take();
5512 if let Some(set) = self.common(state)?.set {
5513 state.get_mut(set)?.ready.remove(self);
5514 }
5515
5516 Ok(event)
5517 }
5518
5519 fn mark_ready(&self, state: &mut ConcurrentState) -> Result<()> {
5523 if let Some(set) = self.common(state)?.set {
5524 let set_state = state.get_mut(set)?;
5525 set_state.ready.insert(*self);
5526
5527 if let Some((thread, mode)) = set_state.waiting.pop_first() {
5528 let wake_on_cancel = state.get_mut(thread.thread)?.wake_on_cancel.take();
5529 assert!(wake_on_cancel.is_none() || wake_on_cancel == WakeOnCancel::Waiting(set));
5530
5531 let item = match mode {
5532 WaitMode::Fiber(fiber) => Some(WorkItem::ResumeFiber {
5533 instance: state.get_mut(thread.task)?.instance,
5534 thread,
5535 fiber,
5536 }),
5537 WaitMode::Callback(instance) => Some(WorkItem::GuestCall {
5538 instance: state.get_mut(thread.task)?.instance,
5539 call: GuestCall {
5540 thread,
5541 kind: GuestCallKind::DeliverEvent {
5542 instance,
5543 set: Some(set),
5544 },
5545 },
5546 }),
5547 };
5548
5549 if let Some(item) = item {
5550 state.push_high_priority(item);
5551 }
5552 }
5553 }
5554 Ok(())
5555 }
5556
5557 fn delete_from(&self, store: &mut StoreOpaque) -> Result<()> {
5559 match self {
5560 Self::Host(task) => {
5561 log::trace!("delete host task {task:?}");
5562 let state = store.concurrent_state_mut()?;
5563 let task = state.delete(*task)?;
5564
5565 state.decrement_group_ref_count(task.group)?;
5566 }
5567 Self::Guest(task) => {
5568 log::trace!("delete guest task {task:?}");
5569 let state = store.concurrent_state_mut()?;
5570 let task = state.delete(*task)?;
5571
5572 state.decrement_group_ref_count(task.group)?;
5573
5574 debug_assert!(task.decremented_interesting_task_count);
5581 }
5582 Self::Transmit(task) => {
5583 store.concurrent_state_mut()?.delete(*task)?;
5584 }
5585 }
5586
5587 Ok(())
5588 }
5589}
5590
5591impl fmt::Debug for Waitable {
5592 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
5593 match self {
5594 Self::Host(id) => write!(f, "{id:?}"),
5595 Self::Guest(id) => write!(f, "{id:?}"),
5596 Self::Transmit(id) => write!(f, "{id:?}"),
5597 }
5598 }
5599}
5600
5601#[derive(Default)]
5603struct WaitableSet {
5604 ready: BTreeSet<Waitable>,
5606 waiting: BTreeMap<QualifiedThreadId, WaitMode>,
5608 is_sync_call_set: bool,
5611}
5612
5613impl TableDebug for WaitableSet {
5614 fn type_name() -> &'static str {
5615 "WaitableSet"
5616 }
5617}
5618
5619type RawLower =
5621 Box<dyn FnOnce(&mut dyn VMStore, &mut [MaybeUninit<ValRaw>]) -> Result<()> + Send + Sync>;
5622
5623type RawLift = Box<
5625 dyn FnOnce(&mut dyn VMStore, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>> + Send + Sync,
5626>;
5627
5628type LiftedResult = Box<dyn Any + Send + Sync>;
5632
5633struct DummyResult;
5636
5637#[derive(Default)]
5639pub struct ConcurrentInstanceState {
5640 backpressure: u16,
5642 do_not_enter: bool,
5644 do_not_suspend: bool,
5647 pending: BTreeMap<QualifiedThreadId, GuestCallKind>,
5650}
5651
5652impl ConcurrentInstanceState {
5653 pub fn pending_is_empty(&self) -> bool {
5654 self.pending.is_empty()
5655 }
5656}
5657
5658#[derive(Debug, Copy, Clone)]
5659pub(crate) enum CurrentThread {
5660 Guest(QualifiedThreadId),
5663 Host(TableId<HostTask>),
5665 DeferredHost(QualifiedThreadId),
5668 None,
5671}
5672
5673impl CurrentThread {
5674 fn guest(&self) -> Option<&QualifiedThreadId> {
5675 match self {
5676 Self::Guest(id) => Some(id),
5677 _ => None,
5678 }
5679 }
5680
5681 fn guest_task(&self) -> Option<TableId<GuestTask>> {
5682 match self {
5683 Self::Guest(id) => Some(id.task),
5684 _ => None,
5685 }
5686 }
5687
5688 fn is_none(&self) -> bool {
5689 matches!(self, Self::None)
5690 }
5691}
5692
5693impl From<QualifiedThreadId> for CurrentThread {
5694 fn from(id: QualifiedThreadId) -> Self {
5695 Self::Guest(id)
5696 }
5697}
5698
5699impl From<TableId<HostTask>> for CurrentThread {
5700 fn from(id: TableId<HostTask>) -> Self {
5701 Self::Host(id)
5702 }
5703}
5704
5705enum Priority {
5706 Switch,
5707 High,
5708 Low,
5709}
5710
5711pub struct ConcurrentState {
5713 unforced_current_thread: CurrentThread,
5719
5720 deferred_host_call_context: Option<CallContext>,
5726
5727 futures: AlwaysMut<Option<FuturesUnordered<HostTaskFuture>>>,
5732 table: AlwaysMut<ResourceTable>,
5734 switch_item: Option<WorkItem>,
5742 next_switch_item: Option<WorkItem>,
5748 high_priority: VecDeque<WorkItem>,
5750 low_priority: VecDeque<WorkItem>,
5752 suspend_reason: Option<SuspendReason>,
5756 worker: Option<StoreFiber<'static>>,
5760 worker_item: Option<WorkerItem>,
5762
5763 global_error_context_ref_counts:
5776 BTreeMap<TypeComponentGlobalErrorContextTableIndex, GlobalErrorContextRefCount>,
5777
5778 interesting_tasks: usize,
5791
5792 interesting_tasks_empty_waker: Option<Waker>,
5796
5797 ready_for_concurrent_call_waker: Option<Waker>,
5802
5803 event_loop_running: bool,
5805
5806 #[cfg(feature = "task-group-hook")]
5808 task_group_hook: Option<Box<dyn TaskGroupHook>>,
5809}
5810
5811impl Default for ConcurrentState {
5812 fn default() -> Self {
5813 Self {
5814 unforced_current_thread: CurrentThread::None,
5815 deferred_host_call_context: None,
5816 table: AlwaysMut::new(ResourceTable::new()),
5817 futures: AlwaysMut::new(Some(FuturesUnordered::new())),
5818 switch_item: None,
5819 next_switch_item: None,
5820 high_priority: VecDeque::new(),
5821 low_priority: VecDeque::new(),
5822 suspend_reason: None,
5823 worker: None,
5824 worker_item: None,
5825 global_error_context_ref_counts: BTreeMap::new(),
5826 interesting_tasks: 0,
5827 interesting_tasks_empty_waker: None,
5828 ready_for_concurrent_call_waker: None,
5829 event_loop_running: false,
5830 #[cfg(feature = "task-group-hook")]
5831 task_group_hook: None,
5832 }
5833 }
5834}
5835
5836impl ConcurrentState {
5837 pub(crate) fn take_fibers_and_futures(
5854 &mut self,
5855 fibers: &mut Vec<StoreFiber<'static>>,
5856 futures: &mut Vec<FuturesUnordered<HostTaskFuture>>,
5857 ) {
5858 let mut items = Vec::new();
5859 for (_, entry) in self.table.get_mut().iter_mut() {
5860 if let Some(set) = entry.downcast_mut::<WaitableSet>() {
5861 for mode in mem::take(&mut set.waiting).into_values() {
5862 match mode {
5863 WaitMode::Fiber(fiber) => {
5864 fibers.push(fiber);
5865 }
5866 WaitMode::Callback(_) => {}
5867 }
5868 }
5869 } else if let Some(thread) = entry.downcast_mut::<GuestThread>() {
5870 if let GuestThreadState::Suspended(fiber) | GuestThreadState::Ready { fiber, .. } =
5871 mem::replace(&mut thread.state, GuestThreadState::Completed)
5872 {
5873 fibers.push(fiber);
5874 }
5875 } else if let Some(item) = entry.downcast_mut::<Option<WorkItem>>() {
5876 if let Some(item) = item.take() {
5877 items.push(item);
5878 }
5879 }
5880 }
5881
5882 if let Some(fiber) = self.worker.take() {
5883 fibers.push(fiber);
5884 }
5885
5886 let mut handle_item = |item| match item {
5887 WorkItem::ResumeFiber { fiber, .. } => {
5888 fibers.push(fiber);
5889 }
5890 WorkItem::PushFuture(future) => {
5891 self.futures
5892 .get_mut()
5893 .as_mut()
5894 .unwrap()
5895 .push(future.into_inner());
5896 }
5897 WorkItem::ResumeThread { .. }
5898 | WorkItem::GuestCall { .. }
5899 | WorkItem::WorkerFunction(_) => {}
5900 };
5901
5902 for item in items {
5903 handle_item(item);
5904 }
5905 if let Some(item) = self.switch_item.take() {
5906 handle_item(item);
5907 }
5908 if let Some(item) = self.next_switch_item.take() {
5909 handle_item(item);
5910 }
5911 for item in mem::take(&mut self.high_priority) {
5912 handle_item(item);
5913 }
5914 for item in mem::take(&mut self.low_priority) {
5915 handle_item(item);
5916 }
5917
5918 if let Some(them) = self.futures.get_mut().take() {
5919 futures.push(them);
5920 }
5921 }
5922
5923 #[cfg(feature = "gc")]
5924 pub(crate) fn trace_fiber_roots(
5925 &mut self,
5926 modules: &ModuleRegistry,
5927 unwind: &dyn Unwind,
5928 gc_roots_list: &mut GcRootsList,
5929 ) {
5930 let ConcurrentState {
5931 table,
5932 worker,
5933 switch_item,
5934 next_switch_item,
5935 high_priority,
5936 low_priority,
5937
5938 futures: _,
5942
5943 worker_item: _,
5945 unforced_current_thread: _,
5946 deferred_host_call_context: _,
5947 suspend_reason: _,
5948 global_error_context_ref_counts: _,
5949 interesting_tasks: _,
5950 interesting_tasks_empty_waker: _,
5951 ready_for_concurrent_call_waker: _,
5952 event_loop_running: _,
5953 #[cfg(feature = "task-group-hook")]
5954 task_group_hook: _,
5955 } = self;
5956
5957 for (_, entry) in table.get_mut().iter_mut() {
5958 if let Some(set) = entry.downcast_mut::<WaitableSet>() {
5959 for mode in set.waiting.values_mut() {
5960 match mode {
5961 WaitMode::Fiber(fiber) => {
5962 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
5963 }
5964 WaitMode::Callback(_) => {}
5965 }
5966 }
5967 } else if let Some(thread) = entry.downcast_mut::<GuestThread>() {
5968 if let GuestThreadState::Suspended(fiber) | GuestThreadState::Ready { fiber, .. } =
5969 &mut thread.state
5970 {
5971 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
5972 }
5973 } else if let Some(Some(WorkItem::ResumeFiber { fiber, .. })) =
5974 entry.downcast_mut::<Option<WorkItem>>()
5975 {
5976 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
5977 }
5978 }
5979
5980 if let Some(fiber) = worker {
5981 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
5982 }
5983
5984 let mut handle_item = |item: &mut WorkItem| match item {
5985 WorkItem::ResumeFiber { fiber, .. } => {
5986 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
5987 }
5988 WorkItem::PushFuture(_future) => {
5989 }
5992 WorkItem::ResumeThread { .. }
5993 | WorkItem::GuestCall { .. }
5994 | WorkItem::WorkerFunction(_) => {}
5995 };
5996
5997 if let Some(item) = switch_item {
5998 handle_item(item);
5999 }
6000 if let Some(item) = next_switch_item {
6001 handle_item(item);
6002 }
6003 for item in high_priority {
6004 handle_item(item);
6005 }
6006 for item in low_priority {
6007 handle_item(item);
6008 }
6009 }
6010
6011 fn push<V: Send + Sync + 'static>(
6012 &mut self,
6013 value: V,
6014 ) -> Result<TableId<V>, ResourceTableError> {
6015 self.table.get_mut().push(value).map(TableId::from)
6016 }
6017
6018 fn get_mut<V: 'static>(&mut self, id: TableId<V>) -> Result<&mut V, ResourceTableError> {
6019 self.table.get_mut().get_mut(&Resource::from(id))
6020 }
6021
6022 pub fn add_child<T: 'static, U: 'static>(
6023 &mut self,
6024 child: TableId<T>,
6025 parent: TableId<U>,
6026 ) -> Result<(), ResourceTableError> {
6027 self.table
6028 .get_mut()
6029 .add_child(Resource::from(child), Resource::from(parent))
6030 }
6031
6032 pub fn remove_child<T: 'static, U: 'static>(
6033 &mut self,
6034 child: TableId<T>,
6035 parent: TableId<U>,
6036 ) -> Result<(), ResourceTableError> {
6037 self.table
6038 .get_mut()
6039 .remove_child(Resource::from(child), Resource::from(parent))
6040 }
6041
6042 fn delete<V: 'static>(&mut self, id: TableId<V>) -> Result<V, ResourceTableError> {
6043 self.table.get_mut().delete(Resource::from(id))
6044 }
6045
6046 fn push_future(&mut self, future: HostTaskFuture) {
6047 self.push_high_priority(WorkItem::PushFuture(AlwaysMut::new(future)));
6054 }
6055
6056 fn set_switch_item(&mut self, item: WorkItem) -> Result<()> {
6057 log::trace!("set switch item: {item:?}");
6058
6059 if self.switch_item.is_some() {
6060 bail_bug!("switch item already set");
6061 }
6062
6063 self.switch_item = Some(item);
6064
6065 Ok(())
6066 }
6067
6068 fn take_next_switch_item(&mut self) -> Result<()> {
6069 if let Some(item) = self.next_switch_item.take() {
6070 self.set_switch_item(item)?;
6071 }
6072 Ok(())
6073 }
6074
6075 fn push_high_priority(&mut self, item: WorkItem) {
6076 log::trace!("push high priority: {item:?}");
6077 self.high_priority.push_front(item);
6078 }
6079
6080 fn push_low_priority(&mut self, item: WorkItem) {
6081 log::trace!("push low priority: {item:?}");
6082 self.low_priority.push_front(item);
6083 }
6084
6085 fn push_work_item(&mut self, item: WorkItem, priority: Priority) -> Result<()> {
6086 match priority {
6087 Priority::Switch => self.set_switch_item(item)?,
6088 Priority::High => self.push_high_priority(item),
6089 Priority::Low => self.push_low_priority(item),
6090 }
6091
6092 Ok(())
6093 }
6094
6095 fn promote_instance_local_thread_work_item(
6096 &mut self,
6097 current_instance: RuntimeInstance,
6098 ) -> Result<bool> {
6099 log::trace!("promote thread work items for {current_instance:?}");
6100
6101 self.promote_work_item_matching(|item: &WorkItem| {
6102 let result = match item {
6103 WorkItem::ResumeThread { instance, .. }
6104 | WorkItem::ResumeFiber { instance, .. }
6105 | WorkItem::GuestCall { instance, .. } => *instance == current_instance,
6106 _ => false,
6107 };
6108
6109 log::trace!("candidate {item:?}: {result}");
6110 result
6111 })
6112 }
6113
6114 fn promote_thread_work_item(&mut self, thread: QualifiedThreadId) -> Result<bool> {
6115 self.promote_work_item_matching(|item: &WorkItem| match item {
6116 WorkItem::ResumeThread {
6117 thread: item_thread,
6118 ..
6119 }
6120 | WorkItem::GuestCall {
6121 call:
6122 GuestCall {
6123 thread: item_thread,
6124 ..
6125 },
6126 ..
6127 } => *item_thread == thread,
6128 _ => false,
6129 })
6130 }
6131
6132 fn promote_work_item_matching<F>(&mut self, mut predicate: F) -> Result<bool>
6133 where
6134 F: FnMut(&WorkItem) -> bool,
6135 {
6136 for item in mem::take(&mut self.high_priority).into_iter().rev() {
6141 if self.switch_item.is_none() && predicate(&item) {
6142 self.set_switch_item(item)?;
6143 } else {
6144 self.push_high_priority(item);
6145 }
6146 }
6147
6148 if self.switch_item.is_none() {
6149 for item in mem::take(&mut self.low_priority).into_iter().rev() {
6150 if self.switch_item.is_none() && predicate(&item) {
6151 self.set_switch_item(item)?;
6152 } else {
6153 self.push_low_priority(item);
6154 }
6155 }
6156 }
6157
6158 Ok(self.switch_item.is_some())
6159 }
6160
6161 pub fn call_context(&mut self, task: Scope) -> Result<&mut CallContext> {
6164 match task {
6165 Scope::HostId(task) => {
6166 let task: TableId<HostTask> = TableId::new(task);
6167 Ok(&mut self.get_mut(task)?.call_context)
6168 }
6169 Scope::Id(task) => {
6170 let task: TableId<GuestTask> = TableId::new(task);
6171 Ok(&mut self.get_mut(task)?.call_context)
6172 }
6173 }
6174 }
6175
6176 pub(crate) fn deferred_host_call_context(&mut self) -> Option<&mut CallContext> {
6177 self.deferred_host_call_context.as_mut()
6178 }
6179
6180 fn futures_mut(&mut self) -> Result<&mut FuturesUnordered<HostTaskFuture>> {
6181 match self.futures.get_mut().as_mut() {
6182 Some(f) => Ok(f),
6183 None => bail_bug!("futures field of concurrent state is currently taken"),
6184 }
6185 }
6186
6187 pub(crate) fn table(&mut self) -> &mut ResourceTable {
6188 self.table.get_mut()
6189 }
6190
6191 fn debug_assert_deferred_host_invariant(&self) {
6192 debug_assert_eq!(
6193 self.deferred_host_call_context.is_some(),
6194 matches!(self.unforced_current_thread, CurrentThread::DeferredHost(_)),
6195 "a deferred host thread and call context must exist together",
6196 );
6197 }
6198
6199 fn materialize_host_task(&mut self) -> Result<CurrentThread> {
6200 self.debug_assert_deferred_host_invariant();
6201 let caller = match self.unforced_current_thread {
6202 CurrentThread::DeferredHost(caller) => caller,
6203 thread => return Ok(thread),
6204 };
6205
6206 let task = HostTask::new(self, HostTaskState::CalleeStarted, caller)?;
6208 let task = self.push(task)?;
6209 let call_context = self
6210 .deferred_host_call_context
6211 .take()
6212 .expect("deferred host call context should be present");
6213 self.get_mut(task)
6214 .expect("newly inserted host task should be present")
6215 .call_context = call_context;
6216 self.unforced_current_thread = CurrentThread::Host(task);
6217 self.debug_assert_deferred_host_invariant();
6218 log::trace!("new host task materialized {task:?}");
6219 Ok(CurrentThread::Host(task))
6220 }
6221
6222 fn materialize_current_host_task_id(&mut self) -> Result<Option<TableId<HostTask>>> {
6223 match self.materialize_host_task()? {
6224 CurrentThread::Host(id) => Ok(Some(id)),
6225 CurrentThread::None => Ok(None),
6226 CurrentThread::Guest(_) => {
6227 bail_bug!("tried to materialize a host task id from a guest thread")
6228 }
6229 CurrentThread::DeferredHost(_) => {
6230 bail_bug!(
6231 "current thread is a deferred host thread which should have been materialized"
6232 )
6233 }
6234 }
6235 }
6236
6237 pub(crate) fn materialize_current_scope(&mut self) -> Result<Scope> {
6238 match self.materialize_host_task()? {
6239 CurrentThread::Host(id) => Ok(Scope::HostId(id.rep())),
6240 _ => bail_bug!("current scope is not a deferred host scope"),
6241 }
6242 }
6243}
6244
6245fn for_any_lower<
6248 F: FnOnce(&mut dyn VMStore, &mut [MaybeUninit<ValRaw>]) -> Result<()> + Send + Sync,
6249>(
6250 fun: F,
6251) -> F {
6252 fun
6253}
6254
6255fn for_any_lift<
6257 F: FnOnce(&mut dyn VMStore, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>> + Send + Sync,
6258>(
6259 fun: F,
6260) -> F {
6261 fun
6262}
6263
6264fn check_ambient_store(id: StoreId) {
6265 let message = "\
6266 `Future`s which depend on asynchronous component tasks, streams, or \
6267 futures to complete may only be polled from the event loop of the \
6268 store to which they belong. Please use \
6269 `StoreContextMut::{run_concurrent,spawn}` to poll or await them.\
6270 ";
6271 tls::try_get(|store| {
6272 let matched = match store {
6273 tls::TryGet::Some(store) => store.id() == id,
6274 tls::TryGet::Taken | tls::TryGet::None => false,
6275 };
6276
6277 if !matched {
6278 panic!("{message}")
6279 }
6280 });
6281}
6282
6283fn unpack_callback_code(code: u32) -> (u32, u32) {
6284 (code & 0xF, code >> 4)
6285}
6286
6287struct WaitableCheckParams {
6291 set: TableId<WaitableSet>,
6292 options: OptionsIndex,
6293 payload: u32,
6294}
6295
6296enum WaitableCheck {
6299 Wait,
6300 Poll,
6301}
6302
6303pub(crate) struct PreparedCall<R> {
6305 handle: Func,
6307 thread: QualifiedThreadId,
6309 param_count: usize,
6311 rx: oneshot::Receiver<LiftedResult>,
6314 runtime_instance: RuntimeInstance,
6316 _phantom: PhantomData<R>,
6317}
6318
6319impl<R> PreparedCall<R> {
6320 pub(crate) fn task_id(&self) -> TaskId {
6322 TaskId {
6323 task: self.thread.task,
6324 runtime_instance: self.runtime_instance,
6325 }
6326 }
6327}
6328
6329pub(crate) struct TaskId {
6331 task: TableId<GuestTask>,
6332 runtime_instance: RuntimeInstance,
6333}
6334
6335impl TaskId {
6336 pub(crate) fn host_future_dropped(&self, store: &mut StoreOpaque) -> Result<()> {
6342 let task = store.concurrent_state_mut()?.get_mut(self.task)?;
6343 let delete = if !task.already_lowered_parameters() {
6344 store.cancel_guest_subtask_without_lowered_parameters(
6345 self.runtime_instance,
6346 self.task,
6347 )?;
6348 true
6349 } else {
6350 task.host_future_state = HostFutureState::Dropped;
6351 task.ready_to_delete()
6352 };
6353 if delete {
6354 Waitable::Guest(self.task).delete_from(store)?
6355 }
6356 Ok(())
6357 }
6358}
6359
6360pub(crate) fn prepare_call<T, R>(
6366 mut store: StoreContextMut<T>,
6367 handle: Func,
6368 param_count: usize,
6369 host_future_present: bool,
6370 lower_params: impl FnOnce(StoreContextMut<T>, &mut [MaybeUninit<ValRaw>]) -> Result<()>
6371 + Send
6372 + Sync
6373 + 'static,
6374 lift_result: impl FnOnce(&mut StoreOpaque, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>>
6375 + Send
6376 + Sync
6377 + 'static,
6378) -> Result<PreparedCall<R>> {
6379 if !store.0.may_enter() {
6380 bail!(Trap::CannotEnterComponent);
6381 }
6382
6383 let (options, _flags, ty, raw_options) = handle.abi_info(store.0);
6384
6385 let instance = handle.instance().id().get(store.0);
6386 let options = &instance.component().env_component().options[options];
6387 let ty = &instance.component().types()[ty];
6388 let async_typed = ty.async_;
6389 let async_lifted = raw_options.async_;
6390 let task_return_type = ty.results;
6391 let component_instance = raw_options.instance;
6392 let callback = options.callback.map(|i| instance.runtime_callback(i));
6393 let memory = options
6394 .memory()
6395 .map(|i| instance.runtime_memory(i))
6396 .map(SendSyncPtr::new);
6397 let string_encoding = options.string_encoding;
6398 let token = StoreToken::new(store.as_context_mut());
6399 let caller = store.0.materialize_host_task_id()?;
6400 let state = store.0.concurrent_state_mut()?;
6401
6402 let (tx, rx) = oneshot::channel();
6403
6404 let instance = handle.instance().runtime_instance(component_instance);
6405 let thread = GuestTask::new(
6406 state,
6407 Box::new(for_any_lower(move |store, params| {
6408 lower_params(token.as_context_mut(store), params)
6409 })),
6410 LiftResult {
6411 lift: Box::new(for_any_lift(move |store, result| {
6412 lift_result(store, result)
6413 })),
6414 ty: task_return_type,
6415 memory,
6416 string_encoding,
6417 },
6418 Caller::Host {
6419 tx: Some(tx),
6420 host_future_present,
6421 caller,
6422 },
6423 callback.map(|callback| {
6424 let callback = SendSyncPtr::new(callback);
6425 let instance = handle.instance();
6426 Box::new(move |store: &mut dyn VMStore, event, handle| {
6427 let store = token.as_context_mut(store);
6428 unsafe { instance.call_callback(store, callback, event, handle) }
6431 }) as CallbackFn
6432 }),
6433 instance,
6434 async_typed,
6435 async_lifted,
6436 )?;
6437
6438 Ok(PreparedCall {
6439 handle,
6440 thread,
6441 param_count,
6442 runtime_instance: instance,
6443 rx,
6444 _phantom: PhantomData,
6445 })
6446}
6447
6448pub(crate) struct StagedCall<R> {
6449 store: StoreId,
6450 rx: oneshot::Receiver<LiftedResult>,
6451 _marker: PhantomData<fn() -> R>,
6452 group: TaskGroupId,
6453}
6454
6455impl<R> StagedCall<R> {
6456 pub(crate) fn new<T: 'static>(
6463 mut store: StoreContextMut<T>,
6464 prepared: PreparedCall<R>,
6465 ) -> Result<StagedCall<R>> {
6466 let PreparedCall {
6467 handle,
6468 thread,
6469 param_count,
6470 rx,
6471 ..
6472 } = prepared;
6473
6474 stage_call0(store.as_context_mut(), handle, thread, param_count)?;
6475
6476 Ok(StagedCall {
6477 store: store.0.id(),
6478 rx,
6479 _marker: PhantomData,
6480 group: store.0.concurrent_state_mut()?.get_mut(thread.task)?.group,
6481 })
6482 }
6483}
6484
6485impl<R> Future for StagedCall<R>
6486where
6487 R: 'static,
6488{
6489 type Output = Result<R>;
6490
6491 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
6492 check_ambient_store(self.store);
6493 Pin::new(&mut self.rx).poll(cx).map(|result| match result {
6494 Ok(r) => match r.downcast() {
6495 Ok(r) => Ok(*r),
6496 Err(_) => bail_bug!("wrong type of value produced"),
6497 },
6498 Err(oneshot::Canceled) => bail_bug!("channel erroneously dropped"),
6499 })
6500 }
6501}
6502
6503fn stage_call0<T: 'static>(
6506 store: StoreContextMut<T>,
6507 handle: Func,
6508 guest_thread: QualifiedThreadId,
6509 param_count: usize,
6510) -> Result<()> {
6511 let (_options, _, _ty, raw_options) = handle.abi_info(store.0);
6512 let is_concurrent = raw_options.async_;
6513 let callback = raw_options.callback;
6514 let instance = handle.instance();
6515 let callee = handle.lifted_core_func(store.0);
6516 let post_return = raw_options
6517 .post_return
6518 .map(|i| instance.id().get(store.0).runtime_post_return(i));
6519 let callback = callback.map(|i| {
6520 let instance = instance.id().get(store.0);
6521 SendSyncPtr::new(instance.runtime_callback(i))
6522 });
6523
6524 log::trace!("queueing call {guest_thread:?}");
6525
6526 unsafe {
6530 instance.stage_call(
6531 store,
6532 guest_thread,
6533 SendSyncPtr::new(callee),
6534 param_count,
6535 1,
6536 is_concurrent,
6537 callback,
6538 post_return.map(SendSyncPtr::new),
6539 true,
6540 )
6541 }
6542}