1use self::error_contexts::GlobalErrorContextRefCount;
54use crate::component::func::{Func, call_post_return};
55use crate::component::{
56 HasData, HasSelf, Instance, Resource, ResourceTable, ResourceTableError, RuntimeInstance,
57};
58use crate::fiber::{self, StoreFiber, StoreFiberYield};
59use crate::hash_set::HashSet;
60#[cfg(feature = "gc")]
61use crate::module::ModuleRegistry;
62use crate::prelude::*;
63use crate::store::{Store, StoreId, StoreInner, StoreOpaque, StoreToken};
64#[cfg(feature = "gc")]
65use crate::vm::GcRootsList;
66use crate::vm::component::{CallContext, ComponentInstance, CurrentScope, InstanceState, Scope};
67use crate::vm::{
68 AlwaysMut, SendSyncPtr, UncaughtException, VMFuncRef, VMLazyThread, VMMemoryDefinition, VMStore,
69};
70use crate::{
71 AsContext, AsContextMut, FuncType, Result, StoreContext, StoreContextMut, ValRaw, ValType, bail,
72};
73use crate::{Instance as ModuleInstance, bail_bug};
74use alloc::borrow::ToOwned;
75use alloc::collections::{BTreeMap, BTreeSet, VecDeque};
76use core::any::Any;
77use core::cell::UnsafeCell;
78use core::fmt;
79use core::future;
80use core::future::Future;
81use core::marker::PhantomData;
82use core::mem::{self, ManuallyDrop, MaybeUninit};
83use core::ops::DerefMut;
84use core::pin::{Pin, pin};
85use core::ptr::{self, NonNull};
86use core::task::{Context, Poll, Waker};
87use futures::channel::oneshot;
88use futures::stream::{FuturesUnordered, StreamExt};
89use futures_and_streams::{FlatAbi, ReturnCode, TransmitHandle, TransmitIndex};
90use table::{TableDebug, TableId};
91use wasmtime_environ::component::{
92 CanonicalAbiInfo, CanonicalOptions, CanonicalOptionsDataModel, MAX_FLAT_PARAMS,
93 MAX_FLAT_RESULTS, OptionsIndex, PREPARE_ASYNC_NO_RESULT, PREPARE_ASYNC_WITH_RESULT,
94 RuntimeComponentInstanceIndex, RuntimeTableIndex, StringEncoding,
95 TypeComponentGlobalErrorContextTableIndex, TypeComponentLocalErrorContextTableIndex,
96 TypeFuncIndex, TypeFutureTableIndex, TypeStreamTableIndex, TypeTupleIndex,
97};
98use wasmtime_environ::packed_option::ReservedValue;
99use wasmtime_environ::{NUM_COMPONENT_CONTEXT_SLOTS, Trap};
100#[cfg(feature = "gc")]
101use wasmtime_unwinder::Unwind;
102
103pub use abort::JoinHandle;
104pub use func::{FuncCallConcurrent, TypedFuncCallConcurrent};
105pub use future_stream_any::{FutureAny, StreamAny};
106pub use futures_and_streams::{
107 Destination, DirectDestination, DirectSource, ErrorContext, FutureConsumer, FutureProducer,
108 FutureReader, GuardedFutureReader, GuardedStreamReader, ReadBuffer, Source, StreamConsumer,
109 StreamProducer, StreamReader, StreamResult, VecBuffer, WriteBuffer,
110};
111pub(crate) use futures_and_streams::{ResourcePair, lower_error_context_to_index};
112#[cfg(feature = "task-group-hook")]
113pub use task_group_hook::TaskGroupHook;
114pub use task_group_hook::TaskGroupId;
115
116mod abort;
117mod error_contexts;
118mod func;
119mod future_stream_any;
120mod futures_and_streams;
121pub(crate) mod table;
122#[cfg(feature = "task-group-hook")]
123mod task_group_hook;
124#[cfg(not(feature = "task-group-hook"))]
125mod task_group_hook_disabled;
126#[cfg(not(feature = "task-group-hook"))]
127use task_group_hook_disabled as task_group_hook;
128pub(crate) mod tls;
129
130const BLOCKED: u32 = 0xffff_ffff;
133
134#[derive(Clone, Copy, Eq, PartialEq, Debug)]
136pub enum Status {
137 Starting = 0,
138 Started = 1,
139 Returned = 2,
140 StartCancelled = 3,
141 ReturnCancelled = 4,
142}
143
144impl Status {
145 pub fn pack(self, waitable: Option<u32>) -> u32 {
151 assert!(matches!(self, Status::Returned) == waitable.is_none());
152 let waitable = waitable.unwrap_or(0);
153 assert!(waitable < (1 << 28));
154 (waitable << 4) | (self as u32)
155 }
156}
157
158#[derive(Clone, Copy, Debug)]
161enum Event {
162 None,
163 Subtask {
164 status: Status,
165 },
166 StreamRead {
167 code: ReturnCode,
168 pending: Option<(TypeStreamTableIndex, u32)>,
169 },
170 StreamWrite {
171 code: ReturnCode,
172 pending: Option<(TypeStreamTableIndex, u32)>,
173 },
174 FutureRead {
175 code: ReturnCode,
176 pending: Option<(TypeFutureTableIndex, u32)>,
177 },
178 FutureWrite {
179 code: ReturnCode,
180 pending: Option<(TypeFutureTableIndex, u32)>,
181 },
182 Cancelled,
183}
184
185impl Event {
186 fn parts(self) -> (u32, u32) {
191 const EVENT_NONE: u32 = 0;
192 const EVENT_SUBTASK: u32 = 1;
193 const EVENT_STREAM_READ: u32 = 2;
194 const EVENT_STREAM_WRITE: u32 = 3;
195 const EVENT_FUTURE_READ: u32 = 4;
196 const EVENT_FUTURE_WRITE: u32 = 5;
197 const EVENT_CANCELLED: u32 = 6;
198 match self {
199 Event::None => (EVENT_NONE, 0),
200 Event::Cancelled => (EVENT_CANCELLED, 0),
201 Event::Subtask { status } => (EVENT_SUBTASK, status as u32),
202 Event::StreamRead { code, .. } => (EVENT_STREAM_READ, code.encode()),
203 Event::StreamWrite { code, .. } => (EVENT_STREAM_WRITE, code.encode()),
204 Event::FutureRead { code, .. } => (EVENT_FUTURE_READ, code.encode()),
205 Event::FutureWrite { code, .. } => (EVENT_FUTURE_WRITE, code.encode()),
206 }
207 }
208}
209
210mod callback_code {
212 pub const EXIT: u32 = 0;
213 pub const YIELD: u32 = 1;
214 pub const WAIT: u32 = 2;
215}
216
217const START_FLAG_ASYNC_CALLEE: u32 = wasmtime_environ::component::START_FLAG_ASYNC_CALLEE as u32;
219const START_FLAG_ASYNC_CALLER: u32 = wasmtime_environ::component::START_FLAG_ASYNC_CALLER as u32;
220
221pub struct Access<'a, T: 'static, D: HasData + ?Sized = HasSelf<T>> {
227 store: StoreContextMut<'a, T>,
228 get_data: fn(&mut T) -> D::Data<'_>,
229}
230
231impl<'a, T, D> Access<'a, T, D>
232where
233 D: HasData + ?Sized,
234 T: 'static,
235{
236 pub fn new(store: StoreContextMut<'a, T>, get_data: fn(&mut T) -> D::Data<'_>) -> Self {
238 Self { store, get_data }
239 }
240
241 pub fn data_mut(&mut self) -> &mut T {
243 self.store.data_mut()
244 }
245
246 pub fn get(&mut self) -> D::Data<'_> {
248 (self.get_data)(self.data_mut())
249 }
250
251 pub fn spawn(&mut self, task: impl for<'fut> AccessorTask<'fut, T, D>) -> Result<JoinHandle>
255 where
256 T: 'static,
257 {
258 let accessor = Accessor {
259 get_data: self.get_data,
260 token: StoreToken::new(self.store.as_context_mut()),
261 };
262 self.store
263 .as_context_mut()
264 .spawn_with_accessor(accessor, task)
265 }
266
267 pub fn getter(&self) -> fn(&mut T) -> D::Data<'_> {
270 self.get_data
271 }
272}
273
274impl<'a, T, D> AsContext for Access<'a, T, D>
275where
276 D: HasData + ?Sized,
277 T: 'static,
278{
279 type Data = T;
280
281 fn as_context(&self) -> StoreContext<'_, T> {
282 self.store.as_context()
283 }
284}
285
286impl<'a, T, D> AsContextMut for Access<'a, T, D>
287where
288 D: HasData + ?Sized,
289 T: 'static,
290{
291 fn as_context_mut(&mut self) -> StoreContextMut<'_, T> {
292 self.store.as_context_mut()
293 }
294}
295
296pub struct Accessor<T: 'static, D = HasSelf<T>>
356where
357 D: HasData + ?Sized,
358{
359 token: StoreToken<T>,
360 get_data: fn(&mut T) -> D::Data<'_>,
361}
362
363pub trait AsAccessor {
380 type Data: 'static;
382
383 type AccessorData: HasData + ?Sized;
386
387 fn as_accessor(&self) -> &Accessor<Self::Data, Self::AccessorData>;
389}
390
391impl<T: AsAccessor + ?Sized> AsAccessor for &T {
392 type Data = T::Data;
393 type AccessorData = T::AccessorData;
394
395 fn as_accessor(&self) -> &Accessor<Self::Data, Self::AccessorData> {
396 T::as_accessor(self)
397 }
398}
399
400impl<T, D: HasData + ?Sized> AsAccessor for Accessor<T, D> {
401 type Data = T;
402 type AccessorData = D;
403
404 fn as_accessor(&self) -> &Accessor<T, D> {
405 self
406 }
407}
408
409const _: () = {
432 const fn assert<T: Send + Sync>() {}
433 assert::<Accessor<UnsafeCell<u32>>>();
434};
435
436impl<T> Accessor<T> {
437 pub(crate) fn new(token: StoreToken<T>) -> Self {
446 Self {
447 token,
448 get_data: |x| x,
449 }
450 }
451}
452
453impl<T, D> Accessor<T, D>
454where
455 D: HasData + ?Sized,
456{
457 pub fn with<R>(&self, fun: impl FnOnce(Access<'_, T, D>) -> R) -> R {
475 tls::get(|vmstore| {
476 fun(Access {
477 store: self.token.as_context_mut(vmstore),
478 get_data: self.get_data,
479 })
480 })
481 }
482
483 pub fn getter(&self) -> fn(&mut T) -> D::Data<'_> {
486 self.get_data
487 }
488
489 pub fn with_getter<D2: HasData>(
506 &self,
507 get_data: fn(&mut T) -> D2::Data<'_>,
508 ) -> Accessor<T, D2> {
509 Accessor {
510 token: self.token,
511 get_data,
512 }
513 }
514
515 pub fn spawn(&self, task: impl for<'fut> AccessorTask<'fut, T, D>) -> Result<JoinHandle>
531 where
532 T: 'static,
533 {
534 let accessor = self.clone_for_spawn();
535 self.with(|mut access| access.as_context_mut().spawn_with_accessor(accessor, task))
536 }
537
538 fn clone_for_spawn(&self) -> Self {
539 Self {
540 token: self.token,
541 get_data: self.get_data,
542 }
543 }
544
545 pub fn poll_no_interesting_tasks(&self, cx: &mut Context<'_>) -> Poll<()> {
581 self.with(|mut access| {
582 let store = access.as_context_mut().0;
583 let state = store.concurrent_state_mut_without_forcing_current_thread();
584 if state.interesting_tasks == 0 {
585 Poll::Ready(())
586 } else {
587 state.interesting_tasks_empty_waker = Some(cx.waker().clone());
588 Poll::Pending
589 }
590 })
591 }
592
593 pub fn poll_ready_for_concurrent_call(&self, func: Func, cx: &mut Context<'_>) -> Poll<()> {
610 self.with(|mut access| {
611 let store = access.as_context_mut().0;
612 let (_, _, _, raw_options) = func.abi_info(store);
613 let instance = func.instance().runtime_instance(raw_options.instance);
614 let state = store.instance_state(instance).concurrent_state();
615 if state.backpressure == 0 {
616 Poll::Ready(())
617 } else {
618 store
619 .concurrent_state_mut_without_forcing_current_thread()
620 .ready_for_concurrent_call_waker = Some(cx.waker().clone());
621 Poll::Pending
622 }
623 })
624 }
625}
626
627pub trait AccessorTask<'fut, T, D = HasSelf<T>>:
649 AsyncFnOnce(&Accessor<T, D>) -> Result<()> + Send + 'static
650where
651 D: HasData + ?Sized,
652{
653 fn run(self, accessor: &'fut Accessor<T, D>) -> impl Future<Output = Result<()>> + Send + 'fut;
655}
656
657impl<'fut, F, Fut, T, D> AccessorTask<'fut, T, D> for F
658where
659 T: 'static,
660 F: AsyncFnOnce(&Accessor<T, D>) -> Result<()>,
661 F: FnOnce(&'fut Accessor<T, D>) -> Fut + Send + 'static,
662 Fut: Future<Output = Result<()>> + Send + 'fut,
663 D: HasData,
664{
665 fn run(self, accessor: &'fut Accessor<T, D>) -> impl Future<Output = Result<()>> + Send + 'fut {
666 (self)(accessor)
667 }
668}
669
670enum CallerInfo {
673 Async {
675 params: Vec<ValRaw>,
676 has_result: bool,
677 },
678 Sync {
680 params: Vec<ValRaw>,
681 result_count: u32,
682 },
683}
684
685enum WaitMode {
687 Fiber(StoreFiber<'static>),
689 Callback(Instance),
692}
693
694impl fmt::Debug for WaitMode {
695 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
696 match self {
697 Self::Fiber(_) => f.debug_tuple("Fiber").finish(),
698 Self::Callback(instance) => f.debug_tuple("Callback").field(instance).finish(),
699 }
700 }
701}
702
703#[derive(Debug)]
705enum SuspendReason {
706 Waiting {
709 set: TableId<WaitableSet>,
710 thread: QualifiedThreadId,
711 },
712 YieldingToSubtask { thread: QualifiedThreadId },
715 NeedWork,
718 Yielding { thread: QualifiedThreadId },
721 ExplicitlySuspending { thread: QualifiedThreadId },
724}
725
726enum GuestCallKind {
728 DeliverEvent {
731 instance: Instance,
733 set: Option<TableId<WaitableSet>>,
738 },
739 StartImplicit(Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send + Sync>),
745 StartExplicit(Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send + Sync>),
746}
747
748impl fmt::Debug for GuestCallKind {
749 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
750 match self {
751 Self::DeliverEvent { instance, set } => f
752 .debug_struct("DeliverEvent")
753 .field("instance", instance)
754 .field("set", set)
755 .finish(),
756 Self::StartImplicit(_) => f.debug_tuple("StartImplicit").finish(),
757 Self::StartExplicit(_) => f.debug_tuple("StartExplicit").finish(),
758 }
759 }
760}
761
762#[derive(Copy, Clone, Debug)]
764pub enum SuspensionTarget {
765 Resume(u32),
766 Promote(u32),
767 None,
768}
769
770#[derive(Copy, Clone, Debug)]
772pub enum ResumeThread {
773 Promote,
774 Resume,
775 ResumeLater,
776}
777
778#[derive(Debug)]
780struct GuestCall {
781 thread: QualifiedThreadId,
782 kind: GuestCallKind,
783}
784
785impl GuestCall {
786 fn is_ready(&self, store: &mut StoreOpaque) -> Result<bool> {
796 let task = store.concurrent_state_mut()?.get_mut(self.thread.task)?;
797 let async_typed = task.async_typed;
798 let instance = task.instance;
799 let state = store.instance_state(instance).concurrent_state();
800
801 let ready = match &self.kind {
802 GuestCallKind::DeliverEvent { .. } => !state.do_not_enter,
803 GuestCallKind::StartImplicit(_) => {
804 !async_typed || !(state.do_not_enter || state.backpressure > 0)
805 }
806 GuestCallKind::StartExplicit(_) => true,
807 };
808 log::trace!(
809 "call {self:?} ready? {ready} (do_not_enter: {}; backpressure: {})",
810 state.do_not_enter,
811 state.backpressure
812 );
813 Ok(ready)
814 }
815}
816
817enum WorkerItem {
819 GuestCall(GuestCall),
820 Function(AlwaysMut<Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send>>),
821}
822
823enum WorkItem {
826 PushFuture(AlwaysMut<HostTaskFuture>),
828 ResumeFiber {
830 instance: RuntimeInstance,
831 thread: QualifiedThreadId,
832 fiber: StoreFiber<'static>,
833 },
834 ResumeThread {
836 instance: RuntimeInstance,
837 thread: QualifiedThreadId,
838 },
839 GuestCall {
841 instance: RuntimeInstance,
842 call: GuestCall,
843 },
844 WorkerFunction(AlwaysMut<Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send>>),
846}
847
848impl fmt::Debug for WorkItem {
849 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
850 match self {
851 Self::PushFuture(_) => f.debug_tuple("PushFuture").finish(),
852 Self::ResumeFiber {
853 instance, thread, ..
854 } => f
855 .debug_struct("ResumeFiber")
856 .field("instance", instance)
857 .field("thread", thread)
858 .finish(),
859 Self::ResumeThread { instance, thread } => f
860 .debug_struct("ResumeThread")
861 .field("instance", instance)
862 .field("thread", thread)
863 .finish(),
864 Self::GuestCall { instance, call } => f
865 .debug_struct("GuestCall")
866 .field("instance", instance)
867 .field("call", call)
868 .finish(),
869 Self::WorkerFunction(_) => f.debug_tuple("WorkerFunction").finish(),
870 }
871 }
872}
873
874#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
876pub(crate) enum WaitResult {
877 Cancelled,
878 Completed,
879}
880
881pub(crate) fn poll_and_block<R: Send + Sync + 'static>(
889 store: &mut dyn VMStore,
890 host_task: EnteredHostTask,
891 future: impl Future<Output = Result<R>> + Send + 'static,
892) -> Result<R> {
893 let mut future = Box::pin(future);
900 let poll = tls::set(store, || {
901 future
902 .as_mut()
903 .poll(&mut Context::from_waker(&Waker::noop()))
904 });
905
906 let caller = match host_task {
907 Some(caller) => caller,
908 None => bail_bug!("host task wasn't created but should have been"),
909 };
910
911 let task = match poll {
912 Poll::Ready(result) => return result,
914
915 Poll::Pending => {
920 let Some(task) = store.materialize_host_task_id()? else {
921 bail_bug!("current thread is not a host thread")
922 };
923
924 let future = Box::pin(async move {
927 let result = run_with_host_task_set(task, future).await??;
928 tls::get(move |store| {
929 let state = store.concurrent_state_mut()?;
930 let host_state = &mut state.get_mut(task)?.state;
931 assert!(matches!(host_state, HostTaskState::CalleeStarted));
932 *host_state = HostTaskState::CalleeFinished(Box::new(result));
933
934 Waitable::Host(task).set_event(
935 state,
936 Some(Event::Subtask {
937 status: Status::Returned,
938 }),
939 )?;
940
941 Ok(())
942 })
943 }) as HostTaskFuture;
944
945 let caller_instance = store.concurrent_state_mut()?.get_mut(caller.task)?.instance;
946 store.switch_or_trap_if_may_not_suspend(caller_instance)?;
947
948 let state = store.concurrent_state_mut()?;
949 state.push_future(future);
950
951 let set = state.get_mut(caller.thread)?.sync_call_set;
952 Waitable::Host(task).join(state, Some(set))?;
953
954 store.suspend(SuspendReason::Waiting {
955 set,
956 thread: caller,
957 })?;
958
959 Waitable::Host(task).join(store.concurrent_state_mut()?, None)?;
963 task
964 }
965 };
966
967 let host_state = &mut store.concurrent_state_mut()?.get_mut(task)?.state;
969 match mem::replace(host_state, HostTaskState::CalleeDone { cancelled: false }) {
970 HostTaskState::CalleeFinished(result) => Ok(match result.downcast() {
971 Ok(result) => *result,
972 Err(_) => bail_bug!("host task finished with wrong type of result"),
973 }),
974 _ => bail_bug!("unexpected host task state after completion"),
975 }
976}
977
978fn handle_guest_call(store: &mut dyn VMStore, call: GuestCall) -> Result<()> {
980 match call.kind {
981 GuestCallKind::DeliverEvent { instance, set } => {
982 if let Some(set) = set {
985 store.concurrent_state_mut()?.get_mut(set)?.stop_waiting()?;
986 }
987 let (event, waitable) = match instance.get_event(store, call.thread.task, set, true)? {
988 Some(pair) => pair,
989 None => match set {
990 Some(set) => {
995 log::trace!(
996 "event for {:?} on {set:?} no longer present; waiting again",
997 call.thread
998 );
999 return instance.wait_with_callback(
1000 store.store_opaque_mut(),
1001 call.thread,
1002 set,
1003 );
1004 }
1005 None => (Event::None, None),
1010 },
1011 };
1012 let state = store.concurrent_state_mut()?;
1013 let task = state.get_mut(call.thread.task)?;
1014 let runtime_instance = task.instance;
1015 let handle = waitable.map(|(_, v)| v).unwrap_or(0);
1016
1017 log::trace!(
1018 "use callback to deliver event {event:?} to {:?} for {waitable:?}",
1019 call.thread,
1020 );
1021
1022 let old_thread = store.set_thread(call.thread)?;
1023 log::trace!(
1024 "GuestCallKind::DeliverEvent: replaced {old_thread:?} with {:?} as current thread",
1025 call.thread
1026 );
1027
1028 store.enter_instance(runtime_instance);
1029
1030 let Some(callback) = store
1031 .concurrent_state_mut()?
1032 .get_mut(call.thread.task)?
1033 .callback
1034 .take()
1035 else {
1036 bail_bug!("guest task callback field not present")
1037 };
1038
1039 let code = callback(store, event, handle)?;
1040
1041 store
1042 .concurrent_state_mut()?
1043 .get_mut(call.thread.task)?
1044 .callback = Some(callback);
1045
1046 store.exit_instance(runtime_instance)?;
1047
1048 store.set_thread(old_thread)?;
1049
1050 instance.handle_callback_code(store, call.thread, runtime_instance.index, code)?;
1051
1052 log::trace!("GuestCallKind::DeliverEvent: restored {old_thread:?} as current thread");
1053 }
1054 GuestCallKind::StartImplicit(fun) => {
1055 fun(store)?;
1056 }
1057 GuestCallKind::StartExplicit(fun) => {
1058 fun(store)?;
1059 }
1060 }
1061
1062 Ok(())
1063}
1064
1065impl<T> Store<T> {
1066 pub async fn run_concurrent<R>(&mut self, fun: impl AsyncFnOnce(&Accessor<T>) -> R) -> Result<R>
1068 where
1069 T: Send + 'static,
1070 {
1071 ensure!(
1072 self.as_context().0.concurrency_support(),
1073 "cannot use `run_concurrent` when Config::concurrency_support disabled",
1074 );
1075 self.as_context_mut().run_concurrent(fun).await
1076 }
1077
1078 #[doc(hidden)]
1079 pub fn assert_concurrent_state_empty(&mut self) {
1080 self.as_context_mut().assert_concurrent_state_empty();
1081 }
1082
1083 #[doc(hidden)]
1084 pub fn concurrent_state_table_size(&mut self) -> usize {
1085 self.as_context_mut().concurrent_state_table_size()
1086 }
1087
1088 pub fn spawn(
1090 &mut self,
1091 task: impl for<'fut> AccessorTask<'fut, T, HasSelf<T>>,
1092 ) -> Result<JoinHandle>
1093 where
1094 T: 'static,
1095 {
1096 self.as_context_mut().spawn(task)
1097 }
1098}
1099
1100impl<T> StoreContextMut<'_, T> {
1101 #[doc(hidden)]
1112 pub fn assert_concurrent_state_empty(self) {
1113 let store = self.0;
1114 store
1115 .store_data_mut()
1116 .components
1117 .assert_instance_states_empty();
1118 let state = store.concurrent_state_mut().unwrap();
1119 assert!(
1120 state.table.get_mut().is_empty(),
1121 "non-empty table: {:?}",
1122 state.table.get_mut()
1123 );
1124 assert!(state.switch_item.is_none());
1125 assert!(state.next_switch_item.is_none());
1126 assert!(state.high_priority.is_empty());
1127 assert!(state.low_priority.is_empty());
1128 assert!(state.unforced_current_thread.is_none());
1129 assert!(state.deferred_host_call_context.is_none());
1130 assert!(state.futures_mut().unwrap().is_empty());
1131 assert!(state.global_error_context_ref_counts.is_empty());
1132 }
1133
1134 #[doc(hidden)]
1139 pub fn concurrent_state_table_size(&mut self) -> usize {
1140 self.0
1141 .concurrent_state_mut()
1142 .unwrap()
1143 .table
1144 .get_mut()
1145 .iter_mut()
1146 .count()
1147 }
1148
1149 pub fn spawn(mut self, task: impl for<'fut> AccessorTask<'fut, T>) -> Result<JoinHandle>
1159 where
1160 T: 'static,
1161 {
1162 let accessor = Accessor::new(StoreToken::new(self.as_context_mut()));
1163 self.spawn_with_accessor(accessor, task)
1164 }
1165
1166 fn spawn_with_accessor<D>(
1169 self,
1170 accessor: Accessor<T, D>,
1171 task: impl for<'fut> AccessorTask<'fut, T, D>,
1172 ) -> Result<JoinHandle>
1173 where
1174 T: 'static,
1175 D: HasData + ?Sized,
1176 {
1177 let (handle, future) = JoinHandle::run(async move { task.run(&accessor).await });
1181 self.0
1182 .concurrent_state_mut()?
1183 .push_future(Box::pin(async move { future.await.unwrap_or(Ok(())) }));
1184 Ok(handle)
1185 }
1186
1187 pub async fn run_concurrent<R>(self, fun: impl AsyncFnOnce(&Accessor<T>) -> R) -> Result<R>
1271 where
1272 T: Send + 'static,
1273 {
1274 ensure!(
1275 self.0.concurrency_support(),
1276 "cannot use `run_concurrent` when Config::concurrency_support disabled",
1277 );
1278 self.do_run_concurrent(fun, false).await
1279 }
1280
1281 pub(super) async fn run_concurrent_trap_on_idle<R>(
1282 self,
1283 fun: impl AsyncFnOnce(&Accessor<T>) -> R,
1284 ) -> Result<R> {
1285 self.do_run_concurrent(fun, true).await
1286 }
1287
1288 async fn do_run_concurrent<R>(
1289 mut self,
1290 fun: impl AsyncFnOnce(&Accessor<T>) -> R,
1291 trap_on_idle: bool,
1292 ) -> Result<R> {
1293 debug_assert!(self.0.concurrency_support());
1294 let already_running = self
1295 .0
1296 .concurrent_state_mut_already_forced_current_thread()
1297 .event_loop_running;
1298 if already_running {
1299 bail!("Recursive `StoreContextMut::run_concurrent` calls not supported")
1300 }
1301 let token = StoreToken::new(self.as_context_mut());
1302
1303 struct Dropper<'a, T: 'static, V> {
1304 store: StoreContextMut<'a, T>,
1305 value: ManuallyDrop<V>,
1306 }
1307
1308 impl<'a, T, V> Drop for Dropper<'a, T, V> {
1309 fn drop(&mut self) {
1310 self.store
1311 .0
1312 .concurrent_state_mut_already_forced_current_thread()
1313 .event_loop_running = false;
1314
1315 tls::set(self.store.0, || {
1316 unsafe { ManuallyDrop::drop(&mut self.value) }
1321 });
1322 }
1323 }
1324
1325 let accessor = &Accessor::new(token);
1326 self.0
1327 .concurrent_state_mut_already_forced_current_thread()
1328 .event_loop_running = true;
1329 let dropper = &mut Dropper {
1330 store: self,
1331 value: ManuallyDrop::new(fun(accessor)),
1332 };
1333 let future = unsafe { Pin::new_unchecked(dropper.value.deref_mut()) };
1335
1336 let result = dropper
1337 .store
1338 .as_context_mut()
1339 .poll_until(future, trap_on_idle)
1340 .await;
1341
1342 if result.is_err() {
1343 dropper.store.0.set_trapped();
1344 }
1345
1346 result
1347 }
1348
1349 async fn poll_until<R>(
1355 mut self,
1356 mut future: Pin<&mut impl Future<Output = R>>,
1357 trap_on_idle: bool,
1358 ) -> Result<R> {
1359 struct Reset<'a, T: 'static> {
1360 store: StoreContextMut<'a, T>,
1361 futures: Option<FuturesUnordered<HostTaskFuture>>,
1362 }
1363
1364 impl<'a, T> Drop for Reset<'a, T> {
1365 fn drop(&mut self) {
1366 if let Some(futures) = self.futures.take() {
1367 *self
1368 .store
1369 .0
1370 .concurrent_state_mut_already_forced_current_thread()
1371 .futures
1372 .get_mut() = Some(futures);
1373 }
1374 }
1375 }
1376
1377 const MAX_TURNS_WITHOUT_YIELD: usize = 128;
1381 let mut turns_without_yield = 0;
1382
1383 loop {
1384 let futures = self.0.concurrent_state_mut()?.futures.get_mut().take();
1388 let mut reset = Reset {
1389 store: self.as_context_mut(),
1390 futures,
1391 };
1392 let mut next = match reset.futures.as_mut() {
1393 Some(f) => pin!(f.next()),
1394 None => bail_bug!("concurrent state missing futures field"),
1395 };
1396
1397 enum PollResult<R> {
1398 Complete(R),
1399 ProcessWork {
1400 ready: Option<WorkItem>,
1401 low_priority: bool,
1402 },
1403 }
1404
1405 let result = future::poll_fn(|cx| {
1406 if let Poll::Ready(value) = tls::set(reset.store.0, || future.as_mut().poll(cx)) {
1409 return Poll::Ready(Ok(PollResult::Complete(value)));
1410 }
1411
1412 if reset.store.0.trapped() {
1421 return Poll::Ready(Err(Trap::CannotEnterComponent.into()));
1422 }
1423
1424 let next = match tls::set(reset.store.0, || next.as_mut().poll(cx)) {
1428 Poll::Ready(Some(output)) => {
1429 match output {
1430 Err(e) => return Poll::Ready(Err(e)),
1431 Ok(()) => {}
1432 }
1433 Poll::Ready(true)
1434 }
1435 Poll::Ready(None) => Poll::Ready(false),
1436 Poll::Pending => Poll::Pending,
1437 };
1438
1439 let state = reset.store.0.concurrent_state_mut()?;
1454 let mut ready = state.switch_item.take();
1455 let mut low_priority = false;
1456 if ready.is_none() {
1457 ready = state.high_priority.pop_back();
1458 if ready.is_none() {
1459 ready = state.low_priority.pop_back();
1460 low_priority = true;
1461 }
1462 }
1463 if ready.is_some() {
1464 return Poll::Ready(Ok(PollResult::ProcessWork {
1465 ready,
1466 low_priority,
1467 }));
1468 }
1469
1470 return match next {
1474 Poll::Ready(true) => {
1475 Poll::Ready(Ok(PollResult::ProcessWork {
1481 ready: None,
1482 low_priority: false,
1483 }))
1484 }
1485 Poll::Ready(false) => {
1486 if let Poll::Ready(value) =
1490 tls::set(reset.store.0, || future.as_mut().poll(cx))
1491 {
1492 Poll::Ready(Ok(PollResult::Complete(value)))
1493 } else {
1494 if trap_on_idle {
1500 Poll::Ready(Err(if reset.store.0.any_may_not_suspend()? {
1507 Trap::CannotBlockSyncTask.into()
1508 } else {
1509 Trap::AsyncDeadlock.into()
1511 }))
1512 } else {
1513 Poll::Pending
1517 }
1518 }
1519 }
1520 Poll::Pending => Poll::Pending,
1525 };
1526 })
1527 .await;
1528
1529 drop(reset);
1533
1534 match result? {
1535 PollResult::Complete(value) => break Ok(value),
1538 PollResult::ProcessWork {
1541 ready,
1542 low_priority,
1543 } => {
1544 struct Dispose<'a, T: 'static> {
1545 store: StoreContextMut<'a, T>,
1546 ready: Option<WorkItem>,
1547 }
1548
1549 impl<'a, T> Drop for Dispose<'a, T> {
1550 fn drop(&mut self) {
1551 if let Some(item) = self.ready.take() {
1552 match item {
1553 WorkItem::ResumeFiber { mut fiber, .. } => {
1554 fiber.dispose(self.store.0);
1555 }
1556 WorkItem::PushFuture(future) => {
1557 tls::set(self.store.0, move || drop(future))
1558 }
1559 _ => {}
1560 }
1561 }
1562 }
1563 }
1564
1565 let mut dispose = Dispose {
1566 store: self.as_context_mut(),
1567 ready,
1568 };
1569
1570 if low_priority {
1592 dispose.store.0.yield_now().await;
1593 turns_without_yield = 0;
1594 }
1595
1596 if let Some(item) = dispose.ready.take() {
1597 dispose
1598 .store
1599 .as_context_mut()
1600 .handle_work_item(item)
1601 .await?;
1602 }
1603
1604 turns_without_yield += 1;
1605 if turns_without_yield == MAX_TURNS_WITHOUT_YIELD {
1606 turns_without_yield = 0;
1607 dispose.store.0.yield_now().await;
1608 }
1609 }
1610 }
1611 }
1612 }
1613
1614 async fn handle_work_item(self, item: WorkItem) -> Result<()> {
1616 log::trace!("handle work item {item:?}");
1617 match item {
1618 WorkItem::PushFuture(future) => {
1619 self.0
1620 .concurrent_state_mut()?
1621 .futures_mut()?
1622 .push(future.into_inner());
1623 }
1624 WorkItem::ResumeFiber { fiber, .. } => {
1625 self.0.resume_fiber(fiber).await?;
1626 }
1627 WorkItem::ResumeThread { thread, .. } => {
1628 if let GuestThreadState::Ready { fiber, .. } = mem::replace(
1629 &mut self.0.concurrent_state_mut()?.get_mut(thread.thread)?.state,
1630 GuestThreadState::Running,
1631 ) {
1632 self.0.resume_fiber(fiber).await?;
1633 } else {
1634 bail_bug!("cannot resume non-pending thread {thread:?}");
1635 }
1636 }
1637 WorkItem::GuestCall { call, .. } => {
1638 if call.is_ready(self.0)? {
1639 self.0
1640 .concurrent_state_mut()?
1641 .get_mut(call.thread.thread)?
1642 .wake_on_cancel = WakeOnCancel::None;
1643 self.run_on_worker(WorkerItem::GuestCall(call)).await?;
1644 } else {
1645 let state = self.0.concurrent_state_mut()?;
1646 let task = state.get_mut(call.thread.task)?;
1647 if !task.starting_sent {
1648 task.starting_sent = true;
1649 if let GuestCallKind::StartImplicit(_) = &call.kind {
1650 Waitable::Guest(call.thread.task).set_event(
1651 state,
1652 Some(Event::Subtask {
1653 status: Status::Starting,
1654 }),
1655 )?;
1656 }
1657 }
1658
1659 let instance = state.get_mut(call.thread.task)?.instance;
1660 self.0
1661 .instance_state(instance)
1662 .concurrent_state()
1663 .pending
1664 .insert(call.thread, call.kind);
1665
1666 self.0.concurrent_state_mut()?.take_next_switch_item()?;
1670 }
1671 }
1672 WorkItem::WorkerFunction(fun) => {
1673 self.run_on_worker(WorkerItem::Function(fun)).await?;
1674 }
1675 }
1676
1677 Ok(())
1678 }
1679
1680 async fn run_on_worker(self, item: WorkerItem) -> Result<()> {
1682 let worker = if let Some(fiber) = self.0.concurrent_state_mut()?.worker.take() {
1683 fiber
1684 } else {
1685 unsafe {
1704 fiber::make_fiber_unchecked(self.0, move |store| {
1705 loop {
1706 let Some(item) = store.concurrent_state_mut()?.worker_item.take() else {
1707 bail_bug!("worker_item not present when resuming fiber")
1708 };
1709 match item {
1710 WorkerItem::GuestCall(call) => handle_guest_call(store, call)?,
1711 WorkerItem::Function(fun) => fun.into_inner()(store)?,
1712 }
1713
1714 store.suspend(SuspendReason::NeedWork)?;
1715 }
1716 })?
1717 }
1718 };
1719
1720 let worker_item = &mut self.0.concurrent_state_mut()?.worker_item;
1721 assert!(worker_item.is_none());
1722 *worker_item = Some(item);
1723
1724 self.0.resume_fiber(worker).await
1725 }
1726
1727 pub(crate) fn wrap_call<F, R>(self, closure: F) -> impl Future<Output = Result<R>> + 'static
1732 where
1733 T: 'static,
1734 F: FnOnce(&Accessor<T>) -> Pin<Box<dyn Future<Output = Result<R>> + Send + '_>>
1735 + Send
1736 + Sync
1737 + 'static,
1738 R: Send + Sync + 'static,
1739 {
1740 let token = StoreToken::new(self);
1741 async move {
1742 let mut accessor = Accessor::new(token);
1743 closure(&mut accessor).await
1744 }
1745 }
1746
1747 pub(crate) async fn start_instance(
1748 &mut self,
1749 instance: ModuleInstance,
1750 callee: Option<RuntimeInstance>,
1751 ) -> Result<ModuleInstance> {
1752 let (tx, rx) = oneshot::channel();
1753 let token = StoreToken::new(self.as_context_mut());
1754 self.0.queue_task(move |store| {
1755 _ = tx.send(
1756 super::instance::start_raw(&mut token.as_context_mut(store), instance, callee)
1757 .map(|()| instance),
1758 );
1759 Ok(())
1760 })?;
1761 self.as_context_mut()
1762 .run_concurrent_trap_on_idle(async |_| {
1763 rx.await
1764 .map_err(|_| format_err!("oneshot channel canceled"))
1765 })
1766 .await??
1767 }
1768}
1769
1770pub type EnteredHostTask = Option<QualifiedThreadId>;
1777
1778impl StoreOpaque {
1779 #[inline]
1783 pub(crate) fn current_thread(&mut self) -> Result<CurrentThread> {
1784 if !self.concurrency_support() {
1786 return Ok(CurrentThread::None);
1787 }
1788
1789 if !self
1792 .vm_store_context_mut()
1793 .current_thread_mut()
1794 .is_deferred()
1795 {
1796 return Ok(self
1797 .concurrent_state_mut_already_forced_current_thread()
1798 .unforced_current_thread);
1799 }
1800
1801 self.force_deferred_current_thread()
1802 }
1803
1804 #[cold]
1807 fn force_deferred_current_thread(&mut self) -> Result<CurrentThread> {
1808 let state = self.concurrent_state_mut_without_forcing_current_thread();
1817 let id = match state.unforced_current_thread.guest_task() {
1818 Some(task) => state.get_mut(task)?.instance.instance,
1819 None => bail_bug!("deferred component-model thread with non-guest base"),
1820 };
1821
1822 let mut frames = Vec::new();
1825 let mut cur = *self.vm_store_context_mut().current_thread_mut();
1826 while let Some(ptr) = cur.as_deferred() {
1827 let deferred = unsafe { ptr.as_non_null().as_ref() };
1832 frames.push((
1833 deferred.callee_async != 0,
1834 deferred.callee_instance,
1835 deferred.saved_context,
1836 ));
1837 cur = deferred.parent;
1838 }
1839
1840 *self.vm_store_context_mut().current_thread_mut() = VMLazyThread::forced();
1844
1845 let current_context = *self.vm_store_context_mut().component_context_mut();
1848
1849 for (callee_async, callee_instance, saved_context) in frames.into_iter().rev() {
1853 *self.vm_store_context_mut().component_context_mut() = saved_context;
1857 let callee = RuntimeInstance {
1858 instance: id,
1859 index: RuntimeComponentInstanceIndex::from_u32(callee_instance),
1860 };
1861 self.enter_guest_sync_call(callee_async, callee)?;
1862 }
1863
1864 *self.vm_store_context_mut().component_context_mut() = current_context;
1866
1867 Ok(self
1868 .concurrent_state_mut_without_forcing_current_thread()
1869 .unforced_current_thread)
1870 }
1871
1872 fn current_guest_thread(&mut self) -> Result<QualifiedThreadId> {
1873 match self.current_thread()?.guest() {
1874 Some(id) => Ok(*id),
1875 None => bail_bug!("current thread is not a guest thread"),
1876 }
1877 }
1878
1879 pub(crate) fn current_materialized_host_task(&mut self) -> Result<Option<TableId<HostTask>>> {
1883 match self.current_thread()? {
1884 CurrentThread::Host(id) => Ok(Some(id)),
1885 CurrentThread::DeferredHost(_) | CurrentThread::None => Ok(None),
1886 _ => bail_bug!("current thread is not a host thread"),
1887 }
1888 }
1889
1890 fn materialize_host_task_id(&mut self) -> Result<Option<TableId<HostTask>>> {
1893 Ok(self
1894 .concurrent_state_mut()?
1895 .materialize_current_host_task_id()?)
1896 }
1897
1898 fn enter_sync_call(&mut self, callee: RuntimeInstance) -> Result<()> {
1899 log::trace!("enter sync-typed call {callee:?}");
1900 let state = self.instance_state(callee).concurrent_state();
1901 let old_do_not_suspend = state.do_not_suspend;
1902 state.do_not_suspend = true;
1903
1904 let thread = self.current_guest_thread()?;
1905 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
1906 if thread.old_do_not_suspend.is_some() {
1907 bail_bug!("current thread already has `old_do_not_suspend` value");
1908 }
1909
1910 thread.old_do_not_suspend = Some(old_do_not_suspend);
1911
1912 Ok(())
1913 }
1914
1915 fn exit_sync_call(&mut self, callee: RuntimeInstance) -> Result<()> {
1916 log::trace!("exit sync-typed call {callee:?}");
1917 let thread = self.current_guest_thread()?;
1918 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
1919 let Some(old_do_not_suspend) = thread.old_do_not_suspend.take() else {
1920 bail_bug!("current thread missing `old_do_not_suspend` value");
1921 };
1922 let state = self.instance_state(callee).concurrent_state();
1923 state.do_not_suspend = old_do_not_suspend;
1924 Ok(())
1925 }
1926
1927 pub(crate) fn enter_guest_sync_call(
1939 &mut self,
1940 callee_async_typed: bool,
1941 callee: RuntimeInstance,
1942 ) -> Result<()> {
1943 log::trace!("enter sync-lifted call {callee:?}");
1944 if !self.concurrency_support() {
1945 return self.enter_call_not_concurrent();
1946 }
1947
1948 let thread = self.current_thread()?;
1949 let caller = if let Some(thread) = thread.guest() {
1950 Caller::Guest { thread: *thread }
1951 } else {
1952 Caller::Host {
1953 tx: None,
1954 host_future_present: false,
1955 caller: self.materialize_host_task_id()?,
1956 }
1957 };
1958
1959 let state = self.concurrent_state_mut()?;
1960 let guest_thread = GuestTask::new(
1961 state,
1962 Box::new(move |_, _| bail_bug!("cannot lower params in sync call")),
1963 LiftResult {
1964 lift: Box::new(move |_, _| bail_bug!("cannot lift result in sync call")),
1965 ty: TypeTupleIndex::reserved_value(),
1966 memory: None,
1967 string_encoding: StringEncoding::Utf8,
1968 },
1969 caller,
1970 None,
1971 callee,
1972 callee_async_typed,
1973 false,
1974 )?;
1975
1976 Instance::from_wasmtime(self, callee.instance).add_guest_thread_to_instance_table(
1977 guest_thread.thread,
1978 self,
1979 callee.index,
1980 )?;
1981 self.set_thread(guest_thread)?;
1982
1983 if !callee_async_typed {
1984 self.enter_sync_call(callee)?;
1985 }
1986
1987 Ok(())
1988 }
1989
1990 pub(crate) fn exit_guest_sync_call(&mut self) -> Result<()> {
1998 if !self.concurrency_support() {
1999 return Ok(self.exit_call_not_concurrent());
2000 }
2001
2002 let thread = match self.current_thread()?.guest() {
2003 Some(t) => *t,
2004 None => bail_bug!("expected task when exiting"),
2005 };
2006 let task = self.concurrent_state_mut()?.get_mut(thread.task)?;
2007 let instance = task.instance;
2008
2009 let caller = match &task.caller {
2010 &Caller::Guest { thread } => thread.into(),
2011 &Caller::Host { caller, .. } => caller
2012 .map(CurrentThread::Host)
2013 .unwrap_or(CurrentThread::None),
2014 };
2015 task.lift_result = None;
2016 task.exited = true;
2017 let async_typed = task.async_typed;
2018
2019 if !async_typed {
2020 self.exit_sync_call(instance)?;
2021 }
2022
2023 self.set_thread(caller)?;
2024
2025 log::trace!("exit sync-lifted call {instance:?}");
2026
2027 if async_typed {
2028 self.switch_or_trap_if_may_not_suspend(instance)?;
2033 }
2034
2035 self.cleanup_thread(thread, instance, CleanupTask::Yes)?;
2036
2037 Ok(())
2038 }
2039
2040 pub(crate) fn host_task_create(&mut self) -> Result<EnteredHostTask> {
2047 if !self.concurrency_support() {
2048 self.enter_call_not_concurrent()?;
2049 return Ok(None);
2050 }
2051 let caller = self.current_guest_thread()?;
2052 log::trace!("new deferred host task with caller {caller:?}");
2053
2054 self.set_thread(CurrentThread::DeferredHost(caller))?;
2055 let state = self.concurrent_state_mut()?;
2056 debug_assert!(state.deferred_host_call_context.is_none());
2057 state.deferred_host_call_context = Some(CallContext::default());
2058 state.debug_assert_deferred_host_invariant();
2059 Ok(Some(caller))
2060 }
2061
2062 pub(crate) fn host_task_delete(
2069 &mut self,
2070 original_task: EnteredHostTask,
2071 materialized_task: Option<TableId<HostTask>>,
2072 ) -> Result<()> {
2073 match original_task {
2074 Some(caller) => {
2075 self.set_thread(caller)?;
2076 if materialized_task.is_none() {
2077 let state = self.concurrent_state_mut()?;
2078 let context = state
2079 .deferred_host_call_context
2080 .take()
2081 .expect("deferred host call context should be present");
2082 debug_assert!(context.is_empty());
2083 state.debug_assert_deferred_host_invariant();
2084 }
2085 log::trace!(
2086 "delete host task with caller {original_task:?} and materialized as {materialized_task:?}"
2087 );
2088 if let Some(task) = materialized_task {
2089 Waitable::Host(task).delete_from(self)?;
2090 }
2091 }
2092 None => {
2093 debug_assert!(materialized_task.is_none());
2094 self.exit_call_not_concurrent();
2095 }
2096 }
2097 Ok(())
2098 }
2099
2100 fn instance_state(&mut self, instance: RuntimeInstance) -> &mut InstanceState {
2103 self.component_instance_mut(instance.instance)
2104 .instance_state(instance.index)
2105 }
2106
2107 pub(crate) fn set_thread(&mut self, thread: impl Into<CurrentThread>) -> Result<CurrentThread> {
2113 let thread = thread.into();
2114 let state = self.concurrent_state_mut()?;
2115 state.debug_assert_deferred_host_invariant();
2116 let old_thread = mem::replace(&mut state.unforced_current_thread, thread);
2117
2118 state.handle_thread_switch(old_thread, thread)?;
2119
2120 if let Some(old_thread) = old_thread.guest() {
2128 let old_context = *self.vm_store_context_mut().component_context_mut();
2129 self.concurrent_state_mut()?
2130 .get_mut(old_thread.thread)?
2131 .context = old_context;
2132 }
2133 if cfg!(debug_assertions) {
2134 *self.vm_store_context_mut().component_context_mut() =
2135 [u32::MAX; NUM_COMPONENT_CONTEXT_SLOTS];
2136 }
2137 if let Some(thread) = thread.guest() {
2138 let thread = self.concurrent_state_mut()?.get_mut(thread.thread)?;
2139 let context = thread.context;
2140 if cfg!(debug_assertions) {
2141 thread.context = [u32::MAX; NUM_COMPONENT_CONTEXT_SLOTS];
2142 }
2143 *self.vm_store_context_mut().component_context_mut() = context;
2144 }
2145
2146 *self.vm_store_context_mut().current_thread_mut() = if thread.is_none() {
2148 VMLazyThread::none()
2149 } else {
2150 VMLazyThread::forced()
2151 };
2152
2153 Ok(old_thread)
2154 }
2155
2156 fn switch_or_trap_if_may_not_suspend(&mut self, instance: RuntimeInstance) -> Result<()> {
2158 if self.switch_if_may_not_suspend(instance)? {
2159 Ok(())
2160 } else {
2161 Err(Trap::CannotBlockSyncTask.into())
2162 }
2163 }
2164
2165 fn switch_if_may_not_suspend(&mut self, instance: RuntimeInstance) -> Result<bool> {
2169 self.concurrent_state_mut()?;
2173
2174 Ok(!self.concurrency_support()
2175 || !self
2176 .instance_state(instance)
2177 .concurrent_state()
2178 .do_not_suspend
2179 || self
2180 .concurrent_state_mut()?
2181 .promote_instance_local_thread_work_item(instance)?)
2182 }
2183
2184 fn enter_instance(&mut self, instance: RuntimeInstance) {
2188 log::trace!("enter {instance:?}");
2189 self.instance_state(instance)
2190 .concurrent_state()
2191 .do_not_enter = true;
2192 }
2193
2194 fn exit_instance(&mut self, instance: RuntimeInstance) -> Result<()> {
2198 log::trace!("exit {instance:?}");
2199 self.instance_state(instance)
2200 .concurrent_state()
2201 .do_not_enter = false;
2202 self.partition_pending(instance)
2203 }
2204
2205 fn partition_pending(&mut self, instance: RuntimeInstance) -> Result<()> {
2213 for (thread, kind) in
2214 mem::take(&mut self.instance_state(instance).concurrent_state().pending).into_iter()
2215 {
2216 let call = GuestCall { thread, kind };
2217 if call.is_ready(self)? {
2218 self.concurrent_state_mut()?
2219 .push_high_priority(WorkItem::GuestCall { instance, call });
2220 } else {
2221 self.instance_state(instance)
2222 .concurrent_state()
2223 .pending
2224 .insert(call.thread, call.kind);
2225 }
2226 }
2227
2228 if let Some(waker) = self
2229 .concurrent_state_mut()?
2230 .ready_for_concurrent_call_waker
2231 .take()
2232 {
2233 waker.wake();
2234 }
2235
2236 Ok(())
2237 }
2238
2239 pub(crate) fn backpressure_modify(
2241 &mut self,
2242 caller_instance: RuntimeInstance,
2243 modify: impl FnOnce(u16) -> Option<u16>,
2244 ) -> Result<()> {
2245 let state = self.instance_state(caller_instance).concurrent_state();
2246 let old = state.backpressure;
2247 let new = modify(old).ok_or_else(|| Trap::BackpressureOverflow)?;
2248 state.backpressure = new;
2249
2250 if old > 0 && new == 0 {
2251 self.partition_pending(caller_instance)?;
2254 }
2255
2256 Ok(())
2257 }
2258
2259 async fn resume_fiber(&mut self, fiber: StoreFiber<'static>) -> Result<()> {
2262 let old_thread = self.current_thread()?;
2263 log::trace!("resume_fiber: save current thread {old_thread:?}");
2264
2265 let fiber = fiber::resolve_or_release(self, fiber).await?;
2266
2267 self.set_thread(old_thread)?;
2268
2269 let state = self.concurrent_state_mut()?;
2270
2271 if let Some(ot) = old_thread.guest() {
2272 state.get_mut(ot.thread)?.state = GuestThreadState::Running;
2273 }
2274 log::trace!("resume_fiber: restore current thread {old_thread:?}");
2275
2276 if let Some(mut fiber) = fiber {
2277 log::trace!("resume_fiber: suspend reason {:?}", &state.suspend_reason);
2278 let reason = match state.suspend_reason.take() {
2280 Some(r) => r,
2281 None => bail_bug!("suspend reason missing when resuming fiber"),
2282 };
2283 match reason {
2284 SuspendReason::NeedWork => {
2285 if state.worker.is_none() {
2286 state.worker = Some(fiber);
2287 } else {
2288 fiber.dispose(self);
2289 }
2290 }
2291 SuspendReason::Yielding { thread } => {
2292 state.get_mut(thread.thread)?.state = GuestThreadState::Ready { fiber };
2293 let instance = state.get_mut(thread.task)?.instance;
2294 state.push_low_priority(WorkItem::ResumeThread { instance, thread });
2295 }
2296 SuspendReason::ExplicitlySuspending { thread } => {
2297 state.get_mut(thread.thread)?.state = GuestThreadState::Suspended(fiber);
2298 }
2299 SuspendReason::Waiting { set, thread } => {
2300 let old = state
2301 .get_mut(set)?
2302 .waiting
2303 .insert(thread, WaitMode::Fiber(fiber));
2304 assert!(old.is_none());
2305 }
2306 SuspendReason::YieldingToSubtask { thread } => {
2307 let item = WorkItem::ResumeFiber {
2316 instance: state.get_mut(thread.task)?.instance,
2317 thread,
2318 fiber,
2319 };
2320
2321 if state.next_switch_item.replace(item).is_some() {
2322 bail_bug!(
2325 "`ConcurrentState::next_switch_item` was already `Some(_)` when \
2326 a thread wanted to wait on a subtask"
2327 );
2328 }
2329 }
2330 };
2331 } else {
2332 log::trace!("resume_fiber: fiber has exited");
2333 }
2334
2335 Ok(())
2336 }
2337
2338 fn suspend(&mut self, reason: SuspendReason) -> Result<()> {
2344 log::trace!("suspend fiber: {reason:?}");
2345
2346 let state = self.concurrent_state_mut()?;
2347
2348 let (save_and_restore_thread, save_and_restore_next_switch_item) = match &reason {
2355 SuspendReason::Yielding { .. }
2356 | SuspendReason::Waiting { .. }
2357 | SuspendReason::ExplicitlySuspending { .. } => {
2358 if state.switch_item.is_none() {
2361 state.take_next_switch_item()?;
2362 }
2363
2364 (true, false)
2365 }
2366 SuspendReason::YieldingToSubtask { .. } => (true, true),
2367 SuspendReason::NeedWork => (false, false),
2368 };
2369
2370 let old_next_switch_item = if save_and_restore_next_switch_item {
2371 let item = state.next_switch_item.take();
2372 Some(state.push(item)?)
2376 } else {
2377 None
2378 };
2379
2380 let old_guest_thread = if save_and_restore_thread {
2381 self.current_thread()?
2382 } else {
2383 CurrentThread::None
2384 };
2385
2386 let waiting_set = match &reason {
2387 SuspendReason::Waiting { set, .. } => Some(*set),
2388 _ => None,
2389 };
2390
2391 let suspend_reason = &mut self.concurrent_state_mut()?.suspend_reason;
2392 assert!(suspend_reason.is_none());
2393 *suspend_reason = Some(reason);
2394
2395 if !self.fiber_async_state_mut().can_block() {
2398 return Err(format_err!("future dropped"));
2399 }
2400
2401 if let Some(set) = waiting_set {
2404 self.concurrent_state_mut()?.get_mut(set)?.num_waiting += 1;
2405 }
2406
2407 self.with_blocking(|_, cx| cx.suspend(StoreFiberYield::ReleaseStore))?;
2408
2409 if let Some(set) = waiting_set {
2410 self.concurrent_state_mut()?.get_mut(set)?.stop_waiting()?;
2411 }
2412
2413 if save_and_restore_thread {
2414 self.set_thread(old_guest_thread)?;
2415 }
2416
2417 if let Some(item) = old_next_switch_item {
2418 let state = self.concurrent_state_mut()?;
2419 state.next_switch_item = state.delete(item)?;
2420 }
2421
2422 Ok(())
2423 }
2424
2425 fn wait_for_event(
2426 &mut self,
2427 caller_instance: RuntimeInstance,
2428 waitable: Waitable,
2429 ) -> Result<()> {
2430 let caller = self.current_guest_thread()?;
2431 let state = self.concurrent_state_mut()?;
2432
2433 waitable.trap_if_in_waitable_set(state)?;
2434
2435 let set = state.get_mut(caller.thread)?.sync_call_set;
2436 waitable.join(state, Some(set))?;
2437
2438 self.switch_or_trap_if_may_not_suspend(caller_instance)?;
2439
2440 self.suspend(SuspendReason::Waiting {
2441 set,
2442 thread: caller,
2443 })?;
2444 let state = self.concurrent_state_mut()?;
2445
2446 waitable.join(state, None)
2447 }
2448
2449 fn cleanup_thread(
2471 &mut self,
2472 guest_thread: QualifiedThreadId,
2473 runtime_instance: RuntimeInstance,
2474 cleanup_task: CleanupTask,
2475 ) -> Result<()> {
2476 let state = self.concurrent_state_mut()?;
2477 state.take_next_switch_item()?;
2480 let thread_data = state.get_mut(guest_thread.thread)?;
2481 let sync_call_set = thread_data.sync_call_set;
2482 if let Some(guest_id) = thread_data.instance_rep {
2483 self.instance_state(runtime_instance)
2484 .thread_handle_table()
2485 .guest_thread_remove(guest_id)?;
2486 }
2487 let state = self.concurrent_state_mut()?;
2488
2489 for waitable in mem::take(&mut state.get_mut(sync_call_set)?.ready) {
2491 if let Some(Event::Subtask {
2492 status: Status::Returned | Status::ReturnCancelled,
2493 }) = waitable.common(self.concurrent_state_mut()?)?.event
2494 {
2495 waitable.delete_from(self)?;
2496 }
2497 }
2498
2499 let state = self.concurrent_state_mut()?;
2500 state.delete(guest_thread.thread)?;
2501 state.delete(sync_call_set)?;
2502 let task = state.get_mut(guest_thread.task)?;
2503 task.threads.remove(&guest_thread.thread);
2504
2505 if task.threads.is_empty() && !task.returned_or_cancelled() {
2506 bail!(Trap::NoAsyncResult);
2507 }
2508 let ready_to_delete = task.ready_to_delete();
2509
2510 if !task.decremented_interesting_task_count && task.exited && task.returned_or_cancelled() {
2511 task.decremented_interesting_task_count = true;
2512
2513 debug_assert!(state.interesting_tasks > 0);
2514 state.interesting_tasks -= 1;
2515 if state.interesting_tasks == 0
2516 && let Some(waker) = state.interesting_tasks_empty_waker.take()
2517 {
2518 waker.wake();
2519 }
2520 }
2521
2522 match cleanup_task {
2523 CleanupTask::Yes => {
2524 if ready_to_delete {
2525 Waitable::Guest(guest_thread.task).delete_from(self)?;
2526 }
2527 }
2528 CleanupTask::No => {}
2529 }
2530
2531 Ok(())
2532 }
2533
2534 fn cancel_guest_subtask_without_lowered_parameters(
2547 &mut self,
2548 caller_instance: RuntimeInstance,
2549 guest_task: TableId<GuestTask>,
2550 ) -> Result<()> {
2551 let concurrent_state = self.concurrent_state_mut()?;
2552 let task = concurrent_state.get_mut(guest_task)?;
2553 assert!(!task.already_lowered_parameters());
2554 task.lower_params = None;
2558 task.lift_result = None;
2559 task.exited = true;
2560 let instance = task.instance;
2561
2562 assert_eq!(1, task.threads.len());
2565 let thread = *task.threads.iter().next().unwrap();
2566 self.cleanup_thread(
2567 QualifiedThreadId {
2568 task: guest_task,
2569 thread,
2570 },
2571 caller_instance,
2572 CleanupTask::No,
2573 )?;
2574
2575 let pending = &mut self.instance_state(instance).concurrent_state().pending;
2577 let pending_count = pending.len();
2578 pending.retain(|thread, _| thread.task != guest_task);
2579 if pending.len() == pending_count {
2581 bail!(Trap::SubtaskCancelAfterTerminal);
2582 }
2583 Ok(())
2584 }
2585
2586 pub(crate) fn current_scope(&mut self) -> Result<Option<CurrentScope>> {
2589 if !self.concurrency_support() {
2590 return Ok(self
2591 .current_scope_id_not_concurrent()?
2592 .map(|id| CurrentScope::Id(Scope::Id(id))));
2593 }
2594
2595 Ok(match self.current_thread()? {
2596 CurrentThread::Guest(id) => Some(CurrentScope::Id(Scope::Id(id.task.rep()))),
2597 CurrentThread::Host(id) => Some(CurrentScope::Id(Scope::HostId(id.rep()))),
2598 CurrentThread::DeferredHost(_) => Some(CurrentScope::DeferredHost),
2599 CurrentThread::None => return Ok(None),
2600 })
2601 }
2602
2603 pub(crate) fn queue_task(
2604 &mut self,
2605 task: impl FnOnce(&mut dyn VMStore) -> Result<()> + Send + 'static,
2606 ) -> Result<()> {
2607 self.concurrent_state_mut()?
2608 .push_high_priority(WorkItem::WorkerFunction(AlwaysMut::new(Box::new(task))));
2609 Ok(())
2610 }
2611
2612 fn any_may_not_suspend(&mut self) -> Result<bool> {
2621 Ok(self
2629 .concurrent_state_mut()?
2630 .table
2631 .get_mut()
2632 .iter_mut()
2633 .filter_map(|(_, entry)| {
2634 if let Some(task) = entry.downcast_ref::<GuestTask>() {
2635 Some(task.instance)
2636 } else {
2637 None
2638 }
2639 })
2640 .collect::<Vec<_>>()
2641 .into_iter()
2642 .any(|instance| {
2643 self.instance_state(instance)
2644 .concurrent_state()
2645 .do_not_suspend
2646 }))
2647 }
2648}
2649
2650enum CleanupTask {
2651 Yes,
2652 No,
2653}
2654
2655impl Instance {
2656 fn get_event(
2659 self,
2660 store: &mut StoreOpaque,
2661 guest_task: TableId<GuestTask>,
2662 set: Option<TableId<WaitableSet>>,
2663 cancellable: bool,
2664 ) -> Result<Option<(Event, Option<(Waitable, u32)>)>> {
2665 let state = store.concurrent_state_mut()?;
2666
2667 let task = state.get_mut(guest_task)?;
2668 let event = &mut task.event;
2669 if let Some(ev) = event
2670 && (cancellable || !matches!(ev, Event::Cancelled))
2671 {
2672 log::trace!("deliver event {ev:?} to {guest_task:?}");
2673
2674 if matches!(ev, Event::Cancelled) {
2675 task.cancel_request_delivered = true;
2676 }
2677
2678 let ev = *ev;
2679 *event = None;
2680 return Ok(Some((ev, None)));
2681 }
2682
2683 let set = match set {
2684 Some(set) => set,
2685 None => return Ok(None),
2686 };
2687 let waitable = match state.get_mut(set)?.ready.pop_first() {
2688 Some(v) => v,
2689 None => return Ok(None),
2690 };
2691
2692 let common = waitable.common(state)?;
2693 let handle = match common.handle {
2694 Some(h) => h,
2695 None => bail_bug!("handle not set when delivering event"),
2696 };
2697 let event = match common.event.take() {
2698 Some(e) => e,
2699 None => bail_bug!("event not set when delivering event"),
2700 };
2701
2702 log::trace!(
2703 "deliver event {event:?} to {guest_task:?} for {waitable:?} (handle {handle}); set {set:?}"
2704 );
2705
2706 waitable.on_delivery(store, self, event)?;
2707
2708 Ok(Some((event, Some((waitable, handle)))))
2709 }
2710
2711 fn handle_callback_code(
2717 self,
2718 store: &mut StoreOpaque,
2719 guest_thread: QualifiedThreadId,
2720 runtime_instance: RuntimeComponentInstanceIndex,
2721 code: u32,
2722 ) -> Result<()> {
2723 let (code, set) = unpack_callback_code(code);
2724
2725 log::trace!("received callback code from {guest_thread:?}: {code} (set: {set})");
2726
2727 let state = store.concurrent_state_mut()?;
2728
2729 state.take_next_switch_item()?;
2730
2731 let get_set = |store: &mut StoreOpaque, handle| -> Result<_> {
2732 let set = store
2733 .instance_state(self.runtime_instance(runtime_instance))
2734 .handle_table()
2735 .waitable_set_rep(handle)?;
2736
2737 Ok(TableId::<WaitableSet>::new(set))
2738 };
2739
2740 match code {
2741 callback_code::EXIT => {
2742 log::trace!("implicit thread {guest_thread:?} completed");
2743 let task = store.concurrent_state_mut()?.get_mut(guest_thread.task)?;
2744 task.exited = true;
2745 task.callback = None;
2746
2747 let runtime_instance = self.runtime_instance(runtime_instance);
2748
2749 store.switch_or_trap_if_may_not_suspend(runtime_instance)?;
2754
2755 store.cleanup_thread(guest_thread, runtime_instance, CleanupTask::Yes)?;
2756 }
2757 callback_code::YIELD => {
2758 let old = state
2761 .get_mut(guest_thread.thread)?
2762 .wake_on_cancel
2763 .replace(WakeOnCancel::Yielding);
2764 if !old.is_none() {
2765 bail_bug!("thread unexpectedly had wake_on_cancel set");
2766 }
2767
2768 let call = GuestCall {
2775 thread: guest_thread,
2776 kind: GuestCallKind::DeliverEvent {
2777 instance: self,
2778 set: None,
2779 },
2780 };
2781 state.push_low_priority(WorkItem::GuestCall {
2784 instance: self.runtime_instance(runtime_instance),
2785 call,
2786 });
2787 }
2788 callback_code::WAIT => {
2789 let set = get_set(store, set)?;
2790 self.wait_with_callback(store, guest_thread, set)?;
2791 }
2792 _ => bail!(Trap::UnsupportedCallbackCode),
2793 }
2794
2795 Ok(())
2796 }
2797
2798 fn wait_with_callback(
2804 self,
2805 store: &mut StoreOpaque,
2806 guest_thread: QualifiedThreadId,
2807 set: TableId<WaitableSet>,
2808 ) -> Result<()> {
2809 let state = store.concurrent_state_mut()?;
2810 state.get_mut(set)?.num_waiting += 1;
2813
2814 if state.get_mut(guest_thread.task)?.event.is_some()
2815 || !state.get_mut(set)?.ready.is_empty()
2816 {
2817 let instance = state.get_mut(guest_thread.task)?.instance;
2819 state.push_high_priority(WorkItem::GuestCall {
2820 instance,
2821 call: GuestCall {
2822 thread: guest_thread,
2823 kind: GuestCallKind::DeliverEvent {
2824 instance: self,
2825 set: Some(set),
2826 },
2827 },
2828 });
2829 return Ok(());
2830 }
2831
2832 let instance = store
2838 .concurrent_state_mut()?
2839 .get_mut(guest_thread.task)?
2840 .instance;
2841 store.switch_or_trap_if_may_not_suspend(instance)?;
2842
2843 let state = store.concurrent_state_mut()?;
2849 let old = state
2850 .get_mut(guest_thread.thread)?
2851 .wake_on_cancel
2852 .replace(WakeOnCancel::Waiting(set));
2853 if !old.is_none() {
2854 bail_bug!("thread unexpectedly had wake_on_cancel set");
2855 }
2856 let old = state
2857 .get_mut(set)?
2858 .waiting
2859 .insert(guest_thread, WaitMode::Callback(self));
2860 if !old.is_none() {
2861 bail_bug!("set's waiting set already had this thread registered");
2862 }
2863 Ok(())
2864 }
2865
2866 unsafe fn stage_call<T: 'static>(
2873 self,
2874 mut store: StoreContextMut<T>,
2875 guest_thread: QualifiedThreadId,
2876 callee: SendSyncPtr<VMFuncRef>,
2877 param_count: usize,
2878 result_count: usize,
2879 async_: bool,
2880 callback: Option<SendSyncPtr<VMFuncRef>>,
2881 post_return: Option<SendSyncPtr<VMFuncRef>>,
2882 host_caller: bool,
2883 ) -> Result<()> {
2884 unsafe fn make_call<T: 'static>(
2899 store: StoreContextMut<T>,
2900 guest_thread: QualifiedThreadId,
2901 callee: SendSyncPtr<VMFuncRef>,
2902 param_count: usize,
2903 result_count: usize,
2904 ) -> impl FnOnce(&mut dyn VMStore) -> Result<[MaybeUninit<ValRaw>; MAX_FLAT_PARAMS]>
2905 + Send
2906 + Sync
2907 + 'static
2908 + use<T> {
2909 let token = StoreToken::new(store);
2910 move |store: &mut dyn VMStore| {
2911 let mut storage = [MaybeUninit::uninit(); MAX_FLAT_PARAMS];
2912
2913 store
2914 .concurrent_state_mut()?
2915 .get_mut(guest_thread.thread)?
2916 .state = GuestThreadState::Running;
2917 let task = store.concurrent_state_mut()?.get_mut(guest_thread.task)?;
2918 let lower = match task.lower_params.take() {
2919 Some(l) => l,
2920 None => bail_bug!("lower_params missing"),
2921 };
2922
2923 lower(store, &mut storage[..param_count])?;
2924
2925 let mut store = token.as_context_mut(store);
2926
2927 unsafe {
2930 crate::Func::call_unchecked_raw(
2931 &mut store,
2932 callee.as_non_null(),
2933 NonNull::new(
2934 &mut storage[..param_count.max(result_count)]
2935 as *mut [MaybeUninit<ValRaw>] as _,
2936 )
2937 .unwrap(),
2938 UncaughtException::Trap,
2939 )?;
2940 }
2941
2942 Ok(storage)
2943 }
2944 }
2945
2946 let call = unsafe {
2950 make_call(
2951 store.as_context_mut(),
2952 guest_thread,
2953 callee,
2954 param_count,
2955 result_count,
2956 )
2957 };
2958
2959 let callee_instance = store
2960 .0
2961 .concurrent_state_mut()?
2962 .get_mut(guest_thread.task)?
2963 .instance;
2964
2965 let fun = if callback.is_some() {
2966 assert!(async_);
2967
2968 Box::new(move |store: &mut dyn VMStore| {
2969 self.add_guest_thread_to_instance_table(
2970 guest_thread.thread,
2971 store,
2972 callee_instance.index,
2973 )?;
2974 let old_thread = store.set_thread(guest_thread)?;
2975 log::trace!(
2976 "stackless call: replaced {old_thread:?} with {guest_thread:?} as current thread"
2977 );
2978
2979 store.enter_instance(callee_instance);
2980
2981 let storage = call(store)?;
2988
2989 store.exit_instance(callee_instance)?;
2990
2991 store.set_thread(old_thread)?;
2992 let state = store.concurrent_state_mut()?;
2993 if let Some(t) = old_thread.guest() {
2994 state.get_mut(t.thread)?.state = GuestThreadState::Running;
2995 }
2996 log::trace!("stackless call: restored {old_thread:?} as current thread");
2997
2998 let code = unsafe { storage[0].assume_init() }.get_i32() as u32;
3001
3002 self.handle_callback_code(store, guest_thread, callee_instance.index, code)
3003 }) as Box<dyn FnOnce(&mut dyn VMStore) -> Result<()> + Send + Sync>
3004 } else {
3005 let token = StoreToken::new(store.as_context_mut());
3006 Box::new(move |store: &mut dyn VMStore| {
3007 self.add_guest_thread_to_instance_table(
3008 guest_thread.thread,
3009 store,
3010 callee_instance.index,
3011 )?;
3012 let old_thread = store.set_thread(guest_thread)?;
3013 log::trace!(
3014 "sync/async-stackful call: replaced {old_thread:?} with {guest_thread:?} as current thread",
3015 );
3016 let flags = self.id().get(store).instance_flags(callee_instance.index);
3017
3018 let callee_async_typed = store
3019 .concurrent_state_mut()?
3020 .get_mut(guest_thread.task)?
3021 .async_typed;
3022
3023 if !async_ && callee_async_typed {
3027 store.enter_instance(callee_instance);
3028 }
3029
3030 if !callee_async_typed {
3031 store.enter_sync_call(callee_instance)?;
3032 }
3033
3034 let storage = call(store)?;
3041
3042 if !callee_async_typed {
3043 store.exit_sync_call(callee_instance)?;
3044 }
3045
3046 if !async_ {
3047 if callee_async_typed {
3053 store.exit_instance(callee_instance)?;
3054 }
3055
3056 let lift = {
3057 let state = store.concurrent_state_mut()?;
3058 if !state.get_mut(guest_thread.task)?.result.is_none() {
3059 bail_bug!("task has already produced a result");
3060 }
3061
3062 match state.get_mut(guest_thread.task)?.lift_result.take() {
3063 Some(lift) => lift,
3064 None => bail_bug!("lift_result field is missing"),
3065 }
3066 };
3067
3068 let result = (lift.lift)(store, unsafe {
3071 mem::transmute::<&[MaybeUninit<ValRaw>], &[ValRaw]>(
3072 &storage[..result_count],
3073 )
3074 })?;
3075
3076 let post_return_arg = match result_count {
3077 0 => ValRaw::i32(0),
3078 1 => unsafe { storage[0].assume_init() },
3081 _ => unreachable!(),
3082 };
3083
3084 unsafe {
3085 call_post_return(
3086 token.as_context_mut(store),
3087 post_return.map(|v| v.as_non_null()),
3088 post_return_arg,
3089 flags,
3090 )?;
3091 }
3092
3093 self.task_complete(store, guest_thread.task, result, Status::Returned)?;
3094 }
3095
3096 store.set_thread(old_thread)?;
3097
3098 store
3099 .concurrent_state_mut()?
3100 .get_mut(guest_thread.task)?
3101 .exited = true;
3102
3103 log::trace!(
3104 "clean up thread; async lifted? {async_} async typed? {callee_async_typed}"
3105 );
3106
3107 if callee_async_typed {
3108 store.switch_or_trap_if_may_not_suspend(callee_instance)?;
3113 }
3114
3115 store.cleanup_thread(guest_thread, callee_instance, CleanupTask::Yes)?;
3117 Ok(())
3118 })
3119 };
3120
3121 store.0.concurrent_state_mut()?.push_work_item(
3122 WorkItem::GuestCall {
3123 instance: callee_instance,
3124 call: GuestCall {
3125 thread: guest_thread,
3126 kind: GuestCallKind::StartImplicit(fun),
3127 },
3128 },
3129 if host_caller {
3130 Priority::High
3131 } else {
3132 Priority::Switch
3133 },
3134 )?;
3135
3136 Ok(())
3137 }
3138
3139 unsafe fn prepare_call<T: 'static>(
3152 self,
3153 mut store: StoreContextMut<T>,
3154 start: NonNull<VMFuncRef>,
3155 return_: NonNull<VMFuncRef>,
3156 caller_instance: RuntimeComponentInstanceIndex,
3157 callee_instance: RuntimeComponentInstanceIndex,
3158 task_return_type: TypeTupleIndex,
3159 callee_async_typed: bool,
3160 memory: *mut VMMemoryDefinition,
3161 string_encoding: StringEncoding,
3162 caller_info: CallerInfo,
3163 ) -> Result<()> {
3164 enum ResultInfo {
3165 Heap { results: u32 },
3166 Stack { result_count: u32 },
3167 }
3168
3169 let result_info = match &caller_info {
3170 CallerInfo::Async {
3171 has_result: true,
3172 params,
3173 } => ResultInfo::Heap {
3174 results: match params.last() {
3175 Some(r) => r.get_u32(),
3176 None => bail_bug!("retptr missing"),
3177 },
3178 },
3179 CallerInfo::Async {
3180 has_result: false, ..
3181 } => ResultInfo::Stack { result_count: 0 },
3182 CallerInfo::Sync {
3183 result_count,
3184 params,
3185 } if *result_count > u32::try_from(MAX_FLAT_RESULTS)? => ResultInfo::Heap {
3186 results: match params.last() {
3187 Some(r) => r.get_u32(),
3188 None => bail_bug!("arg ptr missing"),
3189 },
3190 },
3191 CallerInfo::Sync { result_count, .. } => ResultInfo::Stack {
3192 result_count: *result_count,
3193 },
3194 };
3195
3196 let sync_caller = matches!(caller_info, CallerInfo::Sync { .. });
3197
3198 let start = SendSyncPtr::new(start);
3202 let return_ = SendSyncPtr::new(return_);
3203 let token = StoreToken::new(store.as_context_mut());
3204 let old_thread = store.0.current_guest_thread()?;
3205
3206 let state = store.0.concurrent_state_mut()?;
3207
3208 debug_assert_eq!(
3209 state.get_mut(old_thread.task)?.instance,
3210 self.runtime_instance(caller_instance)
3211 );
3212
3213 let guest_thread = GuestTask::new(
3214 state,
3215 Box::new(move |store, dst| {
3216 let mut store = token.as_context_mut(store);
3217 assert!(dst.len() <= MAX_FLAT_PARAMS);
3218 let mut src = [MaybeUninit::uninit(); MAX_FLAT_PARAMS + 1];
3220 let count = match caller_info {
3221 CallerInfo::Async { params, has_result } => {
3225 let params = ¶ms[..params.len() - usize::from(has_result)];
3226 for (param, src) in params.iter().zip(&mut src) {
3227 src.write(*param);
3228 }
3229 params.len()
3230 }
3231
3232 CallerInfo::Sync { params, .. } => {
3234 for (param, src) in params.iter().zip(&mut src) {
3235 src.write(*param);
3236 }
3237 params.len()
3238 }
3239 };
3240 unsafe {
3247 crate::Func::call_unchecked_raw(
3248 &mut store,
3249 start.as_non_null(),
3250 NonNull::new(
3251 &mut src[..count.max(dst.len())] as *mut [MaybeUninit<ValRaw>] as _,
3252 )
3253 .unwrap(),
3254 UncaughtException::Trap,
3255 )?;
3256 }
3257 dst.copy_from_slice(&src[..dst.len()]);
3258 let task = store.0.current_guest_thread()?.task;
3259 let state = store.0.concurrent_state_mut()?;
3260 Waitable::Guest(task).set_event(
3261 state,
3262 Some(Event::Subtask {
3263 status: Status::Started,
3264 }),
3265 )?;
3266 Ok(())
3267 }),
3268 LiftResult {
3269 lift: Box::new(move |store, src| {
3270 let mut store = token.as_context_mut(store);
3273 let mut my_src = src.to_owned(); if let ResultInfo::Heap { results } = &result_info {
3275 my_src.push(ValRaw::u32(*results));
3276 }
3277
3278 unsafe {
3285 crate::Func::call_unchecked_raw(
3286 &mut store,
3287 return_.as_non_null(),
3288 my_src.as_mut_slice().into(),
3289 UncaughtException::Trap,
3290 )?;
3291 }
3292
3293 let thread = store.0.current_guest_thread()?;
3294 let state = store.0.concurrent_state_mut()?;
3295 if sync_caller {
3296 state.get_mut(thread.task)?.sync_result = SyncResult::Produced(
3297 if let ResultInfo::Stack { result_count } = &result_info {
3298 match result_count {
3299 0 => None,
3300 1 => Some(my_src[0]),
3301 _ => unreachable!(),
3302 }
3303 } else {
3304 None
3305 },
3306 );
3307 }
3308 Ok(Box::new(DummyResult) as Box<dyn Any + Send + Sync>)
3309 }),
3310 ty: task_return_type,
3311 memory: NonNull::new(memory).map(SendSyncPtr::new),
3312 string_encoding,
3313 },
3314 Caller::Guest { thread: old_thread },
3315 None,
3316 self.runtime_instance(callee_instance),
3317 callee_async_typed,
3318 false,
3321 )?;
3322
3323 store.0.set_thread(guest_thread)?;
3326 log::trace!("pushed {guest_thread:?} as current thread; old thread was {old_thread:?}");
3327
3328 Ok(())
3329 }
3330
3331 unsafe fn call_callback<T>(
3336 self,
3337 mut store: StoreContextMut<T>,
3338 function: SendSyncPtr<VMFuncRef>,
3339 event: Event,
3340 handle: u32,
3341 ) -> Result<u32> {
3342 let (ordinal, result) = event.parts();
3343 let params = &mut [
3344 ValRaw::u32(ordinal),
3345 ValRaw::u32(handle),
3346 ValRaw::u32(result),
3347 ];
3348 unsafe {
3353 crate::Func::call_unchecked_raw(
3354 &mut store,
3355 function.as_non_null(),
3356 params.as_mut_slice().into(),
3357 UncaughtException::Trap,
3358 )?;
3359 }
3360 Ok(params[0].get_u32())
3361 }
3362
3363 unsafe fn start_call<T: 'static>(
3381 self,
3382 mut store: StoreContextMut<T>,
3383 callback: *mut VMFuncRef,
3384 post_return: *mut VMFuncRef,
3385 callee: NonNull<VMFuncRef>,
3386 param_count: u32,
3387 result_count: u32,
3388 flags: u32,
3389 storage: &mut [MaybeUninit<ValRaw>],
3390 ) -> Result<()> {
3391 let token = StoreToken::new(store.as_context_mut());
3392 let async_caller = (flags & START_FLAG_ASYNC_CALLER) != 0;
3393 let guest_thread = store.0.current_guest_thread()?;
3394 let state = store.0.concurrent_state_mut()?;
3395
3396 if !state.event_loop_running {
3397 bail_bug!("Instance::start_call called without a running event loop");
3398 }
3399
3400 let callee = SendSyncPtr::new(callee);
3401 let param_count = usize::try_from(param_count)?;
3402 assert!(param_count <= MAX_FLAT_PARAMS);
3403 let result_count = usize::try_from(result_count)?;
3404 assert!(result_count <= MAX_FLAT_RESULTS);
3405
3406 let task = state.get_mut(guest_thread.task)?;
3407 let callee_async_typed = task.async_typed;
3408 let callee_instance = task.instance;
3409
3410 task.async_lifted = (flags & START_FLAG_ASYNC_CALLEE) != 0;
3411
3412 if let Some(callback) = NonNull::new(callback) {
3413 let callback = SendSyncPtr::new(callback);
3417 task.callback = Some(Box::new(move |store, event, handle| {
3418 let store = token.as_context_mut(store);
3419 unsafe { self.call_callback::<T>(store, callback, event, handle) }
3420 }));
3421 }
3422
3423 let Caller::Guest { thread: caller } = &task.caller else {
3424 bail_bug!("start_call unexpectedly invoked for host->guest call");
3427 };
3428 let caller = *caller;
3429 let caller_instance = state.get_mut(caller.task)?.instance;
3430
3431 unsafe {
3433 self.stage_call(
3434 store.as_context_mut(),
3435 guest_thread,
3436 callee,
3437 param_count,
3438 result_count,
3439 (flags & START_FLAG_ASYNC_CALLEE) != 0,
3440 NonNull::new(callback).map(SendSyncPtr::new),
3441 NonNull::new(post_return).map(SendSyncPtr::new),
3442 false,
3443 )?;
3444 }
3445
3446 let old_do_not_suspend = if callee_async_typed {
3447 let state = store.0.instance_state(callee_instance).concurrent_state();
3454 let old_do_not_suspend = state.do_not_suspend;
3455 state.do_not_suspend = false;
3456 Some(old_do_not_suspend)
3457 } else {
3458 None
3459 };
3460
3461 let state = store.0.concurrent_state_mut()?;
3462
3463 let guest_waitable = Waitable::Guest(guest_thread.task);
3466 let old_set = guest_waitable.common(state)?.set;
3467 let set = state.get_mut(caller.thread)?.sync_call_set;
3468 guest_waitable.join(state, Some(set))?;
3469
3470 store.0.set_thread(CurrentThread::None)?;
3471
3472 let mut yielded = false;
3488 let (status, waitable) = loop {
3489 store.0.suspend(if yielded {
3490 SuspendReason::Waiting {
3491 set,
3492 thread: caller,
3493 }
3494 } else {
3495 yielded = true;
3496 SuspendReason::YieldingToSubtask { thread: caller }
3497 })?;
3498
3499 if let Some(old_do_not_suspend) = old_do_not_suspend {
3500 store
3501 .0
3502 .instance_state(callee_instance)
3503 .concurrent_state()
3504 .do_not_suspend = old_do_not_suspend;
3505 }
3506
3507 let state = store.0.concurrent_state_mut()?;
3508
3509 log::trace!("taking event for {:?}", guest_thread.task);
3510 let event = guest_waitable.take_event(state)?;
3511 let Some(Event::Subtask { status }) = event else {
3512 bail_bug!("subtasks should only get subtask events, got {event:?}")
3513 };
3514
3515 log::trace!("status {status:?} for {:?}", guest_thread.task);
3516
3517 if status == Status::Returned {
3518 break (status, None);
3520 } else if async_caller {
3521 let handle = store
3525 .0
3526 .instance_state(caller_instance)
3527 .handle_table()
3528 .subtask_insert_guest(guest_thread.task.rep())?;
3529 store
3530 .0
3531 .concurrent_state_mut()?
3532 .get_mut(guest_thread.task)?
3533 .common
3534 .handle = Some(handle);
3535 break (status, Some(handle));
3536 } else {
3537 store.0.switch_or_trap_if_may_not_suspend(caller_instance)?;
3541 }
3542 };
3543
3544 guest_waitable.join(store.0.concurrent_state_mut()?, old_set)?;
3545
3546 store.0.set_thread(caller)?;
3548 store
3549 .0
3550 .concurrent_state_mut()?
3551 .get_mut(caller.thread)?
3552 .state = GuestThreadState::Running;
3553 log::trace!("popped current thread {guest_thread:?}; new thread is {caller:?}");
3554
3555 if async_caller {
3556 let Some(slot) = storage.first_mut() else {
3557 bail_bug!("no storage for async call status");
3558 };
3559 *slot = MaybeUninit::new(ValRaw::u32(status.pack(waitable)));
3560 } else {
3561 let state = store.0.concurrent_state_mut()?;
3564 let task = state.get_mut(guest_thread.task)?;
3565 if let Some(result) = task.sync_result.take()? {
3566 if let Some(result) = result {
3567 storage[0] = MaybeUninit::new(result);
3568 }
3569
3570 if task.exited && task.ready_to_delete() {
3571 Waitable::Guest(guest_thread.task).delete_from(store.0)?;
3572 }
3573 }
3574 }
3575
3576 Ok(())
3577 }
3578
3579 pub(crate) fn first_poll<T: 'static, R: Send + 'static>(
3595 self,
3596 mut store: StoreContextMut<'_, T>,
3597 host_task: EnteredHostTask,
3598 result_may_require_realloc: bool,
3599 future: impl Future<Output = Result<R>> + Send + 'static,
3600 lower: impl FnOnce(StoreContextMut<T>, Option<R>, bool, Option<TableId<HostTask>>) -> Result<()>
3601 + Send
3602 + 'static,
3603 ) -> Result<u32> {
3604 let token = StoreToken::new(store.as_context_mut());
3605
3606 let (join_handle, future) = JoinHandle::run(future);
3609 let mut future = Box::pin(future);
3610
3611 let poll = tls::set(store.0, || {
3616 future
3617 .as_mut()
3618 .poll(&mut Context::from_waker(&Waker::noop()))
3619 });
3620
3621 match poll {
3622 Poll::Ready(result) => {
3624 let result = result.transpose()?;
3625 let task = store.0.current_materialized_host_task()?;
3628 lower(store.as_context_mut(), result, true, task)?;
3629 return Ok(Status::Returned.pack(None));
3630 }
3631
3632 Poll::Pending => {}
3634 }
3635
3636 let Some(task) = store.0.materialize_host_task_id()? else {
3640 bail_bug!("current thread is not a host thread")
3641 };
3642 {
3643 let state = &mut store.0.concurrent_state_mut()?.get_mut(task)?.state;
3644 assert!(matches!(state, HostTaskState::CalleeStarted));
3645 *state = HostTaskState::CalleeRunning(join_handle);
3646 }
3647
3648 let future = Box::pin(async move {
3656 let result = match run_with_host_task_set(task, future).await? {
3657 Some(result) => Some(result?),
3658 None => None,
3659 };
3660 let on_complete = move |store: &mut dyn VMStore| {
3661 let mut store = token.as_context_mut(store);
3665 let old = store.0.set_thread(task)?;
3666
3667 let status = if result.is_some() {
3668 Status::Returned
3669 } else {
3670 Status::ReturnCancelled
3671 };
3672
3673 lower(store.as_context_mut(), result, false, Some(task))?;
3674 let state = store.0.concurrent_state_mut()?;
3675 match &mut state.get_mut(task)?.state {
3676 pending @ HostTaskState::CalleeCancelling => {
3679 *pending = HostTaskState::CalleeDone { cancelled: true };
3680 }
3681
3682 other => *other = HostTaskState::CalleeDone { cancelled: false },
3684 }
3685 Waitable::Host(task).set_event(state, Some(Event::Subtask { status }))?;
3686
3687 store.0.set_thread(old)?;
3688 Ok(())
3689 };
3690
3691 tls::get(move |store| {
3692 if result_may_require_realloc {
3693 store
3698 .concurrent_state_mut()?
3699 .push_high_priority(WorkItem::WorkerFunction(AlwaysMut::new(Box::new(
3700 on_complete,
3701 ))));
3702 Ok(())
3703 } else {
3704 on_complete(store)
3707 }
3708 })
3709 });
3710
3711 let caller = match host_task {
3714 Some(caller) => caller,
3715 None => bail_bug!("host task wasn't created but should have been"),
3716 };
3717 let state = store.0.concurrent_state_mut()?;
3718 state.push_future(future);
3719 let instance = state.get_mut(caller.task)?.instance;
3720 let handle = store
3721 .0
3722 .instance_state(instance)
3723 .handle_table()
3724 .subtask_insert_host(task.rep())?;
3725 store.0.concurrent_state_mut()?.get_mut(task)?.common.handle = Some(handle);
3726 log::trace!("assign {task:?} handle {handle} for {caller:?} instance {instance:?}");
3727
3728 store.0.set_thread(caller)?;
3732 Ok(Status::Started.pack(Some(handle)))
3733 }
3734
3735 pub(crate) fn task_return(
3738 self,
3739 store: &mut dyn VMStore,
3740 ty: TypeTupleIndex,
3741 options: OptionsIndex,
3742 storage: &[ValRaw],
3743 ) -> Result<()> {
3744 let guest_thread = store.current_guest_thread()?;
3745 let state = store.concurrent_state_mut()?;
3746 if !state.get_mut(guest_thread.task)?.async_lifted {
3747 bail!(Trap::TaskReturnOrCancelSyncLifted);
3748 }
3749 let lift = state
3750 .get_mut(guest_thread.task)?
3751 .lift_result
3752 .take()
3753 .ok_or_else(|| Trap::TaskCancelOrReturnTwice)?;
3754 if !state.get_mut(guest_thread.task)?.result.is_none() {
3755 bail_bug!("task result unexpectedly already set");
3756 }
3757
3758 let CanonicalOptions {
3759 string_encoding,
3760 data_model,
3761 ..
3762 } = &self.id().get(store).component().env_component().options[options];
3763
3764 let invalid = ty != lift.ty
3765 || string_encoding != &lift.string_encoding
3766 || match data_model {
3767 CanonicalOptionsDataModel::LinearMemory(opts) => match opts.memory {
3768 Some(memory) => {
3769 let expected = lift.memory.map(|v| v.as_ptr()).unwrap_or(ptr::null_mut());
3770 let actual = self.id().get(store).runtime_memory(memory);
3771 expected != actual.as_ptr()
3772 }
3773 None => false,
3776 },
3777 CanonicalOptionsDataModel::Gc { .. } => true,
3779 };
3780
3781 if invalid {
3782 bail!(Trap::TaskReturnInvalid);
3783 }
3784
3785 log::trace!("task.return for {guest_thread:?}");
3786
3787 let result = (lift.lift)(store, storage)?;
3788 self.task_complete(store, guest_thread.task, result, Status::Returned)
3789 }
3790
3791 pub(crate) fn task_cancel(self, store: &mut StoreOpaque) -> Result<()> {
3793 let guest_thread = store.current_guest_thread()?;
3794 let state = store.concurrent_state_mut()?;
3795 let task = state.get_mut(guest_thread.task)?;
3796 if !task.async_lifted {
3797 bail!(Trap::TaskReturnOrCancelSyncLifted);
3798 }
3799 if !task.cancel_request_delivered {
3800 bail!(Trap::TaskCancelNotCancelled);
3801 }
3802 _ = task
3803 .lift_result
3804 .take()
3805 .ok_or_else(|| Trap::TaskCancelOrReturnTwice)?;
3806
3807 if !task.result.is_none() {
3808 bail_bug!("task result should not bet set yet");
3809 }
3810
3811 log::trace!("task.cancel for {guest_thread:?}");
3812
3813 self.task_complete(
3814 store,
3815 guest_thread.task,
3816 Box::new(DummyResult),
3817 Status::ReturnCancelled,
3818 )
3819 }
3820
3821 fn task_complete(
3827 self,
3828 store: &mut StoreOpaque,
3829 guest_task: TableId<GuestTask>,
3830 result: Box<dyn Any + Send + Sync>,
3831 status: Status,
3832 ) -> Result<()> {
3833 store
3834 .component_resource_tables(Some(self))?
3835 .validate_scope_exit()?;
3836
3837 let state = store.concurrent_state_mut()?;
3838 let task = state.get_mut(guest_task)?;
3839
3840 task.event = None;
3844
3845 if let Caller::Host { tx, .. } = &mut task.caller {
3846 if let Some(tx) = tx.take() {
3847 _ = tx.send(result);
3848 }
3849 } else {
3850 task.result = Some(result);
3851 Waitable::Guest(guest_task).set_event(state, Some(Event::Subtask { status }))?;
3852 }
3853
3854 Ok(())
3855 }
3856
3857 pub(crate) fn waitable_set_new(
3859 self,
3860 store: &mut StoreOpaque,
3861 caller_instance: RuntimeComponentInstanceIndex,
3862 ) -> Result<u32> {
3863 let set = store.concurrent_state_mut()?.push(WaitableSet::default())?;
3864 let handle = store
3865 .instance_state(self.runtime_instance(caller_instance))
3866 .handle_table()
3867 .waitable_set_insert(set.rep())?;
3868 log::trace!("new waitable set {set:?} (handle {handle})");
3869 Ok(handle)
3870 }
3871
3872 pub(crate) fn waitable_set_drop(
3874 self,
3875 store: &mut StoreOpaque,
3876 caller_instance: RuntimeComponentInstanceIndex,
3877 set: u32,
3878 ) -> Result<()> {
3879 let rep = store
3880 .instance_state(self.runtime_instance(caller_instance))
3881 .handle_table()
3882 .waitable_set_remove(set)?;
3883
3884 log::trace!("drop waitable set {rep} (handle {set})");
3885
3886 let set = store
3890 .concurrent_state_mut()?
3891 .get_mut(TableId::<WaitableSet>::new(rep))?;
3892 if set.num_waiting > 0 {
3893 bail!(Trap::WaitableSetDropHasWaiters);
3894 }
3895
3896 store
3897 .concurrent_state_mut()?
3898 .delete(TableId::<WaitableSet>::new(rep))?;
3899
3900 Ok(())
3901 }
3902
3903 pub(crate) fn waitable_join(
3905 self,
3906 store: &mut StoreOpaque,
3907 caller_instance: RuntimeComponentInstanceIndex,
3908 waitable_handle: u32,
3909 set_handle: u32,
3910 ) -> Result<()> {
3911 let mut instance = self.id().get_mut(store);
3912 let waitable =
3913 Waitable::from_instance(instance.as_mut(), caller_instance, waitable_handle)?;
3914
3915 let set = if set_handle == 0 {
3916 None
3917 } else {
3918 let set = instance.instance_states().0[caller_instance]
3919 .handle_table()
3920 .waitable_set_rep(set_handle)?;
3921
3922 let state = store.concurrent_state_mut()?;
3923 if let Some(old) = waitable.common(state)?.set
3924 && state.get_mut(old)?.is_sync_call_set
3925 {
3926 bail!(Trap::WaitableSyncAndAsync);
3927 }
3928
3929 Some(TableId::<WaitableSet>::new(set))
3930 };
3931
3932 log::trace!(
3933 "waitable {waitable:?} (handle {waitable_handle}) join set {set:?} (handle {set_handle})",
3934 );
3935
3936 waitable.join(store.concurrent_state_mut()?, set)
3937 }
3938
3939 pub(crate) fn subtask_drop(
3941 self,
3942 store: &mut StoreOpaque,
3943 caller_instance: RuntimeComponentInstanceIndex,
3944 task_id: u32,
3945 ) -> Result<()> {
3946 self.waitable_join(store, caller_instance, task_id, 0)?;
3947
3948 let (rep, is_host) = store
3949 .instance_state(self.runtime_instance(caller_instance))
3950 .handle_table()
3951 .subtask_remove(task_id)?;
3952
3953 let concurrent_state = store.concurrent_state_mut()?;
3954 let (waitable, delete) = if is_host {
3955 let id = TableId::<HostTask>::new(rep);
3956 let task = concurrent_state.get_mut(id)?;
3957 match &task.state {
3958 HostTaskState::CalleeRunning(_) | HostTaskState::CalleeCancelling => {
3959 bail!(Trap::SubtaskDropNotResolved)
3960 }
3961 HostTaskState::CalleeDone { .. } => {}
3962 HostTaskState::CalleeStarted | HostTaskState::CalleeFinished(_) => {
3963 bail_bug!("invalid state for callee in `subtask.drop`")
3964 }
3965 }
3966
3967 (Waitable::Host(id), true)
3968 } else {
3969 let id = TableId::<GuestTask>::new(rep);
3970 let task = concurrent_state.get_mut(id)?;
3971 if task.lift_result.is_some() {
3972 bail!(Trap::SubtaskDropNotResolved);
3973 }
3974 (
3975 Waitable::Guest(id),
3976 concurrent_state.get_mut(id)?.ready_to_delete(),
3977 )
3978 };
3979
3980 waitable.common(concurrent_state)?.handle = None;
3981
3982 if waitable.take_event(concurrent_state)?.is_some() {
3985 bail!(Trap::SubtaskDropNotResolved);
3986 }
3987
3988 if delete {
3989 waitable.delete_from(store)?;
3990 }
3991
3992 log::trace!("subtask_drop {waitable:?} (handle {task_id})");
3993 Ok(())
3994 }
3995
3996 pub(crate) fn waitable_set_wait(
3998 self,
3999 store: &mut StoreOpaque,
4000 options: OptionsIndex,
4001 set: u32,
4002 payload: u32,
4003 ) -> Result<u32> {
4004 let &CanonicalOptions {
4005 instance: caller_instance,
4006 ..
4007 } = &self.id().get(store).component().env_component().options[options];
4008 let caller = self.runtime_instance(caller_instance);
4009 let rep = store
4010 .instance_state(self.runtime_instance(caller_instance))
4011 .handle_table()
4012 .waitable_set_rep(set)?;
4013
4014 self.waitable_check(
4015 store,
4016 caller,
4017 WaitableCheck::Wait,
4018 WaitableCheckParams {
4019 set: TableId::new(rep),
4020 options,
4021 payload,
4022 },
4023 )
4024 }
4025
4026 pub(crate) fn waitable_set_poll(
4028 self,
4029 store: &mut StoreOpaque,
4030 options: OptionsIndex,
4031 set: u32,
4032 payload: u32,
4033 ) -> Result<u32> {
4034 let &CanonicalOptions {
4035 instance: caller_instance,
4036 ..
4037 } = &self.id().get(store).component().env_component().options[options];
4038 let caller = self.runtime_instance(caller_instance);
4039 let rep = store
4040 .instance_state(caller)
4041 .handle_table()
4042 .waitable_set_rep(set)?;
4043
4044 self.waitable_check(
4045 store,
4046 caller,
4047 WaitableCheck::Poll,
4048 WaitableCheckParams {
4049 set: TableId::new(rep),
4050 options,
4051 payload,
4052 },
4053 )
4054 }
4055
4056 pub(crate) fn thread_index(&self, store: &mut dyn VMStore) -> Result<u32> {
4058 let thread_id = store.current_guest_thread()?.thread;
4059 match store
4060 .concurrent_state_mut()?
4061 .get_mut(thread_id)?
4062 .instance_rep
4063 {
4064 Some(r) => Ok(r),
4065 None => bail_bug!("thread should have instance_rep by now"),
4066 }
4067 }
4068
4069 pub(crate) fn thread_new_indirect<T: 'static>(
4071 self,
4072 mut store: StoreContextMut<T>,
4073 runtime_instance: RuntimeComponentInstanceIndex,
4074 _func_ty_idx: TypeFuncIndex, start_func_table_idx: RuntimeTableIndex,
4076 start_func_idx: u32,
4077 context: i32,
4078 ) -> Result<u32> {
4079 log::trace!("creating new thread");
4080
4081 let start_func_ty = FuncType::new(store.engine(), [ValType::I32], []);
4082 let (instance, registry) = self.id().get_mut_and_registry(store.0);
4083 let callee = instance
4084 .index_runtime_func_table(registry, start_func_table_idx, start_func_idx as u64)?
4085 .ok_or_else(|| Trap::ThreadNewIndirectUninitialized)?;
4086 if callee.type_index(store.0) != start_func_ty.type_index() {
4087 bail!(Trap::ThreadNewIndirectInvalidType);
4088 }
4089
4090 let token = StoreToken::new(store.as_context_mut());
4091 let start_func = Box::new(
4092 move |store: &mut dyn VMStore, guest_thread: QualifiedThreadId| -> Result<()> {
4093 let old_thread = store.set_thread(guest_thread)?;
4094 log::trace!(
4095 "thread start: replaced {old_thread:?} with {guest_thread:?} as current thread"
4096 );
4097
4098 let mut store = token.as_context_mut(store);
4099 let mut params = [ValRaw::i32(context)];
4100 unsafe { callee.call_unchecked(store.as_context_mut(), &mut params)? };
4103
4104 store.0.set_thread(old_thread)?;
4105
4106 let runtime_instance = self.runtime_instance(runtime_instance);
4107
4108 store
4111 .0
4112 .switch_or_trap_if_may_not_suspend(runtime_instance)?;
4113
4114 store
4115 .0
4116 .cleanup_thread(guest_thread, runtime_instance, CleanupTask::Yes)?;
4117
4118 log::trace!("explicit thread {guest_thread:?} completed");
4119 let state = store.0.concurrent_state_mut()?;
4120 if let Some(t) = old_thread.guest() {
4121 state.get_mut(t.thread)?.state = GuestThreadState::Running;
4122 }
4123 log::trace!("thread start: restored {old_thread:?} as current thread");
4124
4125 Ok(())
4126 },
4127 );
4128
4129 let current_thread = store.0.current_guest_thread()?;
4130 let state = store.0.concurrent_state_mut()?;
4131 let parent_task = current_thread.task;
4132
4133 let new_thread = GuestThread::new_explicit(state, parent_task, start_func)?;
4134 let thread_id = state.push(new_thread)?;
4135 state.get_mut(parent_task)?.threads.insert(thread_id);
4136
4137 log::trace!("new thread with id {thread_id:?} created");
4138
4139 self.add_guest_thread_to_instance_table(thread_id, store.0, runtime_instance)
4140 }
4141
4142 pub(crate) fn resume_thread(
4143 self,
4144 store: &mut StoreOpaque,
4145 runtime_instance: RuntimeComponentInstanceIndex,
4146 thread_idx: u32,
4147 how: ResumeThread,
4148 ) -> Result<bool> {
4149 let thread_id =
4150 GuestThread::from_instance(self.id().get_mut(store), runtime_instance, thread_idx)?;
4151 let state = store.concurrent_state_mut()?;
4152 let guest_thread = QualifiedThreadId::qualify(state, thread_id)?;
4153
4154 if store.current_guest_thread()? == guest_thread {
4155 bail!(Trap::CannotResumeThread);
4156 }
4157
4158 let state = store.concurrent_state_mut()?;
4159 let thread = state.get_mut(guest_thread.thread)?;
4160 let priority = match how {
4161 ResumeThread::Promote | ResumeThread::Resume => Priority::Switch,
4162 ResumeThread::ResumeLater => Priority::Low,
4163 };
4164
4165 match (&how, &thread.state) {
4166 (ResumeThread::Promote, GuestThreadState::Ready { .. }) => {}
4168 (ResumeThread::Promote, _) => return Ok(false),
4169
4170 (
4173 ResumeThread::Resume | ResumeThread::ResumeLater,
4174 GuestThreadState::NotStartedExplicit(_) | GuestThreadState::Suspended(_),
4175 ) => {}
4176 (ResumeThread::Resume | ResumeThread::ResumeLater, _) => {
4177 bail!(Trap::CannotResumeThread)
4178 }
4179 }
4180
4181 match mem::replace(&mut thread.state, GuestThreadState::Running) {
4182 GuestThreadState::NotStartedExplicit(start_func) => {
4183 log::trace!("starting thread {guest_thread:?}");
4184 let guest_call = WorkItem::GuestCall {
4185 instance: self.runtime_instance(runtime_instance),
4186 call: GuestCall {
4187 thread: guest_thread,
4188 kind: GuestCallKind::StartExplicit(Box::new(move |store| {
4189 start_func(store, guest_thread)
4190 })),
4191 },
4192 };
4193 store
4194 .concurrent_state_mut()?
4195 .push_work_item(guest_call, priority)?;
4196 }
4197 GuestThreadState::Suspended(fiber) => {
4198 log::trace!("resuming thread {thread_id:?} that was suspended");
4199 store.concurrent_state_mut()?.push_work_item(
4200 WorkItem::ResumeFiber {
4201 instance: self.runtime_instance(runtime_instance),
4202 thread: guest_thread,
4203 fiber,
4204 },
4205 priority,
4206 )?;
4207 }
4208 GuestThreadState::Ready { fiber } => {
4209 log::trace!("resuming thread {thread_id:?} that was ready");
4210 thread.state = GuestThreadState::Ready { fiber };
4211 store
4212 .concurrent_state_mut()?
4213 .promote_thread_work_item(guest_thread)?;
4214 }
4215 other @ (GuestThreadState::NotStartedImplicit
4216 | GuestThreadState::Running
4217 | GuestThreadState::Completed) => {
4218 thread.state = other;
4219 }
4220 }
4221 Ok(true)
4222 }
4223
4224 fn add_guest_thread_to_instance_table(
4225 self,
4226 thread_id: TableId<GuestThread>,
4227 store: &mut StoreOpaque,
4228 runtime_instance: RuntimeComponentInstanceIndex,
4229 ) -> Result<u32> {
4230 let guest_id = store
4231 .instance_state(self.runtime_instance(runtime_instance))
4232 .thread_handle_table()
4233 .guest_thread_insert(thread_id.rep())?;
4234 store
4235 .concurrent_state_mut()?
4236 .get_mut(thread_id)?
4237 .instance_rep = Some(guest_id);
4238 Ok(guest_id)
4239 }
4240
4241 pub(crate) fn suspension_intrinsic(
4245 self,
4246 store: &mut StoreOpaque,
4247 caller: RuntimeComponentInstanceIndex,
4248 yielding: bool,
4249 to_thread: SuspensionTarget,
4250 ) -> Result<WaitResult> {
4251 let check_suspend = match to_thread {
4252 SuspensionTarget::Promote(thread) => {
4253 !self.resume_thread(store, caller, thread, ResumeThread::Promote)?
4254 }
4255 SuspensionTarget::Resume(thread) => {
4256 if !self.resume_thread(store, caller, thread, ResumeThread::Resume)? {
4257 bail_bug!(
4258 "`resume_thread` should only ever return false \
4259 when `ResumeThread::Promote` is passed to it"
4260 );
4261 }
4262 false
4263 }
4264 SuspensionTarget::None => true,
4265 };
4266
4267 if check_suspend && !store.switch_if_may_not_suspend(self.runtime_instance(caller))? {
4268 return if yielding {
4269 Ok(WaitResult::Completed)
4270 } else {
4271 Err(Trap::CannotBlockSyncTask.into())
4272 };
4273 }
4274
4275 let guest_thread = store.current_guest_thread()?;
4276
4277 let reason = if yielding {
4278 SuspendReason::Yielding {
4279 thread: guest_thread,
4280 }
4281 } else {
4282 SuspendReason::ExplicitlySuspending {
4283 thread: guest_thread,
4284 }
4285 };
4286
4287 store.suspend(reason)?;
4288
4289 Ok(WaitResult::Completed)
4290 }
4291
4292 fn waitable_check(
4294 self,
4295 store: &mut StoreOpaque,
4296 caller: RuntimeInstance,
4297 check: WaitableCheck,
4298 params: WaitableCheckParams,
4299 ) -> Result<u32> {
4300 let guest_thread = store.current_guest_thread()?;
4301
4302 log::trace!("waitable check for {guest_thread:?}; set {:?}", params.set);
4303
4304 match &check {
4307 WaitableCheck::Wait => {
4308 let set = params.set;
4309
4310 loop {
4315 let state = store.concurrent_state_mut()?;
4316 let task = state.get_mut(guest_thread.task)?;
4317 if !(task.event.is_none() || matches!(task.event, Some(Event::Cancelled)))
4318 || !state.get_mut(set)?.ready.is_empty()
4319 {
4320 break;
4321 }
4322
4323 store.switch_or_trap_if_may_not_suspend(caller)?;
4324
4325 store.suspend(SuspendReason::Waiting {
4326 set,
4327 thread: guest_thread,
4328 })?;
4329 }
4330 }
4331 WaitableCheck::Poll => {}
4332 }
4333
4334 log::trace!(
4335 "waitable check for {guest_thread:?}; set {:?}, part two",
4336 params.set
4337 );
4338
4339 let event = self.get_event(store, guest_thread.task, Some(params.set), false)?;
4341
4342 let (ordinal, handle, result) = match &check {
4343 WaitableCheck::Wait => {
4344 let (event, waitable) = match event {
4345 Some(p) => p,
4346 None => bail_bug!("event expected to be present"),
4347 };
4348 let handle = waitable.map(|(_, v)| v).unwrap_or(0);
4349 let (ordinal, result) = event.parts();
4350 (ordinal, handle, result)
4351 }
4352 WaitableCheck::Poll => {
4353 if let Some((event, waitable)) = event {
4354 let handle = waitable.map(|(_, v)| v).unwrap_or(0);
4355 let (ordinal, result) = event.parts();
4356 (ordinal, handle, result)
4357 } else {
4358 log::trace!(
4359 "no events ready to deliver via waitable-set.poll to {:?}; set {:?}",
4360 guest_thread.task,
4361 params.set
4362 );
4363 let (ordinal, result) = Event::None.parts();
4364 (ordinal, 0, result)
4365 }
4366 }
4367 };
4368 let memory = self.options_memory_mut(store, params.options);
4369 let ptr = crate::component::func::validate_inbounds_dynamic(
4370 &CanonicalAbiInfo::POINTER_PAIR,
4371 memory,
4372 &ValRaw::u32(params.payload),
4373 )?;
4374 memory[ptr + 0..][..4].copy_from_slice(&handle.to_le_bytes());
4375 memory[ptr + 4..][..4].copy_from_slice(&result.to_le_bytes());
4376 Ok(ordinal)
4377 }
4378
4379 pub(crate) fn subtask_cancel(
4381 self,
4382 store: &mut StoreOpaque,
4383 caller_instance: RuntimeComponentInstanceIndex,
4384 async_: bool,
4385 task_id: u32,
4386 ) -> Result<u32> {
4387 let (rep, is_host) = store
4388 .instance_state(self.runtime_instance(caller_instance))
4389 .handle_table()
4390 .subtask_rep(task_id)?;
4391 let waitable = if is_host {
4392 Waitable::Host(TableId::<HostTask>::new(rep))
4393 } else {
4394 Waitable::Guest(TableId::<GuestTask>::new(rep))
4395 };
4396 let concurrent_state = store.concurrent_state_mut()?;
4397
4398 log::trace!("subtask_cancel {waitable:?} (handle {task_id}; async {async_})");
4399
4400 waitable.trap_if_in_waitable_set(concurrent_state)?;
4401
4402 let needs_block;
4403 if let Waitable::Host(host_task) = waitable {
4404 let state = &mut concurrent_state.get_mut(host_task)?.state;
4405 match state {
4406 HostTaskState::CalleeRunning(handle) => {
4413 handle.abort();
4414 *state = HostTaskState::CalleeCancelling;
4415 needs_block = true;
4416 }
4417
4418 HostTaskState::CalleeCancelling | HostTaskState::CalleeDone { cancelled: true } => {
4421 bail!(Trap::SubtaskCancelAfterTerminal);
4422 }
4423 HostTaskState::CalleeDone { cancelled: false } => {
4424 *state = HostTaskState::CalleeDone { cancelled: true };
4427 needs_block = false;
4428 }
4429
4430 HostTaskState::CalleeStarted | HostTaskState::CalleeFinished(_) => {
4433 bail_bug!("invalid states for host callee")
4434 }
4435 }
4436 } else {
4437 let guest_task = TableId::<GuestTask>::new(rep);
4438 let task = concurrent_state.get_mut(guest_task)?;
4439 if !task.already_lowered_parameters() {
4440 store.cancel_guest_subtask_without_lowered_parameters(
4441 self.runtime_instance(caller_instance),
4442 guest_task,
4443 )?;
4444 return Ok(Status::StartCancelled as u32);
4445 } else if !task.returned_or_cancelled() {
4446 task.event = Some(Event::Cancelled);
4454 let runtime_instance = task.instance;
4455 for thread in task.threads.clone() {
4456 let thread = QualifiedThreadId {
4457 task: guest_task,
4458 thread,
4459 };
4460 let concurrent_state = store.concurrent_state_mut()?;
4461 let thread_mut = concurrent_state.get_mut(thread.thread)?;
4462
4463 let yield_ = |store: &mut StoreOpaque| {
4464 let state = store.instance_state(runtime_instance).concurrent_state();
4469 let old_do_not_suspend = state.do_not_suspend;
4470 state.do_not_suspend = false;
4471
4472 let caller = store.current_guest_thread()?;
4473
4474 let state = store.concurrent_state_mut()?;
4479 let set = state.get_mut(caller.thread)?.sync_call_set;
4480 waitable.join(state, Some(set))?;
4481
4482 store.suspend(SuspendReason::YieldingToSubtask { thread: caller })?;
4483
4484 let state = store.concurrent_state_mut()?;
4485 waitable.join(state, None)?;
4486
4487 store
4488 .instance_state(runtime_instance)
4489 .concurrent_state()
4490 .do_not_suspend = old_do_not_suspend;
4491
4492 Ok::<(), crate::Error>(())
4493 };
4494
4495 match thread_mut.wake_on_cancel.take() {
4496 WakeOnCancel::Waiting(set) => {
4497 let item = match concurrent_state.get_mut(set)?.waiting.remove(&thread)
4499 {
4500 Some(WaitMode::Callback(instance)) => WorkItem::GuestCall {
4501 instance: runtime_instance,
4502 call: GuestCall {
4503 thread,
4504 kind: GuestCallKind::DeliverEvent {
4505 instance,
4506 set: Some(set),
4507 },
4508 },
4509 },
4510 other => bail_bug!(
4511 "expected `Some(WaitMode::Callback(_))`; got `{other:?}`"
4512 ),
4513 };
4514 concurrent_state.set_switch_item(item)?;
4515
4516 yield_(store)?;
4517
4518 break;
4519 }
4520 WakeOnCancel::Yielding => {
4521 if concurrent_state.promote_thread_work_item(thread)? {
4522 yield_(store)?;
4523 break;
4524 } else if store
4525 .instance_state(runtime_instance)
4526 .concurrent_state()
4527 .pending
4528 .contains_key(&thread)
4529 {
4530 store
4540 .concurrent_state_mut()?
4541 .get_mut(thread.thread)?
4542 .wake_on_cancel = WakeOnCancel::Yielding;
4543 } else {
4544 bail_bug!("thread with `WakeOnCancel::Yielding` not promotable");
4545 }
4546 }
4547 WakeOnCancel::None => {}
4548 }
4549 }
4550
4551 needs_block = !store
4554 .concurrent_state_mut()?
4555 .get_mut(guest_task)?
4556 .returned_or_cancelled()
4557 } else {
4558 needs_block = false;
4559 }
4560 };
4561
4562 if needs_block {
4566 if async_ {
4567 return Ok(BLOCKED);
4568 }
4569
4570 let old_next_switch_item = {
4573 let state = store.concurrent_state_mut()?;
4574 let item = state.next_switch_item.take();
4575 state.push(item)?
4579 };
4580
4581 store.wait_for_event(self.runtime_instance(caller_instance), waitable)?;
4584
4585 let state = store.concurrent_state_mut()?;
4586 state.next_switch_item = state.delete(old_next_switch_item)?;
4587
4588 }
4590
4591 let event = waitable.take_event(store.concurrent_state_mut()?)?;
4592 if let Some(Event::Subtask {
4593 status: status @ (Status::Returned | Status::ReturnCancelled),
4594 }) = event
4595 {
4596 Ok(status as u32)
4597 } else {
4598 bail!(Trap::SubtaskCancelAfterTerminal);
4599 }
4600 }
4601}
4602
4603pub trait VMComponentAsyncStore {
4611 unsafe fn prepare_call(
4617 &mut self,
4618 instance: Instance,
4619 memory: *mut VMMemoryDefinition,
4620 start: NonNull<VMFuncRef>,
4621 return_: NonNull<VMFuncRef>,
4622 caller_instance: RuntimeComponentInstanceIndex,
4623 callee_instance: RuntimeComponentInstanceIndex,
4624 task_return_type: TypeTupleIndex,
4625 callee_async: bool,
4626 string_encoding: StringEncoding,
4627 result_count: u32,
4628 storage: *mut ValRaw,
4629 storage_len: usize,
4630 ) -> Result<()>;
4631
4632 unsafe fn start_call(
4635 &mut self,
4636 instance: Instance,
4637 callback: *mut VMFuncRef,
4638 post_return: *mut VMFuncRef,
4639 callee: NonNull<VMFuncRef>,
4640 param_count: u32,
4641 result_count: u32,
4642 flags: u32,
4643 storage: *mut MaybeUninit<ValRaw>,
4644 storage_len: usize,
4645 ) -> Result<()>;
4646
4647 fn future_write(
4649 &mut self,
4650 instance: Instance,
4651 caller: RuntimeComponentInstanceIndex,
4652 ty: TypeFutureTableIndex,
4653 options: OptionsIndex,
4654 future: u32,
4655 address: u32,
4656 ) -> Result<u32>;
4657
4658 fn future_read(
4660 &mut self,
4661 instance: Instance,
4662 caller: RuntimeComponentInstanceIndex,
4663 ty: TypeFutureTableIndex,
4664 options: OptionsIndex,
4665 future: u32,
4666 address: u32,
4667 ) -> Result<u32>;
4668
4669 fn future_drop_writable(
4671 &mut self,
4672 instance: Instance,
4673 ty: TypeFutureTableIndex,
4674 writer: u32,
4675 ) -> Result<()>;
4676
4677 fn stream_write(
4679 &mut self,
4680 instance: Instance,
4681 caller: RuntimeComponentInstanceIndex,
4682 ty: TypeStreamTableIndex,
4683 options: OptionsIndex,
4684 stream: u32,
4685 address: u32,
4686 count: u32,
4687 ) -> Result<u32>;
4688
4689 fn stream_read(
4691 &mut self,
4692 instance: Instance,
4693 caller: RuntimeComponentInstanceIndex,
4694 ty: TypeStreamTableIndex,
4695 options: OptionsIndex,
4696 stream: u32,
4697 address: u32,
4698 count: u32,
4699 ) -> Result<u32>;
4700
4701 fn flat_stream_write(
4704 &mut self,
4705 instance: Instance,
4706 caller: RuntimeComponentInstanceIndex,
4707 ty: TypeStreamTableIndex,
4708 options: OptionsIndex,
4709 payload_size: u32,
4710 payload_align: u32,
4711 stream: u32,
4712 address: u32,
4713 count: u32,
4714 ) -> Result<u32>;
4715
4716 fn flat_stream_read(
4719 &mut self,
4720 instance: Instance,
4721 caller: RuntimeComponentInstanceIndex,
4722 ty: TypeStreamTableIndex,
4723 options: OptionsIndex,
4724 payload_size: u32,
4725 payload_align: u32,
4726 stream: u32,
4727 address: u32,
4728 count: u32,
4729 ) -> Result<u32>;
4730
4731 fn stream_drop_writable(
4733 &mut self,
4734 instance: Instance,
4735 ty: TypeStreamTableIndex,
4736 writer: u32,
4737 ) -> Result<()>;
4738
4739 fn error_context_debug_message(
4741 &mut self,
4742 instance: Instance,
4743 ty: TypeComponentLocalErrorContextTableIndex,
4744 options: OptionsIndex,
4745 err_ctx_handle: u32,
4746 debug_msg_address: u32,
4747 ) -> Result<()>;
4748
4749 fn thread_new_indirect(
4751 &mut self,
4752 instance: Instance,
4753 caller: RuntimeComponentInstanceIndex,
4754 func_ty_idx: TypeFuncIndex,
4755 start_func_table_idx: RuntimeTableIndex,
4756 start_func_idx: u32,
4757 context: i32,
4758 ) -> Result<u32>;
4759}
4760
4761impl<T: 'static> VMComponentAsyncStore for StoreInner<T> {
4763 unsafe fn prepare_call(
4764 &mut self,
4765 instance: Instance,
4766 memory: *mut VMMemoryDefinition,
4767 start: NonNull<VMFuncRef>,
4768 return_: NonNull<VMFuncRef>,
4769 caller_instance: RuntimeComponentInstanceIndex,
4770 callee_instance: RuntimeComponentInstanceIndex,
4771 task_return_type: TypeTupleIndex,
4772 callee_async: bool,
4773 string_encoding: StringEncoding,
4774 result_count_or_max_if_async: u32,
4775 storage: *mut ValRaw,
4776 storage_len: usize,
4777 ) -> Result<()> {
4778 let params = unsafe { core::slice::from_raw_parts(storage, storage_len) }.to_vec();
4782
4783 unsafe {
4784 instance.prepare_call(
4785 StoreContextMut(self),
4786 start,
4787 return_,
4788 caller_instance,
4789 callee_instance,
4790 task_return_type,
4791 callee_async,
4792 memory,
4793 string_encoding,
4794 match result_count_or_max_if_async {
4795 PREPARE_ASYNC_NO_RESULT => CallerInfo::Async {
4796 params,
4797 has_result: false,
4798 },
4799 PREPARE_ASYNC_WITH_RESULT => CallerInfo::Async {
4800 params,
4801 has_result: true,
4802 },
4803 result_count => CallerInfo::Sync {
4804 params,
4805 result_count,
4806 },
4807 },
4808 )
4809 }
4810 }
4811
4812 unsafe fn start_call(
4813 &mut self,
4814 instance: Instance,
4815 callback: *mut VMFuncRef,
4816 post_return: *mut VMFuncRef,
4817 callee: NonNull<VMFuncRef>,
4818 param_count: u32,
4819 result_count: u32,
4820 flags: u32,
4821 storage: *mut MaybeUninit<ValRaw>,
4822 storage_len: usize,
4823 ) -> Result<()> {
4824 unsafe {
4825 instance.start_call(
4826 StoreContextMut(self),
4827 callback,
4828 post_return,
4829 callee,
4830 param_count,
4831 result_count,
4832 flags,
4833 core::slice::from_raw_parts_mut(storage, storage_len),
4837 )
4838 }
4839 }
4840
4841 fn future_write(
4842 &mut self,
4843 instance: Instance,
4844 caller: RuntimeComponentInstanceIndex,
4845 ty: TypeFutureTableIndex,
4846 options: OptionsIndex,
4847 future: u32,
4848 address: u32,
4849 ) -> Result<u32> {
4850 instance
4851 .guest_write(
4852 StoreContextMut(self),
4853 caller,
4854 TransmitIndex::Future(ty),
4855 options,
4856 None,
4857 future,
4858 address,
4859 1,
4860 )
4861 .map(|result| result.encode())
4862 }
4863
4864 fn future_read(
4865 &mut self,
4866 instance: Instance,
4867 caller: RuntimeComponentInstanceIndex,
4868 ty: TypeFutureTableIndex,
4869 options: OptionsIndex,
4870 future: u32,
4871 address: u32,
4872 ) -> Result<u32> {
4873 instance
4874 .guest_read(
4875 StoreContextMut(self),
4876 caller,
4877 TransmitIndex::Future(ty),
4878 options,
4879 None,
4880 future,
4881 address,
4882 1,
4883 )
4884 .map(|result| result.encode())
4885 }
4886
4887 fn stream_write(
4888 &mut self,
4889 instance: Instance,
4890 caller: RuntimeComponentInstanceIndex,
4891 ty: TypeStreamTableIndex,
4892 options: OptionsIndex,
4893 stream: u32,
4894 address: u32,
4895 count: u32,
4896 ) -> Result<u32> {
4897 instance
4898 .guest_write(
4899 StoreContextMut(self),
4900 caller,
4901 TransmitIndex::Stream(ty),
4902 options,
4903 None,
4904 stream,
4905 address,
4906 count,
4907 )
4908 .map(|result| result.encode())
4909 }
4910
4911 fn stream_read(
4912 &mut self,
4913 instance: Instance,
4914 caller: RuntimeComponentInstanceIndex,
4915 ty: TypeStreamTableIndex,
4916 options: OptionsIndex,
4917 stream: u32,
4918 address: u32,
4919 count: u32,
4920 ) -> Result<u32> {
4921 instance
4922 .guest_read(
4923 StoreContextMut(self),
4924 caller,
4925 TransmitIndex::Stream(ty),
4926 options,
4927 None,
4928 stream,
4929 address,
4930 count,
4931 )
4932 .map(|result| result.encode())
4933 }
4934
4935 fn future_drop_writable(
4936 &mut self,
4937 instance: Instance,
4938 ty: TypeFutureTableIndex,
4939 writer: u32,
4940 ) -> Result<()> {
4941 instance.guest_drop_writable(self, TransmitIndex::Future(ty), writer)
4942 }
4943
4944 fn flat_stream_write(
4945 &mut self,
4946 instance: Instance,
4947 caller: RuntimeComponentInstanceIndex,
4948 ty: TypeStreamTableIndex,
4949 options: OptionsIndex,
4950 payload_size: u32,
4951 payload_align: u32,
4952 stream: u32,
4953 address: u32,
4954 count: u32,
4955 ) -> Result<u32> {
4956 instance
4957 .guest_write(
4958 StoreContextMut(self),
4959 caller,
4960 TransmitIndex::Stream(ty),
4961 options,
4962 Some(FlatAbi {
4963 size: payload_size,
4964 align: payload_align,
4965 }),
4966 stream,
4967 address,
4968 count,
4969 )
4970 .map(|result| result.encode())
4971 }
4972
4973 fn flat_stream_read(
4974 &mut self,
4975 instance: Instance,
4976 caller: RuntimeComponentInstanceIndex,
4977 ty: TypeStreamTableIndex,
4978 options: OptionsIndex,
4979 payload_size: u32,
4980 payload_align: u32,
4981 stream: u32,
4982 address: u32,
4983 count: u32,
4984 ) -> Result<u32> {
4985 instance
4986 .guest_read(
4987 StoreContextMut(self),
4988 caller,
4989 TransmitIndex::Stream(ty),
4990 options,
4991 Some(FlatAbi {
4992 size: payload_size,
4993 align: payload_align,
4994 }),
4995 stream,
4996 address,
4997 count,
4998 )
4999 .map(|result| result.encode())
5000 }
5001
5002 fn stream_drop_writable(
5003 &mut self,
5004 instance: Instance,
5005 ty: TypeStreamTableIndex,
5006 writer: u32,
5007 ) -> Result<()> {
5008 instance.guest_drop_writable(self, TransmitIndex::Stream(ty), writer)
5009 }
5010
5011 fn error_context_debug_message(
5012 &mut self,
5013 instance: Instance,
5014 ty: TypeComponentLocalErrorContextTableIndex,
5015 options: OptionsIndex,
5016 err_ctx_handle: u32,
5017 debug_msg_address: u32,
5018 ) -> Result<()> {
5019 instance.error_context_debug_message(
5020 StoreContextMut(self),
5021 ty,
5022 options,
5023 err_ctx_handle,
5024 debug_msg_address,
5025 )
5026 }
5027
5028 fn thread_new_indirect(
5029 &mut self,
5030 instance: Instance,
5031 caller: RuntimeComponentInstanceIndex,
5032 func_ty_idx: TypeFuncIndex,
5033 start_func_table_idx: RuntimeTableIndex,
5034 start_func_idx: u32,
5035 context: i32,
5036 ) -> Result<u32> {
5037 instance.thread_new_indirect(
5038 StoreContextMut(self),
5039 caller,
5040 func_ty_idx,
5041 start_func_table_idx,
5042 start_func_idx,
5043 context,
5044 )
5045 }
5046}
5047
5048type HostTaskFuture = Pin<Box<dyn Future<Output = Result<()>> + Send + 'static>>;
5049
5050async fn run_with_host_task_set<F>(task: TableId<HostTask>, future: F) -> Result<F::Output>
5053where
5054 F: Future,
5055{
5056 let mut future = pin!(future);
5057 future::poll_fn(|cx| {
5058 let old_thread = match tls::get(|store| store.set_thread(task)) {
5059 Ok(thread) => thread,
5060 Err(error) => return Poll::Ready(Err(error)),
5061 };
5062 let result = future.as_mut().poll(cx);
5063 match tls::get(|store| store.set_thread(old_thread)) {
5064 Ok(_) => result.map(Ok),
5065 Err(error) => Poll::Ready(Err(error)),
5066 }
5067 })
5068 .await
5069}
5070
5071pub(crate) struct HostTask {
5075 common: WaitableCommon,
5076
5077 call_context: CallContext,
5080
5081 state: HostTaskState,
5082
5083 group: TaskGroupId,
5084}
5085
5086enum HostTaskState {
5087 CalleeStarted,
5092
5093 CalleeRunning(JoinHandle),
5098
5099 CalleeCancelling,
5103
5104 CalleeFinished(LiftedResult),
5108
5109 CalleeDone { cancelled: bool },
5112}
5113
5114impl HostTask {
5115 fn new(
5116 concurrent_state: &mut ConcurrentState,
5117 state: HostTaskState,
5118 caller: QualifiedThreadId,
5119 ) -> Result<Self> {
5120 let group = concurrent_state.get_mut(caller.task)?.group;
5121 concurrent_state.increment_group_ref_count(group)?;
5122
5123 Ok(Self {
5124 common: WaitableCommon::default(),
5125 call_context: CallContext::default(),
5126 state,
5127 group,
5128 })
5129 }
5130}
5131
5132impl TableDebug for HostTask {
5133 fn type_name() -> &'static str {
5134 "HostTask"
5135 }
5136}
5137
5138type CallbackFn = Box<dyn Fn(&mut dyn VMStore, Event, u32) -> Result<u32> + Send + Sync + 'static>;
5139
5140enum Caller {
5142 Host {
5144 tx: Option<oneshot::Sender<LiftedResult>>,
5146 host_future_present: bool,
5149 caller: Option<TableId<HostTask>>,
5153 },
5154 Guest {
5156 thread: QualifiedThreadId,
5158 },
5159}
5160
5161struct LiftResult {
5164 lift: RawLift,
5165 ty: TypeTupleIndex,
5166 memory: Option<SendSyncPtr<VMMemoryDefinition>>,
5167 string_encoding: StringEncoding,
5168}
5169
5170#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
5175pub(crate) struct QualifiedThreadId {
5176 task: TableId<GuestTask>,
5177 thread: TableId<GuestThread>,
5178}
5179
5180impl QualifiedThreadId {
5181 fn qualify(
5182 state: &mut ConcurrentState,
5183 thread: TableId<GuestThread>,
5184 ) -> Result<QualifiedThreadId> {
5185 Ok(QualifiedThreadId {
5186 task: state.get_mut(thread)?.parent_task,
5187 thread,
5188 })
5189 }
5190}
5191
5192impl fmt::Debug for QualifiedThreadId {
5193 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
5194 f.debug_tuple("QualifiedThreadId")
5195 .field(&self.task.rep())
5196 .field(&self.thread.rep())
5197 .finish()
5198 }
5199}
5200
5201enum GuestThreadState {
5202 NotStartedImplicit,
5203 NotStartedExplicit(
5204 Box<dyn FnOnce(&mut dyn VMStore, QualifiedThreadId) -> Result<()> + Send + Sync>,
5205 ),
5206 Running,
5207 Suspended(StoreFiber<'static>),
5208 Ready {
5209 fiber: StoreFiber<'static>,
5210 },
5211 Completed,
5212}
5213
5214impl fmt::Debug for GuestThreadState {
5215 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
5216 match self {
5217 Self::NotStartedImplicit => f.debug_tuple("NotStartedImplicit").finish(),
5218 Self::NotStartedExplicit(_) => f.debug_tuple("NotStartedExplicit").finish(),
5219 Self::Running => f.debug_tuple("Running").finish(),
5220 Self::Suspended(_) => f.debug_tuple("Suspended").finish(),
5221 Self::Ready { .. } => f.debug_struct("Ready").finish(),
5222 Self::Completed => f.debug_tuple("Completed").finish(),
5223 }
5224 }
5225}
5226
5227#[derive(Copy, Clone, PartialEq, Eq, Debug)]
5228enum WakeOnCancel {
5229 None,
5230 Waiting(TableId<WaitableSet>),
5231 Yielding,
5232}
5233
5234impl WakeOnCancel {
5235 fn is_none(self) -> bool {
5236 matches!(self, WakeOnCancel::None)
5237 }
5238
5239 fn replace(&mut self, other: WakeOnCancel) -> Self {
5240 let old = *self;
5241 *self = other;
5242 old
5243 }
5244
5245 fn take(&mut self) -> Self {
5246 self.replace(WakeOnCancel::None)
5247 }
5248}
5249
5250pub struct GuestThread {
5251 context: [u32; NUM_COMPONENT_CONTEXT_SLOTS],
5254 parent_task: TableId<GuestTask>,
5256 wake_on_cancel: WakeOnCancel,
5259 state: GuestThreadState,
5261 instance_rep: Option<u32>,
5264 sync_call_set: TableId<WaitableSet>,
5266 old_do_not_suspend: Option<bool>,
5269}
5270
5271impl GuestThread {
5272 fn from_instance(
5275 state: Pin<&mut ComponentInstance>,
5276 caller_instance: RuntimeComponentInstanceIndex,
5277 guest_thread: u32,
5278 ) -> Result<TableId<Self>> {
5279 let rep = state.instance_states().0[caller_instance]
5280 .thread_handle_table()
5281 .guest_thread_rep(guest_thread)?;
5282 Ok(TableId::new(rep))
5283 }
5284
5285 fn new_implicit(state: &mut ConcurrentState, parent_task: TableId<GuestTask>) -> Result<Self> {
5286 let sync_call_set = state.push(WaitableSet {
5287 is_sync_call_set: true,
5288 ..WaitableSet::default()
5289 })?;
5290 Ok(Self {
5291 context: [0; NUM_COMPONENT_CONTEXT_SLOTS],
5292 parent_task,
5293 wake_on_cancel: WakeOnCancel::None,
5294 state: GuestThreadState::NotStartedImplicit,
5295 instance_rep: None,
5296 sync_call_set,
5297 old_do_not_suspend: None,
5298 })
5299 }
5300
5301 fn new_explicit(
5302 state: &mut ConcurrentState,
5303 parent_task: TableId<GuestTask>,
5304 start_func: Box<
5305 dyn FnOnce(&mut dyn VMStore, QualifiedThreadId) -> Result<()> + Send + Sync,
5306 >,
5307 ) -> Result<Self> {
5308 let sync_call_set = state.push(WaitableSet {
5309 is_sync_call_set: true,
5310 ..WaitableSet::default()
5311 })?;
5312 Ok(Self {
5313 context: [0; NUM_COMPONENT_CONTEXT_SLOTS],
5314 parent_task,
5315 wake_on_cancel: WakeOnCancel::None,
5316 state: GuestThreadState::NotStartedExplicit(start_func),
5317 instance_rep: None,
5318 sync_call_set,
5319 old_do_not_suspend: None,
5320 })
5321 }
5322}
5323
5324impl TableDebug for GuestThread {
5325 fn type_name() -> &'static str {
5326 "GuestThread"
5327 }
5328}
5329
5330enum SyncResult {
5331 NotProduced,
5332 Produced(Option<ValRaw>),
5333 Taken,
5334}
5335
5336impl SyncResult {
5337 fn take(&mut self) -> Result<Option<Option<ValRaw>>> {
5338 Ok(match mem::replace(self, SyncResult::Taken) {
5339 SyncResult::NotProduced => None,
5340 SyncResult::Produced(val) => Some(val),
5341 SyncResult::Taken => {
5342 bail_bug!("attempted to take a synchronous result that was already taken")
5343 }
5344 })
5345 }
5346}
5347
5348#[derive(Debug)]
5349enum HostFutureState {
5350 NotApplicable,
5351 Live,
5352 Dropped,
5353}
5354
5355pub(crate) struct GuestTask {
5357 common: WaitableCommon,
5359 lower_params: Option<RawLower>,
5361 lift_result: Option<LiftResult>,
5363 result: Option<LiftedResult>,
5366 callback: Option<CallbackFn>,
5369 caller: Caller,
5371 call_context: CallContext,
5376 sync_result: SyncResult,
5379 cancel_request_delivered: bool,
5383 starting_sent: bool,
5386 instance: RuntimeInstance,
5393 event: Option<Event>,
5395 exited: bool,
5397 threads: HashSet<TableId<GuestThread>>,
5399 host_future_state: HostFutureState,
5402 async_typed: bool,
5405 async_lifted: bool,
5408
5409 decremented_interesting_task_count: bool,
5410
5411 group: TaskGroupId,
5412}
5413
5414impl GuestTask {
5415 fn already_lowered_parameters(&self) -> bool {
5416 self.lower_params.is_none()
5418 }
5419
5420 fn returned_or_cancelled(&self) -> bool {
5421 self.lift_result.is_none()
5423 }
5424
5425 fn ready_to_delete(&self) -> bool {
5426 let threads_completed = self.threads.is_empty();
5427 let has_sync_result = matches!(self.sync_result, SyncResult::Produced(_));
5428 let pending_completion_event = matches!(
5429 self.common.event,
5430 Some(Event::Subtask {
5431 status: Status::Returned | Status::ReturnCancelled
5432 })
5433 );
5434 let ready = threads_completed
5435 && !has_sync_result
5436 && !pending_completion_event
5437 && !matches!(self.host_future_state, HostFutureState::Live);
5438 log::trace!(
5439 "ready to delete? {ready} (threads_completed: {}, has_sync_result: {}, pending_completion_event: {}, host_future_state: {:?})",
5440 threads_completed,
5441 has_sync_result,
5442 pending_completion_event,
5443 self.host_future_state
5444 );
5445 ready
5446 }
5447
5448 fn new(
5449 state: &mut ConcurrentState,
5450 lower_params: RawLower,
5451 lift_result: LiftResult,
5452 caller: Caller,
5453 callback: Option<CallbackFn>,
5454 instance: RuntimeInstance,
5455 async_typed: bool,
5456 async_lifted: bool,
5457 ) -> Result<QualifiedThreadId> {
5458 let host_future_state = match &caller {
5459 Caller::Guest { .. } => HostFutureState::NotApplicable,
5460 Caller::Host {
5461 host_future_present,
5462 ..
5463 } => {
5464 if *host_future_present {
5465 HostFutureState::Live
5466 } else {
5467 HostFutureState::NotApplicable
5468 }
5469 }
5470 };
5471
5472 let group = match caller {
5473 Caller::Guest { thread } => {
5474 let group = state.get_mut(thread.task)?.group;
5475 state.increment_group_ref_count(group)?;
5476 group
5477 }
5478 Caller::Host { .. } => state.make_task_group()?,
5479 };
5480
5481 let task = state.push(Self {
5482 common: WaitableCommon::default(),
5483 lower_params: Some(lower_params),
5484 lift_result: Some(lift_result),
5485 result: None,
5486 callback,
5487 caller,
5488 call_context: CallContext::default(),
5489 sync_result: SyncResult::NotProduced,
5490 cancel_request_delivered: false,
5491 starting_sent: false,
5492 instance,
5493 event: None,
5494 exited: false,
5495 threads: HashSet::new(),
5496 host_future_state,
5497 async_typed,
5498 async_lifted,
5499 decremented_interesting_task_count: false,
5500 group,
5501 })?;
5502 let new_thread = GuestThread::new_implicit(state, task)?;
5503 let thread = state.push(new_thread)?;
5504 state.get_mut(task)?.threads.insert(thread);
5505 state.interesting_tasks += 1;
5506 let thread = QualifiedThreadId { task, thread };
5507 log::trace!("new implicit thread {thread:?} for instance {instance:?}");
5508 Ok(thread)
5509 }
5510}
5511
5512impl TableDebug for GuestTask {
5513 fn type_name() -> &'static str {
5514 "GuestTask"
5515 }
5516}
5517
5518#[derive(Default)]
5520struct WaitableCommon {
5521 event: Option<Event>,
5523 set: Option<TableId<WaitableSet>>,
5525 handle: Option<u32>,
5527}
5528
5529#[derive(Copy, Clone, Ord, PartialOrd, Eq, PartialEq)]
5531enum Waitable {
5532 Host(TableId<HostTask>),
5534 Guest(TableId<GuestTask>),
5536 Transmit(TableId<TransmitHandle>),
5538}
5539
5540impl Waitable {
5541 fn from_instance(
5544 state: Pin<&mut ComponentInstance>,
5545 caller_instance: RuntimeComponentInstanceIndex,
5546 waitable: u32,
5547 ) -> Result<Self> {
5548 use crate::runtime::vm::component::Waitable;
5549
5550 let (waitable, kind) = state.instance_states().0[caller_instance]
5551 .handle_table()
5552 .waitable_rep(waitable)?;
5553
5554 Ok(match kind {
5555 Waitable::Subtask { is_host: true } => Self::Host(TableId::new(waitable)),
5556 Waitable::Subtask { is_host: false } => Self::Guest(TableId::new(waitable)),
5557 Waitable::Stream | Waitable::Future => Self::Transmit(TableId::new(waitable)),
5558 })
5559 }
5560
5561 fn rep(&self) -> u32 {
5563 match self {
5564 Self::Host(id) => id.rep(),
5565 Self::Guest(id) => id.rep(),
5566 Self::Transmit(id) => id.rep(),
5567 }
5568 }
5569
5570 fn join(&self, state: &mut ConcurrentState, set: Option<TableId<WaitableSet>>) -> Result<()> {
5574 log::trace!("waitable {self:?} join set {set:?}");
5575
5576 let old = mem::replace(&mut self.common(state)?.set, set);
5577
5578 if let Some(old) = old {
5579 match *self {
5580 Waitable::Host(id) => state.remove_child(id, old),
5581 Waitable::Guest(id) => state.remove_child(id, old),
5582 Waitable::Transmit(id) => state.remove_child(id, old),
5583 }?;
5584
5585 state.get_mut(old)?.ready.remove(self);
5586 }
5587
5588 if let Some(set) = set {
5589 match *self {
5590 Waitable::Host(id) => state.add_child(id, set),
5591 Waitable::Guest(id) => state.add_child(id, set),
5592 Waitable::Transmit(id) => state.add_child(id, set),
5593 }?;
5594
5595 if self.common(state)?.event.is_some() {
5596 self.mark_ready(state)?;
5597 }
5598 }
5599
5600 Ok(())
5601 }
5602
5603 fn common<'a>(&self, state: &'a mut ConcurrentState) -> Result<&'a mut WaitableCommon> {
5605 Ok(match self {
5606 Self::Host(id) => &mut state.get_mut(*id)?.common,
5607 Self::Guest(id) => &mut state.get_mut(*id)?.common,
5608 Self::Transmit(id) => &mut state.get_mut(*id)?.common,
5609 })
5610 }
5611
5612 fn trap_if_in_waitable_set(&self, state: &mut ConcurrentState) -> Result<()> {
5618 if self.common(state)?.set.is_some() {
5619 bail!(Trap::WaitableSyncAndAsync);
5620 }
5621 Ok(())
5622 }
5623
5624 fn set_event(&self, state: &mut ConcurrentState, event: Option<Event>) -> Result<()> {
5628 log::trace!("set event for {self:?}: {event:?}");
5629 self.common(state)?.event = event;
5630 self.mark_ready(state)
5631 }
5632
5633 fn take_event(&self, state: &mut ConcurrentState) -> Result<Option<Event>> {
5635 let common = self.common(state)?;
5636 let event = common.event.take();
5637 if let Some(set) = self.common(state)?.set {
5638 state.get_mut(set)?.ready.remove(self);
5639 }
5640
5641 Ok(event)
5642 }
5643
5644 fn mark_ready(&self, state: &mut ConcurrentState) -> Result<()> {
5648 if let Some(set) = self.common(state)?.set {
5649 let set_state = state.get_mut(set)?;
5650 set_state.ready.insert(*self);
5651
5652 if let Some((thread, mode)) = set_state.waiting.pop_first() {
5653 let wake_on_cancel = state.get_mut(thread.thread)?.wake_on_cancel.take();
5654 assert!(wake_on_cancel.is_none() || wake_on_cancel == WakeOnCancel::Waiting(set));
5655
5656 let item = match mode {
5657 WaitMode::Fiber(fiber) => Some(WorkItem::ResumeFiber {
5658 instance: state.get_mut(thread.task)?.instance,
5659 thread,
5660 fiber,
5661 }),
5662 WaitMode::Callback(instance) => Some(WorkItem::GuestCall {
5663 instance: state.get_mut(thread.task)?.instance,
5664 call: GuestCall {
5665 thread,
5666 kind: GuestCallKind::DeliverEvent {
5667 instance,
5668 set: Some(set),
5669 },
5670 },
5671 }),
5672 };
5673
5674 if let Some(item) = item {
5675 state.push_high_priority(item);
5676 }
5677 }
5678 }
5679 Ok(())
5680 }
5681
5682 fn delete_from(&self, store: &mut StoreOpaque) -> Result<()> {
5684 match self {
5685 Self::Host(task) => {
5686 log::trace!("delete host task {task:?}");
5687 let state = store.concurrent_state_mut()?;
5688 let task = state.delete(*task)?;
5689
5690 state.decrement_group_ref_count(task.group)?;
5691 }
5692 Self::Guest(task) => {
5693 log::trace!("delete guest task {task:?}");
5694 let state = store.concurrent_state_mut()?;
5695 let task = state.delete(*task)?;
5696
5697 state.decrement_group_ref_count(task.group)?;
5698
5699 debug_assert!(task.decremented_interesting_task_count);
5706 }
5707 Self::Transmit(task) => {
5708 store.concurrent_state_mut()?.delete(*task)?;
5709 }
5710 }
5711
5712 Ok(())
5713 }
5714}
5715
5716impl fmt::Debug for Waitable {
5717 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
5718 match self {
5719 Self::Host(id) => write!(f, "{id:?}"),
5720 Self::Guest(id) => write!(f, "{id:?}"),
5721 Self::Transmit(id) => write!(f, "{id:?}"),
5722 }
5723 }
5724}
5725
5726#[derive(Default)]
5728struct WaitableSet {
5729 ready: BTreeSet<Waitable>,
5731 waiting: BTreeMap<QualifiedThreadId, WaitMode>,
5733 num_waiting: usize,
5736 is_sync_call_set: bool,
5739}
5740
5741impl WaitableSet {
5742 fn stop_waiting(&mut self) -> Result<()> {
5744 self.num_waiting = match self.num_waiting.checked_sub(1) {
5745 Some(n) => n,
5746 None => bail_bug!("waiter not accounted for in waitable set"),
5747 };
5748 Ok(())
5749 }
5750}
5751
5752impl TableDebug for WaitableSet {
5753 fn type_name() -> &'static str {
5754 "WaitableSet"
5755 }
5756}
5757
5758type RawLower =
5760 Box<dyn FnOnce(&mut dyn VMStore, &mut [MaybeUninit<ValRaw>]) -> Result<()> + Send + Sync>;
5761
5762type RawLift = Box<
5764 dyn FnOnce(&mut dyn VMStore, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>> + Send + Sync,
5765>;
5766
5767type LiftedResult = Box<dyn Any + Send + Sync>;
5771
5772struct DummyResult;
5775
5776#[derive(Default)]
5778pub struct ConcurrentInstanceState {
5779 backpressure: u16,
5781 do_not_enter: bool,
5783 do_not_suspend: bool,
5786 pending: BTreeMap<QualifiedThreadId, GuestCallKind>,
5789}
5790
5791impl ConcurrentInstanceState {
5792 pub fn pending_is_empty(&self) -> bool {
5793 self.pending.is_empty()
5794 }
5795}
5796
5797#[derive(Debug, Copy, Clone)]
5798pub(crate) enum CurrentThread {
5799 Guest(QualifiedThreadId),
5802 Host(TableId<HostTask>),
5804 DeferredHost(QualifiedThreadId),
5807 None,
5810}
5811
5812impl CurrentThread {
5813 fn guest(&self) -> Option<&QualifiedThreadId> {
5814 match self {
5815 Self::Guest(id) => Some(id),
5816 _ => None,
5817 }
5818 }
5819
5820 fn guest_task(&self) -> Option<TableId<GuestTask>> {
5821 match self {
5822 Self::Guest(id) => Some(id.task),
5823 _ => None,
5824 }
5825 }
5826
5827 fn is_none(&self) -> bool {
5828 matches!(self, Self::None)
5829 }
5830}
5831
5832impl From<QualifiedThreadId> for CurrentThread {
5833 fn from(id: QualifiedThreadId) -> Self {
5834 Self::Guest(id)
5835 }
5836}
5837
5838impl From<TableId<HostTask>> for CurrentThread {
5839 fn from(id: TableId<HostTask>) -> Self {
5840 Self::Host(id)
5841 }
5842}
5843
5844enum Priority {
5845 Switch,
5846 High,
5847 Low,
5848}
5849
5850pub struct ConcurrentState {
5852 unforced_current_thread: CurrentThread,
5858
5859 deferred_host_call_context: Option<CallContext>,
5865
5866 futures: AlwaysMut<Option<FuturesUnordered<HostTaskFuture>>>,
5871 table: AlwaysMut<ResourceTable>,
5873 switch_item: Option<WorkItem>,
5881 next_switch_item: Option<WorkItem>,
5887 high_priority: VecDeque<WorkItem>,
5889 low_priority: VecDeque<WorkItem>,
5891 suspend_reason: Option<SuspendReason>,
5895 worker: Option<StoreFiber<'static>>,
5899 worker_item: Option<WorkerItem>,
5901
5902 global_error_context_ref_counts:
5915 BTreeMap<TypeComponentGlobalErrorContextTableIndex, GlobalErrorContextRefCount>,
5916
5917 interesting_tasks: usize,
5930
5931 interesting_tasks_empty_waker: Option<Waker>,
5935
5936 ready_for_concurrent_call_waker: Option<Waker>,
5941
5942 event_loop_running: bool,
5944
5945 #[cfg(feature = "task-group-hook")]
5947 task_group_hook: Option<Box<dyn TaskGroupHook>>,
5948}
5949
5950impl Default for ConcurrentState {
5951 fn default() -> Self {
5952 Self {
5953 unforced_current_thread: CurrentThread::None,
5954 deferred_host_call_context: None,
5955 table: AlwaysMut::new(ResourceTable::new()),
5956 futures: AlwaysMut::new(Some(FuturesUnordered::new())),
5957 switch_item: None,
5958 next_switch_item: None,
5959 high_priority: VecDeque::new(),
5960 low_priority: VecDeque::new(),
5961 suspend_reason: None,
5962 worker: None,
5963 worker_item: None,
5964 global_error_context_ref_counts: BTreeMap::new(),
5965 interesting_tasks: 0,
5966 interesting_tasks_empty_waker: None,
5967 ready_for_concurrent_call_waker: None,
5968 event_loop_running: false,
5969 #[cfg(feature = "task-group-hook")]
5970 task_group_hook: None,
5971 }
5972 }
5973}
5974
5975impl ConcurrentState {
5976 pub(crate) fn take_fibers_and_futures(
5993 &mut self,
5994 fibers: &mut Vec<StoreFiber<'static>>,
5995 futures: &mut Vec<FuturesUnordered<HostTaskFuture>>,
5996 ) {
5997 let mut items = Vec::new();
5998 for (_, entry) in self.table.get_mut().iter_mut() {
5999 if let Some(set) = entry.downcast_mut::<WaitableSet>() {
6000 for mode in mem::take(&mut set.waiting).into_values() {
6001 match mode {
6002 WaitMode::Fiber(fiber) => {
6003 fibers.push(fiber);
6004 }
6005 WaitMode::Callback(_) => {}
6006 }
6007 }
6008 } else if let Some(thread) = entry.downcast_mut::<GuestThread>() {
6009 if let GuestThreadState::Suspended(fiber) | GuestThreadState::Ready { fiber, .. } =
6010 mem::replace(&mut thread.state, GuestThreadState::Completed)
6011 {
6012 fibers.push(fiber);
6013 }
6014 } else if let Some(item) = entry.downcast_mut::<Option<WorkItem>>() {
6015 if let Some(item) = item.take() {
6016 items.push(item);
6017 }
6018 }
6019 }
6020
6021 if let Some(fiber) = self.worker.take() {
6022 fibers.push(fiber);
6023 }
6024
6025 let mut handle_item = |item| match item {
6026 WorkItem::ResumeFiber { fiber, .. } => {
6027 fibers.push(fiber);
6028 }
6029 WorkItem::PushFuture(future) => {
6030 self.futures
6031 .get_mut()
6032 .as_mut()
6033 .unwrap()
6034 .push(future.into_inner());
6035 }
6036 WorkItem::ResumeThread { .. }
6037 | WorkItem::GuestCall { .. }
6038 | WorkItem::WorkerFunction(_) => {}
6039 };
6040
6041 for item in items {
6042 handle_item(item);
6043 }
6044 if let Some(item) = self.switch_item.take() {
6045 handle_item(item);
6046 }
6047 if let Some(item) = self.next_switch_item.take() {
6048 handle_item(item);
6049 }
6050 for item in mem::take(&mut self.high_priority) {
6051 handle_item(item);
6052 }
6053 for item in mem::take(&mut self.low_priority) {
6054 handle_item(item);
6055 }
6056
6057 if let Some(them) = self.futures.get_mut().take() {
6058 futures.push(them);
6059 }
6060 }
6061
6062 #[cfg(feature = "gc")]
6063 pub(crate) fn trace_fiber_roots(
6064 &mut self,
6065 modules: &ModuleRegistry,
6066 unwind: &dyn Unwind,
6067 gc_roots_list: &mut GcRootsList,
6068 ) {
6069 let ConcurrentState {
6070 table,
6071 worker,
6072 switch_item,
6073 next_switch_item,
6074 high_priority,
6075 low_priority,
6076
6077 futures: _,
6081
6082 worker_item: _,
6084 unforced_current_thread: _,
6085 deferred_host_call_context: _,
6086 suspend_reason: _,
6087 global_error_context_ref_counts: _,
6088 interesting_tasks: _,
6089 interesting_tasks_empty_waker: _,
6090 ready_for_concurrent_call_waker: _,
6091 event_loop_running: _,
6092 #[cfg(feature = "task-group-hook")]
6093 task_group_hook: _,
6094 } = self;
6095
6096 for (_, entry) in table.get_mut().iter_mut() {
6097 if let Some(set) = entry.downcast_mut::<WaitableSet>() {
6098 for mode in set.waiting.values_mut() {
6099 match mode {
6100 WaitMode::Fiber(fiber) => {
6101 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6102 }
6103 WaitMode::Callback(_) => {}
6104 }
6105 }
6106 } else if let Some(thread) = entry.downcast_mut::<GuestThread>() {
6107 if let GuestThreadState::Suspended(fiber) | GuestThreadState::Ready { fiber, .. } =
6108 &mut thread.state
6109 {
6110 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6111 }
6112 } else if let Some(Some(WorkItem::ResumeFiber { fiber, .. })) =
6113 entry.downcast_mut::<Option<WorkItem>>()
6114 {
6115 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6116 }
6117 }
6118
6119 if let Some(fiber) = worker {
6120 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6121 }
6122
6123 let mut handle_item = |item: &mut WorkItem| match item {
6124 WorkItem::ResumeFiber { fiber, .. } => {
6125 fiber.trace_gc_roots(modules, unwind, gc_roots_list);
6126 }
6127 WorkItem::PushFuture(_future) => {
6128 }
6131 WorkItem::ResumeThread { .. }
6132 | WorkItem::GuestCall { .. }
6133 | WorkItem::WorkerFunction(_) => {}
6134 };
6135
6136 if let Some(item) = switch_item {
6137 handle_item(item);
6138 }
6139 if let Some(item) = next_switch_item {
6140 handle_item(item);
6141 }
6142 for item in high_priority {
6143 handle_item(item);
6144 }
6145 for item in low_priority {
6146 handle_item(item);
6147 }
6148 }
6149
6150 fn push<V: Send + Sync + 'static>(
6151 &mut self,
6152 value: V,
6153 ) -> Result<TableId<V>, ResourceTableError> {
6154 self.table.get_mut().push(value).map(TableId::from)
6155 }
6156
6157 fn get_mut<V: 'static>(&mut self, id: TableId<V>) -> Result<&mut V, ResourceTableError> {
6158 self.table.get_mut().get_mut(&Resource::from(id))
6159 }
6160
6161 pub fn add_child<T: 'static, U: 'static>(
6162 &mut self,
6163 child: TableId<T>,
6164 parent: TableId<U>,
6165 ) -> Result<(), ResourceTableError> {
6166 self.table
6167 .get_mut()
6168 .add_child(Resource::from(child), Resource::from(parent))
6169 }
6170
6171 pub fn remove_child<T: 'static, U: 'static>(
6172 &mut self,
6173 child: TableId<T>,
6174 parent: TableId<U>,
6175 ) -> Result<(), ResourceTableError> {
6176 self.table
6177 .get_mut()
6178 .remove_child(Resource::from(child), Resource::from(parent))
6179 }
6180
6181 fn delete<V: 'static>(&mut self, id: TableId<V>) -> Result<V, ResourceTableError> {
6182 self.table.get_mut().delete(Resource::from(id))
6183 }
6184
6185 fn push_future(&mut self, future: HostTaskFuture) {
6186 self.push_high_priority(WorkItem::PushFuture(AlwaysMut::new(future)));
6193 }
6194
6195 fn set_switch_item(&mut self, item: WorkItem) -> Result<()> {
6196 log::trace!("set switch item: {item:?}");
6197
6198 if self.switch_item.is_some() {
6199 bail_bug!("switch item already set");
6200 }
6201
6202 self.switch_item = Some(item);
6203
6204 Ok(())
6205 }
6206
6207 fn take_next_switch_item(&mut self) -> Result<()> {
6208 if let Some(item) = self.next_switch_item.take() {
6209 self.set_switch_item(item)?;
6210 }
6211 Ok(())
6212 }
6213
6214 fn push_high_priority(&mut self, item: WorkItem) {
6215 log::trace!("push high priority: {item:?}");
6216 self.high_priority.push_front(item);
6217 }
6218
6219 fn push_low_priority(&mut self, item: WorkItem) {
6220 log::trace!("push low priority: {item:?}");
6221 self.low_priority.push_front(item);
6222 }
6223
6224 fn push_work_item(&mut self, item: WorkItem, priority: Priority) -> Result<()> {
6225 match priority {
6226 Priority::Switch => self.set_switch_item(item)?,
6227 Priority::High => self.push_high_priority(item),
6228 Priority::Low => self.push_low_priority(item),
6229 }
6230
6231 Ok(())
6232 }
6233
6234 fn promote_instance_local_thread_work_item(
6235 &mut self,
6236 current_instance: RuntimeInstance,
6237 ) -> Result<bool> {
6238 log::trace!("promote thread work items for {current_instance:?}");
6239
6240 self.promote_work_item_matching(|item: &WorkItem| {
6241 let result = match item {
6242 WorkItem::ResumeThread { instance, .. }
6243 | WorkItem::ResumeFiber { instance, .. }
6244 | WorkItem::GuestCall { instance, .. } => *instance == current_instance,
6245 _ => false,
6246 };
6247
6248 log::trace!("candidate {item:?}: {result}");
6249 result
6250 })
6251 }
6252
6253 fn promote_thread_work_item(&mut self, thread: QualifiedThreadId) -> Result<bool> {
6254 self.promote_work_item_matching(|item: &WorkItem| match item {
6255 WorkItem::ResumeThread {
6256 thread: item_thread,
6257 ..
6258 }
6259 | WorkItem::GuestCall {
6260 call:
6261 GuestCall {
6262 thread: item_thread,
6263 ..
6264 },
6265 ..
6266 } => *item_thread == thread,
6267 _ => false,
6268 })
6269 }
6270
6271 fn promote_work_item_matching<F>(&mut self, mut predicate: F) -> Result<bool>
6272 where
6273 F: FnMut(&WorkItem) -> bool,
6274 {
6275 for item in mem::take(&mut self.high_priority).into_iter().rev() {
6280 if self.switch_item.is_none() && predicate(&item) {
6281 self.set_switch_item(item)?;
6282 } else {
6283 self.push_high_priority(item);
6284 }
6285 }
6286
6287 if self.switch_item.is_none() {
6288 for item in mem::take(&mut self.low_priority).into_iter().rev() {
6289 if self.switch_item.is_none() && predicate(&item) {
6290 self.set_switch_item(item)?;
6291 } else {
6292 self.push_low_priority(item);
6293 }
6294 }
6295 }
6296
6297 Ok(self.switch_item.is_some())
6298 }
6299
6300 pub fn call_context(&mut self, task: Scope) -> Result<&mut CallContext> {
6303 match task {
6304 Scope::HostId(task) => {
6305 let task: TableId<HostTask> = TableId::new(task);
6306 Ok(&mut self.get_mut(task)?.call_context)
6307 }
6308 Scope::Id(task) => {
6309 let task: TableId<GuestTask> = TableId::new(task);
6310 Ok(&mut self.get_mut(task)?.call_context)
6311 }
6312 }
6313 }
6314
6315 pub(crate) fn deferred_host_call_context(&mut self) -> Option<&mut CallContext> {
6316 self.deferred_host_call_context.as_mut()
6317 }
6318
6319 fn futures_mut(&mut self) -> Result<&mut FuturesUnordered<HostTaskFuture>> {
6320 match self.futures.get_mut().as_mut() {
6321 Some(f) => Ok(f),
6322 None => bail_bug!("futures field of concurrent state is currently taken"),
6323 }
6324 }
6325
6326 pub(crate) fn table(&mut self) -> &mut ResourceTable {
6327 self.table.get_mut()
6328 }
6329
6330 fn debug_assert_deferred_host_invariant(&self) {
6331 debug_assert_eq!(
6332 self.deferred_host_call_context.is_some(),
6333 matches!(self.unforced_current_thread, CurrentThread::DeferredHost(_)),
6334 "a deferred host thread and call context must exist together",
6335 );
6336 }
6337
6338 fn materialize_host_task(&mut self) -> Result<CurrentThread> {
6339 self.debug_assert_deferred_host_invariant();
6340 let caller = match self.unforced_current_thread {
6341 CurrentThread::DeferredHost(caller) => caller,
6342 thread => return Ok(thread),
6343 };
6344
6345 let task = HostTask::new(self, HostTaskState::CalleeStarted, caller)?;
6347 let task = self.push(task)?;
6348 let call_context = self
6349 .deferred_host_call_context
6350 .take()
6351 .expect("deferred host call context should be present");
6352 self.get_mut(task)
6353 .expect("newly inserted host task should be present")
6354 .call_context = call_context;
6355 self.unforced_current_thread = CurrentThread::Host(task);
6356 self.debug_assert_deferred_host_invariant();
6357 log::trace!("new host task materialized {task:?}");
6358 Ok(CurrentThread::Host(task))
6359 }
6360
6361 fn materialize_current_host_task_id(&mut self) -> Result<Option<TableId<HostTask>>> {
6362 match self.materialize_host_task()? {
6363 CurrentThread::Host(id) => Ok(Some(id)),
6364 CurrentThread::None => Ok(None),
6365 CurrentThread::Guest(_) => {
6366 bail_bug!("tried to materialize a host task id from a guest thread")
6367 }
6368 CurrentThread::DeferredHost(_) => {
6369 bail_bug!(
6370 "current thread is a deferred host thread which should have been materialized"
6371 )
6372 }
6373 }
6374 }
6375
6376 pub(crate) fn materialize_current_scope(&mut self) -> Result<Scope> {
6377 match self.materialize_host_task()? {
6378 CurrentThread::Host(id) => Ok(Scope::HostId(id.rep())),
6379 _ => bail_bug!("current scope is not a deferred host scope"),
6380 }
6381 }
6382}
6383
6384fn for_any_lower<
6387 F: FnOnce(&mut dyn VMStore, &mut [MaybeUninit<ValRaw>]) -> Result<()> + Send + Sync,
6388>(
6389 fun: F,
6390) -> F {
6391 fun
6392}
6393
6394fn for_any_lift<
6396 F: FnOnce(&mut dyn VMStore, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>> + Send + Sync,
6397>(
6398 fun: F,
6399) -> F {
6400 fun
6401}
6402
6403fn check_ambient_store(id: StoreId) {
6404 let message = "\
6405 `Future`s which depend on asynchronous component tasks, streams, or \
6406 futures to complete may only be polled from the event loop of the \
6407 store to which they belong. Please use \
6408 `StoreContextMut::{run_concurrent,spawn}` to poll or await them.\
6409 ";
6410 tls::try_get(|store| {
6411 let matched = match store {
6412 tls::TryGet::Some(store) => store.id() == id,
6413 tls::TryGet::Taken | tls::TryGet::None => false,
6414 };
6415
6416 if !matched {
6417 panic!("{message}")
6418 }
6419 });
6420}
6421
6422fn unpack_callback_code(code: u32) -> (u32, u32) {
6423 (code & 0xF, code >> 4)
6424}
6425
6426struct WaitableCheckParams {
6430 set: TableId<WaitableSet>,
6431 options: OptionsIndex,
6432 payload: u32,
6433}
6434
6435enum WaitableCheck {
6438 Wait,
6439 Poll,
6440}
6441
6442pub(crate) struct PreparedCall<R> {
6444 handle: Func,
6446 thread: QualifiedThreadId,
6448 param_count: usize,
6450 rx: oneshot::Receiver<LiftedResult>,
6453 runtime_instance: RuntimeInstance,
6455 _phantom: PhantomData<R>,
6456}
6457
6458impl<R> PreparedCall<R> {
6459 pub(crate) fn task_id(&self) -> TaskId {
6461 TaskId {
6462 task: self.thread.task,
6463 runtime_instance: self.runtime_instance,
6464 }
6465 }
6466}
6467
6468pub(crate) struct TaskId {
6470 task: TableId<GuestTask>,
6471 runtime_instance: RuntimeInstance,
6472}
6473
6474impl TaskId {
6475 pub(crate) fn host_future_dropped(&self, store: &mut StoreOpaque) -> Result<()> {
6481 let task = store.concurrent_state_mut()?.get_mut(self.task)?;
6482 let delete = if !task.already_lowered_parameters() {
6483 store.cancel_guest_subtask_without_lowered_parameters(
6484 self.runtime_instance,
6485 self.task,
6486 )?;
6487 true
6488 } else {
6489 task.host_future_state = HostFutureState::Dropped;
6490 task.ready_to_delete()
6491 };
6492 if delete {
6493 Waitable::Guest(self.task).delete_from(store)?
6494 }
6495 Ok(())
6496 }
6497}
6498
6499pub(crate) fn prepare_call<T, R>(
6505 mut store: StoreContextMut<T>,
6506 handle: Func,
6507 param_count: usize,
6508 host_future_present: bool,
6509 lower_params: impl FnOnce(StoreContextMut<T>, &mut [MaybeUninit<ValRaw>]) -> Result<()>
6510 + Send
6511 + Sync
6512 + 'static,
6513 lift_result: impl FnOnce(&mut StoreOpaque, &[ValRaw]) -> Result<Box<dyn Any + Send + Sync>>
6514 + Send
6515 + Sync
6516 + 'static,
6517) -> Result<PreparedCall<R>> {
6518 if !store.0.may_enter() {
6519 bail!(Trap::CannotEnterComponent);
6520 }
6521
6522 let (options, _flags, ty, raw_options) = handle.abi_info(store.0);
6523
6524 let instance = handle.instance().id().get(store.0);
6525 let options = &instance.component().env_component().options[options];
6526 let ty = &instance.component().types()[ty];
6527 let async_typed = ty.async_;
6528 let async_lifted = raw_options.async_;
6529 let task_return_type = ty.results;
6530 let component_instance = raw_options.instance;
6531 let callback = options.callback.map(|i| instance.runtime_callback(i));
6532 let memory = options
6533 .memory()
6534 .map(|i| instance.runtime_memory(i))
6535 .map(SendSyncPtr::new);
6536 let string_encoding = options.string_encoding;
6537 let token = StoreToken::new(store.as_context_mut());
6538 let caller = store.0.materialize_host_task_id()?;
6539 let state = store.0.concurrent_state_mut()?;
6540
6541 let (tx, rx) = oneshot::channel();
6542
6543 let instance = handle.instance().runtime_instance(component_instance);
6544 let thread = GuestTask::new(
6545 state,
6546 Box::new(for_any_lower(move |store, params| {
6547 lower_params(token.as_context_mut(store), params)
6548 })),
6549 LiftResult {
6550 lift: Box::new(for_any_lift(move |store, result| {
6551 lift_result(store, result)
6552 })),
6553 ty: task_return_type,
6554 memory,
6555 string_encoding,
6556 },
6557 Caller::Host {
6558 tx: Some(tx),
6559 host_future_present,
6560 caller,
6561 },
6562 callback.map(|callback| {
6563 let callback = SendSyncPtr::new(callback);
6564 let instance = handle.instance();
6565 Box::new(move |store: &mut dyn VMStore, event, handle| {
6566 let store = token.as_context_mut(store);
6567 unsafe { instance.call_callback(store, callback, event, handle) }
6570 }) as CallbackFn
6571 }),
6572 instance,
6573 async_typed,
6574 async_lifted,
6575 )?;
6576
6577 Ok(PreparedCall {
6578 handle,
6579 thread,
6580 param_count,
6581 runtime_instance: instance,
6582 rx,
6583 _phantom: PhantomData,
6584 })
6585}
6586
6587pub(crate) struct StagedCall<R> {
6588 store: StoreId,
6589 rx: oneshot::Receiver<LiftedResult>,
6590 _marker: PhantomData<fn() -> R>,
6591 group: TaskGroupId,
6592}
6593
6594impl<R> StagedCall<R> {
6595 pub(crate) fn new<T: 'static>(
6602 mut store: StoreContextMut<T>,
6603 prepared: PreparedCall<R>,
6604 ) -> Result<StagedCall<R>> {
6605 let PreparedCall {
6606 handle,
6607 thread,
6608 param_count,
6609 rx,
6610 ..
6611 } = prepared;
6612
6613 stage_call0(store.as_context_mut(), handle, thread, param_count)?;
6614
6615 Ok(StagedCall {
6616 store: store.0.id(),
6617 rx,
6618 _marker: PhantomData,
6619 group: store.0.concurrent_state_mut()?.get_mut(thread.task)?.group,
6620 })
6621 }
6622}
6623
6624impl<R> Future for StagedCall<R>
6625where
6626 R: 'static,
6627{
6628 type Output = Result<R>;
6629
6630 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
6631 check_ambient_store(self.store);
6632 Pin::new(&mut self.rx).poll(cx).map(|result| match result {
6633 Ok(r) => match r.downcast() {
6634 Ok(r) => Ok(*r),
6635 Err(_) => bail_bug!("wrong type of value produced"),
6636 },
6637 Err(oneshot::Canceled) => bail_bug!("channel erroneously dropped"),
6638 })
6639 }
6640}
6641
6642fn stage_call0<T: 'static>(
6645 store: StoreContextMut<T>,
6646 handle: Func,
6647 guest_thread: QualifiedThreadId,
6648 param_count: usize,
6649) -> Result<()> {
6650 let (_options, _, _ty, raw_options) = handle.abi_info(store.0);
6651 let is_concurrent = raw_options.async_;
6652 let callback = raw_options.callback;
6653 let instance = handle.instance();
6654 let callee = handle.lifted_core_func(store.0);
6655 let post_return = raw_options
6656 .post_return
6657 .map(|i| instance.id().get(store.0).runtime_post_return(i));
6658 let callback = callback.map(|i| {
6659 let instance = instance.id().get(store.0);
6660 SendSyncPtr::new(instance.runtime_callback(i))
6661 });
6662
6663 log::trace!("queueing call {guest_thread:?}");
6664
6665 unsafe {
6669 instance.stage_call(
6670 store,
6671 guest_thread,
6672 SendSyncPtr::new(callee),
6673 param_count,
6674 1,
6675 is_concurrent,
6676 callback,
6677 post_return.map(SendSyncPtr::new),
6678 true,
6679 )
6680 }
6681}
6682
6683#[cfg(all(test, feature = "cranelift", feature = "wat"))]
6684mod tests {
6685 use super::*;
6686 use crate::component::{Component, Linker};
6687 use crate::store::AsStoreOpaque;
6688 use crate::{Config, Engine};
6689
6690 fn host_subtask(
6691 state: HostTaskState,
6692 event: Option<Event>,
6693 ) -> Result<(Store<()>, Instance, TableId<HostTask>, u32)> {
6694 let mut config = Config::new();
6695 config.wasm_component_model_async(true);
6696 let engine = Engine::new(&config)?;
6697 let component = Component::new(&engine, "(component)")?;
6698 let mut store = Store::new(&engine, ());
6699 let instance = Linker::new(&engine).instantiate(&mut store, &component)?;
6700 let store_opaque = store.as_store_opaque();
6701 let concurrent_state = store_opaque.concurrent_state_mut()?;
6702 let group = concurrent_state.make_task_group()?;
6705 let task = concurrent_state.push(HostTask {
6706 common: WaitableCommon::default(),
6707 call_context: CallContext::default(),
6708 state,
6709 group,
6710 })?;
6711 let handle = store_opaque
6712 .instance_state(instance.runtime_instance(RuntimeComponentInstanceIndex::from_u32(0)))
6713 .handle_table()
6714 .subtask_insert_host(task.rep())?;
6715 let common = &mut store_opaque.concurrent_state_mut()?.get_mut(task)?.common;
6716 common.handle = Some(handle);
6717 common.event = event;
6718 Ok((store, instance, task, handle))
6719 }
6720
6721 #[test]
6722 fn host_subtask_drop_during_cancellation() -> Result<()> {
6723 for abort_completed in [false, true] {
6724 let (handle, future) = JoinHandle::run(future::pending::<()>());
6725 let mut future = pin!(future);
6726 let (mut store, instance, task, handle) =
6727 host_subtask(HostTaskState::CalleeRunning(handle), None)?;
6728 let store = store.as_store_opaque();
6729 let caller = RuntimeComponentInstanceIndex::from_u32(0);
6730 assert_eq!(
6731 instance.subtask_cancel(store, caller, true, handle)?,
6732 BLOCKED
6733 );
6734 if abort_completed {
6735 assert!(matches!(
6738 future
6739 .as_mut()
6740 .poll(&mut Context::from_waker(Waker::noop())),
6741 Poll::Ready(None),
6742 ));
6743 }
6744 for async_ in [false, true] {
6745 let err = instance
6746 .subtask_cancel(store, caller, async_, handle)
6747 .unwrap_err();
6748 assert_eq!(err.downcast::<Trap>()?, Trap::SubtaskCancelAfterTerminal);
6749 }
6750 let err = instance.subtask_drop(store, caller, handle).unwrap_err();
6751 assert_eq!(err.downcast::<Trap>()?, Trap::SubtaskDropNotResolved);
6752 assert!(store.concurrent_state_mut()?.get_mut(task).is_ok());
6753 }
6754 Ok(())
6755 }
6756
6757 #[test]
6758 fn host_subtask_cancel_after_completion() -> Result<()> {
6759 for async_ in [false, true] {
6760 let (mut store, instance, task, handle) = host_subtask(
6761 HostTaskState::CalleeDone { cancelled: false },
6762 Some(Event::Subtask {
6763 status: Status::Returned,
6764 }),
6765 )?;
6766 let store = store.as_store_opaque();
6767 let caller = RuntimeComponentInstanceIndex::from_u32(0);
6768 assert_eq!(
6769 instance.subtask_cancel(store, caller, async_, handle)?,
6770 Status::Returned as u32,
6771 );
6772 let err = instance
6773 .subtask_cancel(store, caller, async_, handle)
6774 .unwrap_err();
6775 assert_eq!(err.downcast::<Trap>()?, Trap::SubtaskCancelAfterTerminal);
6776 instance.subtask_drop(store, caller, handle)?;
6777 assert!(store.concurrent_state_mut()?.get_mut(task).is_err());
6778 }
6779 Ok(())
6780 }
6781
6782 #[test]
6783 fn host_subtask_drop_requires_terminal_event_delivery() -> Result<()> {
6784 for (cancelled, status) in [
6785 (false, Status::Returned),
6786 (true, Status::Returned),
6787 (true, Status::ReturnCancelled),
6788 ] {
6789 for delivered in [false, true] {
6790 let event = if delivered {
6791 None
6792 } else {
6793 Some(Event::Subtask { status })
6794 };
6795 let (mut store, instance, task, handle) =
6796 host_subtask(HostTaskState::CalleeDone { cancelled }, event)?;
6797 let store = store.as_store_opaque();
6798 let result = instance.subtask_drop(
6799 store,
6800 RuntimeComponentInstanceIndex::from_u32(0),
6801 handle,
6802 );
6803 if delivered {
6804 result?;
6805 assert!(store.concurrent_state_mut()?.get_mut(task).is_err());
6806 } else {
6807 let err = result.unwrap_err();
6808 assert_eq!(err.downcast::<Trap>()?, Trap::SubtaskDropNotResolved);
6809 assert!(store.concurrent_state_mut()?.get_mut(task).is_ok());
6810 }
6811 }
6812 }
6813 Ok(())
6814 }
6815}