1use super::table::{TableDebug, TableId};
2use super::{Event, GlobalErrorContextRefCount, Waitable, WaitableCommon};
3use crate::component::concurrent::{ConcurrentState, QualifiedThreadId, WorkItem, tls};
4use crate::component::func::{self, LiftContext, LowerContext};
5use crate::component::matching::InstanceType;
6use crate::component::types;
7use crate::component::values::ErrorContextAny;
8use crate::component::{
9 AsAccessor, ComponentInstanceId, ComponentType, FutureAny, Instance, Lift, Lower, StreamAny,
10 Val, WasmList,
11};
12use crate::prelude::*;
13use crate::store::{StoreOpaque, StoreToken};
14use crate::try_mutex::{TryMutex, TryMutexGuard};
15use crate::vm::component::{ComponentInstance, HandleTable, TransmitLocalState};
16use crate::vm::{AlwaysMut, VMStore};
17use crate::{AsContext, AsContextMut, StoreContextMut, ValRaw};
18use crate::{
19 Error, Result, Trap, bail, bail_bug, ensure,
20 error::{Context as _, format_err},
21};
22use alloc::sync::Arc;
23use buffers::{Extender, SliceBuffer, UntypedWriteBuffer};
24use core::any::{Any, TypeId};
25use core::fmt;
26use core::future;
27use core::iter;
28use core::marker::PhantomData;
29use core::mem::{self, ManuallyDrop, MaybeUninit};
30use core::ops::{Deref, DerefMut};
31use core::pin::Pin;
32use core::task::{Context, Poll, Waker, ready};
33use futures::channel::oneshot;
34use futures::{FutureExt as _, stream};
35use wasmtime_environ::component::{
36 CanonicalAbiInfo, ComponentTypes, InterfaceType, OptionsIndex, RuntimeComponentInstanceIndex,
37 TypeComponentGlobalErrorContextTableIndex, TypeComponentLocalErrorContextTableIndex,
38 TypeFutureTableIndex, TypeStreamTableIndex,
39};
40
41pub use buffers::{ReadBuffer, VecBuffer, WriteBuffer};
42
43mod buffers;
44
45#[derive(Copy, Clone, Debug)]
48pub enum TransmitKind {
49 Stream,
50 Future,
51}
52
53#[derive(Copy, Clone, Debug, PartialEq)]
55pub enum ReturnCode {
56 Blocked,
57 Completed(ItemCount),
58 Dropped(ItemCount),
59 Cancelled(ItemCount),
60}
61
62impl ReturnCode {
63 pub fn encode(&self) -> u32 {
68 const BLOCKED: u32 = 0xffff_ffff;
69 const COMPLETED: u32 = 0x0;
70 const DROPPED: u32 = 0x1;
71 const CANCELLED: u32 = 0x2;
72 match self {
73 ReturnCode::Blocked => BLOCKED,
74 ReturnCode::Completed(n) => (n.as_u32() << 4) | COMPLETED,
75 ReturnCode::Dropped(n) => (n.as_u32() << 4) | DROPPED,
76 ReturnCode::Cancelled(n) => (n.as_u32() << 4) | CANCELLED,
77 }
78 }
79
80 fn completed(kind: TransmitKind, count: ItemCount) -> Self {
83 Self::Completed(if let TransmitKind::Future = kind {
84 ItemCount::ZERO
85 } else {
86 count
87 })
88 }
89}
90
91#[derive(Copy, Clone, PartialEq, Eq, PartialOrd, Ord)]
98#[repr(transparent)]
99pub struct ItemCount {
100 raw: u32,
101}
102
103impl ItemCount {
104 const MAX: u32 = 1 << 28;
105 const ZERO: ItemCount = ItemCount { raw: 0 };
106
107 fn new(count: u32) -> Result<Self, Trap> {
110 if count < Self::MAX {
111 Ok(Self { raw: count })
112 } else {
113 Err(Trap::StreamOpTooBig)
114 }
115 }
116
117 fn new_usize(count: usize) -> Result<Self, Trap> {
119 let count = u32::try_from(count).map_err(|_| Trap::StreamOpTooBig)?;
120 Self::new(count)
121 }
122
123 fn as_u32(&self) -> u32 {
124 self.raw
125 }
126
127 fn as_usize(&self) -> usize {
128 usize::try_from(self.raw).unwrap()
129 }
130
131 fn inc(&mut self, amt: usize) -> Result<(), Trap> {
134 let amt = u32::try_from(amt).map_err(|_| Trap::StreamOpTooBig)?;
135 let new_raw = self.raw.checked_add(amt).ok_or(Trap::StreamOpTooBig)?;
136 if new_raw < Self::MAX {
137 self.raw = new_raw;
138 Ok(())
139 } else {
140 Err(Trap::StreamOpTooBig)
141 }
142 }
143
144 fn add(&self, other: ItemCount) -> Result<ItemCount> {
150 match self.raw.checked_add(other.raw) {
151 Some(raw) => Ok(ItemCount::new(raw)?),
152 None => bail_bug!("overflow in `ItemCount::add`"),
153 }
154 }
155
156 fn sub(&self, other: ItemCount) -> Result<ItemCount> {
161 match self.raw.checked_sub(other.raw) {
162 Some(raw) => Ok(ItemCount { raw }),
163 None => bail_bug!("underflow in `ItemCount::sub`"),
164 }
165 }
166}
167
168impl fmt::Display for ItemCount {
169 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
170 self.raw.fmt(f)
171 }
172}
173
174impl fmt::Debug for ItemCount {
175 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
176 self.raw.fmt(f)
177 }
178}
179
180impl PartialEq<u32> for ItemCount {
181 fn eq(&self, other: &u32) -> bool {
182 self.raw == *other
183 }
184}
185
186impl PartialOrd<u32> for ItemCount {
187 fn partial_cmp(&self, other: &u32) -> Option<core::cmp::Ordering> {
188 self.raw.partial_cmp(other)
189 }
190}
191
192#[derive(Copy, Clone, Debug)]
197pub enum TransmitIndex {
198 Stream(TypeStreamTableIndex),
199 Future(TypeFutureTableIndex),
200}
201
202impl TransmitIndex {
203 pub fn kind(&self) -> TransmitKind {
204 match self {
205 TransmitIndex::Stream(_) => TransmitKind::Stream,
206 TransmitIndex::Future(_) => TransmitKind::Future,
207 }
208 }
209
210 fn payload<'a>(&self, types: &'a ComponentTypes) -> Option<&'a InterfaceType> {
213 match self {
214 TransmitIndex::Stream(i) => {
215 let ty = types[*i].ty;
216 types[ty].payload.as_ref()
217 }
218 TransmitIndex::Future(i) => {
219 let ty = types[*i].ty;
220 types[ty].payload.as_ref()
221 }
222 }
223 }
224}
225
226fn get_mut_by_index_from(
229 handle_table: &mut HandleTable,
230 ty: TransmitIndex,
231 index: u32,
232) -> Result<(u32, &mut TransmitLocalState)> {
233 match ty {
234 TransmitIndex::Stream(ty) => handle_table.stream_rep(ty, index),
235 TransmitIndex::Future(ty) => handle_table.future_rep(ty, index),
236 }
237}
238
239fn lower<T: func::Lower + Send + 'static, B: WriteBuffer<T>, U: 'static>(
240 mut store: StoreContextMut<U>,
241 instance: Instance,
242 caller_thread: QualifiedThreadId,
243 options: OptionsIndex,
244 ty: TransmitIndex,
245 address: usize,
246 count: usize,
247 buffer: &mut B,
248) -> Result<()> {
249 let count = buffer.remaining().len().min(count);
250
251 let (lower, old_thread) = if T::MAY_REQUIRE_REALLOC {
255 let old_thread = store.0.set_thread(caller_thread)?;
256 (
257 &mut LowerContext::new(store.as_context_mut(), options, instance),
258 Some(old_thread),
259 )
260 } else {
261 (
262 &mut LowerContext::new_without_realloc(store.as_context_mut(), options, instance),
263 None,
264 )
265 };
266
267 if address % usize::try_from(T::ALIGN32)? != 0 {
268 bail!("read pointer not aligned");
269 }
270 lower
271 .as_slice_mut()
272 .get_mut(address..)
273 .and_then(|b| b.get_mut(..T::SIZE32 * count))
274 .ok_or_else(|| crate::format_err!("read pointer out of bounds of memory"))?;
275
276 if let Some(ty) = ty.payload(lower.types) {
277 T::linear_store_list_to_memory(lower, *ty, address, &buffer.remaining()[..count])?;
278 }
279
280 if let Some(old_thread) = old_thread {
281 store.0.set_thread(old_thread)?;
282 }
283
284 buffer.skip(count);
285
286 Ok(())
287}
288
289fn lift<T: func::Lift + Send + 'static, B: ReadBuffer<T>>(
290 lift: &mut LiftContext<'_>,
291 ty: Option<InterfaceType>,
292 buffer: &mut B,
293 address: usize,
294 count: usize,
295) -> Result<()> {
296 let count = count.min(buffer.remaining_capacity());
297 if T::IS_RUST_UNIT_TYPE {
298 buffer.extend(
302 iter::repeat_with(|| unsafe { MaybeUninit::uninit().assume_init() }).take(count),
303 )
304 } else {
305 let ty = match ty {
306 Some(ty) => ty,
307 None => bail_bug!("type required for non-unit lift"),
308 };
309 if address % usize::try_from(T::ALIGN32)? != 0 {
310 bail!("write pointer not aligned");
311 }
312 lift.memory()
313 .get(address..)
314 .and_then(|b| b.get(..T::SIZE32 * count))
315 .ok_or_else(|| crate::format_err!("write pointer out of bounds of memory"))?;
316
317 let list = &WasmList::new(address, count, lift, ty)?;
318 T::linear_lift_into_from_memory(lift, list, &mut Extender(buffer))?
319 }
320 Ok(())
321}
322
323#[derive(Debug, PartialEq, Eq, PartialOrd)]
325pub(super) struct ErrorContextState {
326 pub(crate) debug_msg: String,
328}
329
330#[derive(Debug, Clone, Copy, PartialEq, Eq)]
333pub(super) struct FlatAbi {
334 pub(super) size: u32,
335 pub(super) align: u32,
336}
337
338struct HostBuffer<'a> {
339 dst: &'a mut Vec<u8>,
340 marked_written: &'a mut usize,
341}
342
343impl HostBuffer<'_> {
344 fn reborrow(&mut self) -> HostBuffer<'_> {
345 HostBuffer {
346 dst: &mut *self.dst,
347 marked_written: &mut *self.marked_written,
348 }
349 }
350}
351
352pub struct Destination<'a, T, B> {
354 id: TableId<TransmitState>,
355 buffer: &'a mut B,
356 host_buffer: Option<HostBuffer<'a>>,
357 _phantom: PhantomData<fn() -> T>,
358}
359
360impl<'a, T, B> Destination<'a, T, B> {
361 pub fn reborrow(&mut self) -> Destination<'_, T, B> {
363 Destination {
364 id: self.id,
365 buffer: &mut *self.buffer,
366 host_buffer: self.host_buffer.as_mut().map(|b| b.reborrow()),
367 _phantom: PhantomData,
368 }
369 }
370
371 pub fn take_buffer(&mut self) -> B
377 where
378 B: Default,
379 {
380 mem::take(self.buffer)
381 }
382
383 pub fn set_buffer(&mut self, buffer: B) {
393 *self.buffer = buffer;
394 }
395
396 pub fn remaining(&self, mut store: impl AsContextMut) -> Option<usize> {
413 self.remaining_(store.as_context_mut().0).unwrap()
417 }
418
419 fn remaining_(&self, store: &mut StoreOpaque) -> Result<Option<usize>> {
420 let transmit = store.concurrent_state_mut()?.get_mut(self.id)?;
421
422 if let &ReadState::GuestReady { count, .. } = &transmit.read {
423 let &WriteState::HostReady { guest_offset, .. } = &transmit.write else {
424 bail_bug!("expected WriteState::HostReady")
425 };
426
427 Ok(Some(count.as_usize() - guest_offset.as_usize()))
428 } else {
429 Ok(None)
430 }
431 }
432}
433
434impl<'a, B> Destination<'a, u8, B> {
435 pub fn as_direct<D>(
446 mut self,
447 store: StoreContextMut<'a, D>,
448 capacity: usize,
449 ) -> DirectDestination<'a, D> {
450 if let Some(buffer) = &mut self.host_buffer {
451 *buffer.marked_written = 0;
452 buffer.dst.resize(capacity, 0);
453 }
454
455 DirectDestination {
456 id: self.id,
457 host_buffer: self.host_buffer,
458 store,
459 }
460 }
461}
462
463pub struct DirectDestination<'a, D: 'static> {
466 id: TableId<TransmitState>,
467 host_buffer: Option<HostBuffer<'a>>,
468 store: StoreContextMut<'a, D>,
469}
470
471#[cfg(feature = "std")]
472impl<D: 'static> std::io::Write for DirectDestination<'_, D> {
473 fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
474 let rem = self.remaining();
475 let n = rem.len().min(buf.len());
476 rem[..n].copy_from_slice(&buf[..n]);
477 self.mark_written(n);
478 Ok(n)
479 }
480
481 fn flush(&mut self) -> std::io::Result<()> {
482 Ok(())
483 }
484}
485
486impl<D: 'static> DirectDestination<'_, D> {
487 pub fn remaining(&mut self) -> &mut [u8] {
489 self.remaining_().unwrap()
493 }
494
495 fn remaining_(&mut self) -> Result<&mut [u8]> {
496 if let Some(buffer) = self.host_buffer.as_mut() {
497 return Ok(buffer.dst);
498 }
499 let transmit = self
500 .store
501 .as_context_mut()
502 .0
503 .concurrent_state_mut()?
504 .get_mut(self.id)?;
505
506 let &ReadState::GuestReady {
507 address,
508 count,
509 options,
510 instance,
511 ..
512 } = &transmit.read
513 else {
514 bail_bug!("expected ReadState::GuestReady")
515 };
516
517 let &WriteState::HostReady { guest_offset, .. } = &transmit.write else {
518 bail_bug!("expected WriteState::HostReady")
519 };
520
521 let memory = instance
522 .options_memory_mut(self.store.0, options)
523 .get_mut((address + guest_offset.as_usize())..)
524 .and_then(|b| b.get_mut(..(count.as_usize() - guest_offset.as_usize())));
525 match memory {
526 Some(memory) => Ok(memory),
527 None => bail_bug!("guest buffer unexpectedly out of bounds"),
528 }
529 }
530
531 pub fn mark_written(&mut self, count: usize) {
538 self.mark_written_(count).unwrap()
542 }
543
544 fn mark_written_(&mut self, count: usize) -> Result<()> {
545 if let Some(buffer) = self.host_buffer.as_mut() {
546 *buffer.marked_written = buffer.marked_written.checked_add(count).unwrap();
549 } else {
550 let transmit = self
551 .store
552 .as_context_mut()
553 .0
554 .concurrent_state_mut()?
555 .get_mut(self.id)?;
556
557 let ReadState::GuestReady {
558 count: read_count, ..
559 } = &transmit.read
560 else {
561 bail_bug!("expected ReadState::GuestReady")
562 };
563
564 let WriteState::HostReady { guest_offset, .. } = &mut transmit.write else {
565 bail_bug!("expected WriteState::HostReady");
566 };
567
568 if guest_offset.as_usize() + count > read_count.as_usize() {
569 panic!(
572 "write count ({count}) must be less than or equal to read count ({read_count})"
573 )
574 } else {
575 guest_offset.inc(count)?;
576 }
577 }
578 Ok(())
579 }
580}
581
582#[derive(Copy, Clone, Debug)]
584pub enum StreamResult {
585 Completed,
588 Cancelled,
593 Dropped,
596}
597
598pub trait StreamProducer<D>: Send + 'static {
600 type Item;
602
603 type Buffer: WriteBuffer<Self::Item> + Default;
605
606 fn poll_produce<'a>(
742 self: Pin<&mut Self>,
743 cx: &mut Context<'_>,
744 store: StoreContextMut<'a, D>,
745 destination: Destination<'a, Self::Item, Self::Buffer>,
746 finish: bool,
747 ) -> Poll<Result<StreamResult>>;
748
749 fn try_into(me: Pin<Box<Self>>, _ty: TypeId) -> Result<Box<dyn Any>, Pin<Box<Self>>> {
755 Err(me)
756 }
757}
758
759impl<T, D> StreamProducer<D> for iter::Empty<T>
760where
761 T: Send + Sync + 'static,
762{
763 type Item = T;
764 type Buffer = Option<Self::Item>;
765
766 fn poll_produce<'a>(
767 self: Pin<&mut Self>,
768 _: &mut Context<'_>,
769 _: StoreContextMut<'a, D>,
770 _: Destination<'a, Self::Item, Self::Buffer>,
771 _: bool,
772 ) -> Poll<Result<StreamResult>> {
773 Poll::Ready(Ok(StreamResult::Dropped))
774 }
775}
776
777impl<T, D> StreamProducer<D> for stream::Empty<T>
778where
779 T: Send + Sync + 'static,
780{
781 type Item = T;
782 type Buffer = Option<Self::Item>;
783
784 fn poll_produce<'a>(
785 self: Pin<&mut Self>,
786 _: &mut Context<'_>,
787 _: StoreContextMut<'a, D>,
788 _: Destination<'a, Self::Item, Self::Buffer>,
789 _: bool,
790 ) -> Poll<Result<StreamResult>> {
791 Poll::Ready(Ok(StreamResult::Dropped))
792 }
793}
794
795impl<T, D> StreamProducer<D> for Vec<T>
796where
797 T: Unpin + Send + Sync + 'static,
798{
799 type Item = T;
800 type Buffer = VecBuffer<T>;
801
802 fn poll_produce<'a>(
803 self: Pin<&mut Self>,
804 _: &mut Context<'_>,
805 _: StoreContextMut<'a, D>,
806 mut dst: Destination<'a, Self::Item, Self::Buffer>,
807 _: bool,
808 ) -> Poll<Result<StreamResult>> {
809 dst.set_buffer(mem::take(self.get_mut()).into());
810 Poll::Ready(Ok(StreamResult::Dropped))
811 }
812}
813
814impl<T, D> StreamProducer<D> for Box<[T]>
815where
816 T: Unpin + Send + Sync + 'static,
817{
818 type Item = T;
819 type Buffer = VecBuffer<T>;
820
821 fn poll_produce<'a>(
822 self: Pin<&mut Self>,
823 _: &mut Context<'_>,
824 _: StoreContextMut<'a, D>,
825 mut dst: Destination<'a, Self::Item, Self::Buffer>,
826 _: bool,
827 ) -> Poll<Result<StreamResult>> {
828 dst.set_buffer(mem::take(self.get_mut()).into_vec().into());
829 Poll::Ready(Ok(StreamResult::Dropped))
830 }
831}
832
833#[cfg(feature = "component-model-bytes")]
834impl<D> StreamProducer<D> for bytes::Bytes {
835 type Item = u8;
836 type Buffer = Self;
837
838 fn poll_produce<'a>(
839 self: Pin<&mut Self>,
840 _: &mut Context<'_>,
841 _store: StoreContextMut<'a, D>,
842 mut dst: Destination<'a, Self::Item, Self::Buffer>,
843 _: bool,
844 ) -> Poll<Result<StreamResult>> {
845 dst.set_buffer(mem::take(self.get_mut()));
846 Poll::Ready(Ok(StreamResult::Dropped))
847 }
848}
849
850#[cfg(feature = "component-model-bytes")]
851impl<D> StreamProducer<D> for bytes::BytesMut {
852 type Item = u8;
853 type Buffer = Self;
854
855 fn poll_produce<'a>(
856 self: Pin<&mut Self>,
857 _: &mut Context<'_>,
858 _store: StoreContextMut<'a, D>,
859 mut dst: Destination<'a, Self::Item, Self::Buffer>,
860 _: bool,
861 ) -> Poll<Result<StreamResult>> {
862 dst.set_buffer(mem::take(self.get_mut()));
863 Poll::Ready(Ok(StreamResult::Dropped))
864 }
865}
866
867pub struct Source<'a, T> {
869 id: TableId<TransmitState>,
870 host_buffer: Option<&'a mut dyn WriteBuffer<T>>,
871}
872
873impl<'a, T> Source<'a, T> {
874 pub fn reborrow(&mut self) -> Source<'_, T> {
876 Source {
877 id: self.id,
878 host_buffer: self.host_buffer.as_deref_mut(),
879 }
880 }
881
882 pub fn read<B, S: AsContextMut>(&mut self, mut store: S, buffer: &mut B) -> Result<()>
884 where
885 T: func::Lift + 'static,
886 B: ReadBuffer<T>,
887 {
888 if let Some(input) = &mut self.host_buffer {
889 let count = input.remaining().len().min(buffer.remaining_capacity());
890 buffer.move_from(*input, count);
891 } else {
892 let store = store.as_context_mut();
893 let transmit = store.0.concurrent_state_mut()?.get_mut(self.id)?;
894
895 let &ReadState::HostReady { guest_offset, .. } = &transmit.read else {
896 bail_bug!("expected ReadState::HostReady");
897 };
898
899 let &WriteState::GuestReady {
900 ty,
901 address,
902 count,
903 options,
904 instance,
905 ..
906 } = &transmit.write
907 else {
908 bail_bug!("expected WriteState::GuestReady");
909 };
910
911 let cx = &mut LiftContext::new(store.0.store_opaque_mut(), options, instance)?;
912 let ty = ty.payload(cx.types);
913 let old_remaining = buffer.remaining_capacity();
914 lift::<T, B>(
915 cx,
916 ty.copied(),
917 buffer,
918 address + (T::SIZE32 * guest_offset.as_usize()),
919 count.as_usize() - guest_offset.as_usize(),
920 )?;
921
922 let transmit = store.0.concurrent_state_mut()?.get_mut(self.id)?;
923
924 let ReadState::HostReady { guest_offset, .. } = &mut transmit.read else {
925 bail_bug!("expected ReadState::HostReady");
926 };
927
928 guest_offset.inc(old_remaining - buffer.remaining_capacity())?;
929 }
930
931 Ok(())
932 }
933
934 pub fn remaining(&self, mut store: impl AsContextMut) -> usize
937 where
938 T: 'static,
939 {
940 self.remaining_(store.as_context_mut().0).unwrap()
944 }
945
946 fn remaining_(&self, store: &mut StoreOpaque) -> Result<usize>
947 where
948 T: 'static,
949 {
950 let transmit = store.concurrent_state_mut()?.get_mut(self.id)?;
951
952 if let &WriteState::GuestReady { count, .. } = &transmit.write {
953 let &ReadState::HostReady { guest_offset, .. } = &transmit.read else {
954 bail_bug!("expected ReadState::HostReady")
955 };
956
957 Ok(count.as_usize() - guest_offset.as_usize())
958 } else if let Some(host_buffer) = &self.host_buffer {
959 Ok(host_buffer.remaining().len())
960 } else {
961 bail_bug!("expected either WriteState::GuestReady or host buffer")
962 }
963 }
964}
965
966impl<'a> Source<'a, u8> {
967 pub fn as_direct<D>(self, store: StoreContextMut<'a, D>) -> DirectSource<'a, D> {
969 DirectSource {
970 id: self.id,
971 host_buffer: self.host_buffer,
972 store,
973 }
974 }
975}
976
977pub struct DirectSource<'a, D: 'static> {
980 id: TableId<TransmitState>,
981 host_buffer: Option<&'a mut dyn WriteBuffer<u8>>,
982 store: StoreContextMut<'a, D>,
983}
984
985#[cfg(feature = "std")]
986impl<D: 'static> std::io::Read for DirectSource<'_, D> {
987 fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
988 let rem = self.remaining();
989 let n = rem.len().min(buf.len());
990 buf[..n].copy_from_slice(&rem[..n]);
991 self.mark_read(n);
992 Ok(n)
993 }
994}
995
996impl<D: 'static> DirectSource<'_, D> {
997 pub fn remaining(&mut self) -> &[u8] {
999 self.remaining_().unwrap()
1003 }
1004
1005 fn remaining_(&mut self) -> Result<&[u8]> {
1006 if let Some(buffer) = self.host_buffer.as_deref_mut() {
1007 return Ok(buffer.remaining());
1008 }
1009 let transmit = self
1010 .store
1011 .as_context_mut()
1012 .0
1013 .concurrent_state_mut()?
1014 .get_mut(self.id)?;
1015
1016 let &WriteState::GuestReady {
1017 address,
1018 count,
1019 options,
1020 instance,
1021 ..
1022 } = &transmit.write
1023 else {
1024 bail_bug!("expected WriteState::GuestReady")
1025 };
1026
1027 let &ReadState::HostReady { guest_offset, .. } = &transmit.read else {
1028 bail_bug!("expected ReadState::HostReady")
1029 };
1030
1031 let memory = instance
1032 .options_memory(self.store.0, options)
1033 .get((address + guest_offset.as_usize())..)
1034 .and_then(|b| b.get(..(count.as_usize() - guest_offset.as_usize())));
1035 match memory {
1036 Some(memory) => Ok(memory),
1037 None => bail_bug!("guest buffer unexpectedly out of bounds"),
1038 }
1039 }
1040
1041 pub fn mark_read(&mut self, count: usize) {
1048 self.mark_read_(count).unwrap()
1052 }
1053
1054 fn mark_read_(&mut self, count: usize) -> Result<()> {
1055 if let Some(buffer) = self.host_buffer.as_deref_mut() {
1056 buffer.skip(count);
1057 return Ok(());
1058 }
1059
1060 let transmit = self
1061 .store
1062 .as_context_mut()
1063 .0
1064 .concurrent_state_mut()?
1065 .get_mut(self.id)?;
1066
1067 let WriteState::GuestReady {
1068 count: write_count, ..
1069 } = &transmit.write
1070 else {
1071 bail_bug!("expected WriteState::GuestReady");
1072 };
1073
1074 let ReadState::HostReady { guest_offset, .. } = &mut transmit.read else {
1075 bail_bug!("expected ReadState::HostReady");
1076 };
1077
1078 if guest_offset.as_usize() + count > write_count.as_usize() {
1079 panic!("read count ({count}) must be less than or equal to write count ({write_count})")
1081 } else {
1082 guest_offset.inc(count)?;
1083 }
1084 Ok(())
1085 }
1086}
1087
1088pub trait StreamConsumer<D>: Send + 'static {
1090 type Item;
1092
1093 fn poll_consume(
1176 self: Pin<&mut Self>,
1177 cx: &mut Context<'_>,
1178 store: StoreContextMut<D>,
1179 source: Source<'_, Self::Item>,
1180 finish: bool,
1181 ) -> Poll<Result<StreamResult>>;
1182}
1183
1184pub trait FutureProducer<D>: Send + 'static {
1186 type Item;
1188
1189 fn poll_produce(
1199 self: Pin<&mut Self>,
1200 cx: &mut Context<'_>,
1201 store: StoreContextMut<D>,
1202 finish: bool,
1203 ) -> Poll<Result<Option<Self::Item>>>;
1204}
1205
1206impl<T, E, D, Fut> FutureProducer<D> for Fut
1207where
1208 E: Into<Error>,
1209 Fut: Future<Output = Result<T, E>> + ?Sized + Send + 'static,
1210{
1211 type Item = T;
1212
1213 fn poll_produce<'a>(
1214 self: Pin<&mut Self>,
1215 cx: &mut Context<'_>,
1216 _: StoreContextMut<'a, D>,
1217 finish: bool,
1218 ) -> Poll<Result<Option<T>>> {
1219 match self.poll(cx) {
1220 Poll::Ready(Ok(v)) => Poll::Ready(Ok(Some(v))),
1221 Poll::Ready(Err(err)) => Poll::Ready(Err(err.into())),
1222 Poll::Pending if finish => Poll::Ready(Ok(None)),
1223 Poll::Pending => Poll::Pending,
1224 }
1225 }
1226}
1227
1228pub trait FutureConsumer<D>: Send + 'static {
1230 type Item;
1232
1233 fn poll_consume(
1245 self: Pin<&mut Self>,
1246 cx: &mut Context<'_>,
1247 store: StoreContextMut<D>,
1248 source: Source<'_, Self::Item>,
1249 finish: bool,
1250 ) -> Poll<Result<()>>;
1251}
1252
1253pub struct FutureReader<T> {
1260 id: TableId<TransmitHandle>,
1261 _phantom: PhantomData<T>,
1262}
1263
1264impl<T> FutureReader<T> {
1265 pub fn new<S: AsContextMut>(
1274 mut store: S,
1275 producer: impl FutureProducer<S::Data, Item = T>,
1276 ) -> Result<Self>
1277 where
1278 T: func::Lower + func::Lift + Send + Sync + 'static,
1279 {
1280 ensure!(
1281 store.as_context().0.concurrency_support(),
1282 "concurrency support is not enabled"
1283 );
1284
1285 struct Producer<P>(P);
1286
1287 impl<D, T: func::Lower + 'static, P: FutureProducer<D, Item = T>> StreamProducer<D>
1288 for Producer<P>
1289 {
1290 type Item = P::Item;
1291 type Buffer = Option<P::Item>;
1292
1293 fn poll_produce<'a>(
1294 self: Pin<&mut Self>,
1295 cx: &mut Context<'_>,
1296 store: StoreContextMut<D>,
1297 mut destination: Destination<'a, Self::Item, Self::Buffer>,
1298 finish: bool,
1299 ) -> Poll<Result<StreamResult>> {
1300 let producer = unsafe { self.map_unchecked_mut(|v| &mut v.0) };
1303
1304 Poll::Ready(Ok(
1305 if let Some(value) = ready!(producer.poll_produce(cx, store, finish))? {
1306 destination.set_buffer(Some(value));
1307
1308 StreamResult::Completed
1315 } else {
1316 StreamResult::Cancelled
1317 },
1318 ))
1319 }
1320 }
1321
1322 Ok(Self::new_(
1323 store
1324 .as_context_mut()
1325 .new_transmit(TransmitKind::Future, Producer(producer))?,
1326 ))
1327 }
1328
1329 pub(super) fn new_(id: TableId<TransmitHandle>) -> Self {
1330 Self {
1331 id,
1332 _phantom: PhantomData,
1333 }
1334 }
1335
1336 pub(super) fn id(&self) -> TableId<TransmitHandle> {
1337 self.id
1338 }
1339
1340 pub fn pipe<S: AsContextMut>(
1350 self,
1351 mut store: S,
1352 consumer: impl FutureConsumer<S::Data, Item = T> + Unpin,
1353 ) -> Result<()>
1354 where
1355 T: func::Lift + 'static,
1356 {
1357 struct Consumer<C>(C);
1358
1359 impl<D: 'static, T: func::Lift + 'static, C: FutureConsumer<D, Item = T>> StreamConsumer<D>
1360 for Consumer<C>
1361 {
1362 type Item = T;
1363
1364 fn poll_consume(
1365 self: Pin<&mut Self>,
1366 cx: &mut Context<'_>,
1367 mut store: StoreContextMut<D>,
1368 mut source: Source<Self::Item>,
1369 finish: bool,
1370 ) -> Poll<Result<StreamResult>> {
1371 let consumer = unsafe { self.map_unchecked_mut(|v| &mut v.0) };
1374
1375 ready!(consumer.poll_consume(
1376 cx,
1377 store.as_context_mut(),
1378 source.reborrow(),
1379 finish
1380 ))?;
1381
1382 Poll::Ready(Ok(if source.remaining(store) == 0 {
1383 StreamResult::Completed
1389 } else {
1390 StreamResult::Cancelled
1391 }))
1392 }
1393 }
1394
1395 store
1396 .as_context_mut()
1397 .set_consumer(self.id, TransmitKind::Future, Consumer(consumer))
1398 }
1399
1400 fn lift_from_index(cx: &mut LiftContext<'_>, ty: InterfaceType, index: u32) -> Result<Self> {
1402 let id = lift_index_to_future(cx, ty, index)?;
1403 Ok(Self::new_(id))
1404 }
1405
1406 pub fn close(&mut self, mut store: impl AsContextMut) -> Result<()> {
1424 future_close(store.as_context_mut().0, &mut self.id)
1425 }
1426
1427 pub fn close_with(&mut self, accessor: impl AsAccessor) -> Result<()> {
1429 accessor.as_accessor().with(|access| self.close(access))
1430 }
1431
1432 pub fn guard<A>(self, accessor: A) -> GuardedFutureReader<T, A>
1438 where
1439 A: AsAccessor,
1440 {
1441 GuardedFutureReader::new(accessor, self)
1442 }
1443
1444 pub fn try_into_future_any(self, store: impl AsContextMut) -> Result<FutureAny>
1451 where
1452 T: ComponentType + 'static,
1453 {
1454 FutureAny::try_from_future_reader(store, self)
1455 }
1456
1457 pub fn try_from_future_any(future: FutureAny) -> Result<Self>
1464 where
1465 T: ComponentType + 'static,
1466 {
1467 future.try_into_future_reader()
1468 }
1469}
1470
1471impl<T> fmt::Debug for FutureReader<T> {
1472 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1473 f.debug_struct("FutureReader")
1474 .field("id", &self.id)
1475 .finish()
1476 }
1477}
1478
1479pub(super) fn future_close(
1480 store: &mut StoreOpaque,
1481 id: &mut TableId<TransmitHandle>,
1482) -> Result<()> {
1483 let id = mem::replace(id, TableId::new(u32::MAX));
1484 store.host_drop_reader(id, TransmitKind::Future)
1485}
1486
1487pub(super) fn lift_index_to_future(
1489 cx: &mut LiftContext<'_>,
1490 ty: InterfaceType,
1491 index: u32,
1492) -> Result<TableId<TransmitHandle>> {
1493 match ty {
1494 InterfaceType::Future(src) => {
1495 let (state, instance) = cx.concurrent_state_and_instance_mut();
1496 lift_index_to_transmit(instance, state, TransmitIndex::Future(src), index)
1497 }
1498 _ => func::bad_type_info(),
1499 }
1500}
1501
1502pub(super) fn lower_future_to_index<U>(
1504 id: TableId<TransmitHandle>,
1505 cx: &mut LowerContext<'_, U>,
1506 ty: InterfaceType,
1507) -> Result<u32> {
1508 match ty {
1509 InterfaceType::Future(dst) => {
1510 cx.instance_handle()
1511 .lower_transmit_to_index(cx.store.0, TransmitIndex::Future(dst), id)
1512 }
1513 _ => func::bad_type_info(),
1514 }
1515}
1516
1517unsafe impl<T: ComponentType> ComponentType for FutureReader<T> {
1520 const ABI: CanonicalAbiInfo = CanonicalAbiInfo::SCALAR4;
1521
1522 type Lower = <u32 as func::ComponentType>::Lower;
1523
1524 fn typecheck(ty: &InterfaceType, types: &InstanceType<'_>) -> Result<()> {
1525 match ty {
1526 InterfaceType::Future(ty) => {
1527 let ty = types.types[*ty].ty;
1528 types::typecheck_payload::<T>(types.types[ty].payload.as_ref(), types)
1529 }
1530 other => bail!("expected `future`, found `{}`", func::desc(other)),
1531 }
1532 }
1533}
1534
1535unsafe impl<T: ComponentType> func::Lower for FutureReader<T> {
1537 fn linear_lower_to_flat<U>(
1538 &self,
1539 cx: &mut LowerContext<'_, U>,
1540 ty: InterfaceType,
1541 dst: &mut MaybeUninit<Self::Lower>,
1542 ) -> Result<()> {
1543 lower_future_to_index(self.id, cx, ty)?.linear_lower_to_flat(cx, InterfaceType::U32, dst)
1544 }
1545
1546 fn linear_lower_to_memory<U>(
1547 &self,
1548 cx: &mut LowerContext<'_, U>,
1549 ty: InterfaceType,
1550 offset: usize,
1551 ) -> Result<()> {
1552 lower_future_to_index(self.id, cx, ty)?.linear_lower_to_memory(
1553 cx,
1554 InterfaceType::U32,
1555 offset,
1556 )
1557 }
1558}
1559
1560unsafe impl<T: ComponentType> func::Lift for FutureReader<T> {
1562 fn linear_lift_from_flat(
1563 cx: &mut LiftContext<'_>,
1564 ty: InterfaceType,
1565 src: &Self::Lower,
1566 ) -> Result<Self> {
1567 let index = u32::linear_lift_from_flat(cx, InterfaceType::U32, src)?;
1568 Self::lift_from_index(cx, ty, index)
1569 }
1570
1571 fn linear_lift_from_memory(
1572 cx: &mut LiftContext<'_>,
1573 ty: InterfaceType,
1574 bytes: &[u8],
1575 ) -> Result<Self> {
1576 let index = u32::linear_lift_from_memory(cx, InterfaceType::U32, bytes)?;
1577 Self::lift_from_index(cx, ty, index)
1578 }
1579}
1580
1581pub struct GuardedFutureReader<T, A>
1589where
1590 A: AsAccessor,
1591{
1592 reader: Option<FutureReader<T>>,
1596 accessor: A,
1597}
1598
1599impl<T, A> GuardedFutureReader<T, A>
1600where
1601 A: AsAccessor,
1602{
1603 pub fn new(accessor: A, reader: FutureReader<T>) -> Self {
1611 assert!(
1612 accessor
1613 .as_accessor()
1614 .with(|a| a.as_context().0.concurrency_support())
1615 );
1616 Self {
1617 reader: Some(reader),
1618 accessor,
1619 }
1620 }
1621
1622 pub fn into_future(self) -> FutureReader<T> {
1625 self.into()
1626 }
1627}
1628
1629impl<T, A> From<GuardedFutureReader<T, A>> for FutureReader<T>
1630where
1631 A: AsAccessor,
1632{
1633 fn from(mut guard: GuardedFutureReader<T, A>) -> Self {
1634 guard.reader.take().unwrap()
1635 }
1636}
1637
1638impl<T, A> Drop for GuardedFutureReader<T, A>
1639where
1640 A: AsAccessor,
1641{
1642 fn drop(&mut self) {
1643 if let Some(reader) = &mut self.reader {
1644 let result = reader.close_with(&self.accessor);
1647 debug_assert!(result.is_ok());
1648 }
1649 }
1650}
1651
1652pub struct StreamReader<T> {
1659 id: TableId<TransmitHandle>,
1660 _phantom: PhantomData<T>,
1661}
1662
1663impl<T> StreamReader<T> {
1664 pub fn new<S: AsContextMut>(
1673 mut store: S,
1674 producer: impl StreamProducer<S::Data, Item = T>,
1675 ) -> Result<Self>
1676 where
1677 T: func::Lower + func::Lift + Send + Sync + 'static,
1678 {
1679 ensure!(
1680 store.as_context().0.concurrency_support(),
1681 "concurrency support is not enabled",
1682 );
1683 Ok(Self::new_(
1684 store
1685 .as_context_mut()
1686 .new_transmit(TransmitKind::Stream, producer)?,
1687 ))
1688 }
1689
1690 pub(super) fn new_(id: TableId<TransmitHandle>) -> Self {
1691 Self {
1692 id,
1693 _phantom: PhantomData,
1694 }
1695 }
1696
1697 pub(super) fn id(&self) -> TableId<TransmitHandle> {
1698 self.id
1699 }
1700
1701 pub fn try_into<V: 'static>(mut self, mut store: impl AsContextMut) -> Result<V, Self> {
1723 let store = store.as_context_mut();
1724 let state = store.0.concurrent_state_mut_already_forced_current_thread();
1725 let id = state.get_mut(self.id).unwrap().state;
1726 if let WriteState::HostReady { try_into, .. } = &state.get_mut(id).unwrap().write {
1727 match try_into(TypeId::of::<V>()) {
1728 Some(result) => {
1729 self.close(store).unwrap();
1730 Ok(*result.downcast::<V>().unwrap())
1731 }
1732 None => Err(self),
1733 }
1734 } else {
1735 Err(self)
1736 }
1737 }
1738
1739 pub fn pipe<S: AsContextMut>(
1749 self,
1750 mut store: S,
1751 consumer: impl StreamConsumer<S::Data, Item = T>,
1752 ) -> Result<()>
1753 where
1754 T: 'static,
1755 {
1756 store
1757 .as_context_mut()
1758 .set_consumer(self.id, TransmitKind::Stream, consumer)
1759 }
1760
1761 fn lift_from_index(cx: &mut LiftContext<'_>, ty: InterfaceType, index: u32) -> Result<Self> {
1763 let id = lift_index_to_stream(cx, ty, index)?;
1764 Ok(Self::new_(id))
1765 }
1766
1767 pub fn close(&mut self, mut store: impl AsContextMut) -> Result<()> {
1783 stream_close(store.as_context_mut().0, &mut self.id)
1784 }
1785
1786 pub fn close_with(&mut self, accessor: impl AsAccessor) -> Result<()> {
1788 accessor.as_accessor().with(|access| self.close(access))
1789 }
1790
1791 pub fn guard<A>(self, accessor: A) -> GuardedStreamReader<T, A>
1797 where
1798 A: AsAccessor,
1799 {
1800 GuardedStreamReader::new(accessor, self)
1801 }
1802
1803 pub fn try_into_stream_any(self, store: impl AsContextMut) -> Result<StreamAny>
1810 where
1811 T: ComponentType + 'static,
1812 {
1813 StreamAny::try_from_stream_reader(store, self)
1814 }
1815
1816 pub fn try_from_stream_any(stream: StreamAny) -> Result<Self>
1823 where
1824 T: ComponentType + 'static,
1825 {
1826 stream.try_into_stream_reader()
1827 }
1828}
1829
1830impl<T> fmt::Debug for StreamReader<T> {
1831 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1832 f.debug_struct("StreamReader")
1833 .field("id", &self.id)
1834 .finish()
1835 }
1836}
1837
1838pub(super) fn stream_close(
1839 store: &mut StoreOpaque,
1840 id: &mut TableId<TransmitHandle>,
1841) -> Result<()> {
1842 let id = mem::replace(id, TableId::new(u32::MAX));
1843 store.host_drop_reader(id, TransmitKind::Stream)
1844}
1845
1846pub(super) fn lift_index_to_stream(
1848 cx: &mut LiftContext<'_>,
1849 ty: InterfaceType,
1850 index: u32,
1851) -> Result<TableId<TransmitHandle>> {
1852 match ty {
1853 InterfaceType::Stream(src) => {
1854 let (state, instance) = cx.concurrent_state_and_instance_mut();
1855 lift_index_to_transmit(instance, state, TransmitIndex::Stream(src), index)
1856 }
1857 _ => func::bad_type_info(),
1858 }
1859}
1860
1861pub(super) fn lower_stream_to_index<U>(
1863 id: TableId<TransmitHandle>,
1864 cx: &mut LowerContext<'_, U>,
1865 ty: InterfaceType,
1866) -> Result<u32> {
1867 match ty {
1868 InterfaceType::Stream(dst) => {
1869 cx.instance_handle()
1870 .lower_transmit_to_index(cx.store.0, TransmitIndex::Stream(dst), id)
1871 }
1872 _ => func::bad_type_info(),
1873 }
1874}
1875
1876unsafe impl<T: ComponentType> ComponentType for StreamReader<T> {
1879 const ABI: CanonicalAbiInfo = CanonicalAbiInfo::SCALAR4;
1880
1881 type Lower = <u32 as func::ComponentType>::Lower;
1882
1883 fn typecheck(ty: &InterfaceType, types: &InstanceType<'_>) -> Result<()> {
1884 match ty {
1885 InterfaceType::Stream(ty) => {
1886 let ty = types.types[*ty].ty;
1887 types::typecheck_payload::<T>(types.types[ty].payload.as_ref(), types)
1888 }
1889 other => bail!("expected `stream`, found `{}`", func::desc(other)),
1890 }
1891 }
1892}
1893
1894unsafe impl<T: ComponentType> func::Lower for StreamReader<T> {
1896 fn linear_lower_to_flat<U>(
1897 &self,
1898 cx: &mut LowerContext<'_, U>,
1899 ty: InterfaceType,
1900 dst: &mut MaybeUninit<Self::Lower>,
1901 ) -> Result<()> {
1902 lower_stream_to_index(self.id, cx, ty)?.linear_lower_to_flat(cx, InterfaceType::U32, dst)
1903 }
1904
1905 fn linear_lower_to_memory<U>(
1906 &self,
1907 cx: &mut LowerContext<'_, U>,
1908 ty: InterfaceType,
1909 offset: usize,
1910 ) -> Result<()> {
1911 lower_stream_to_index(self.id, cx, ty)?.linear_lower_to_memory(
1912 cx,
1913 InterfaceType::U32,
1914 offset,
1915 )
1916 }
1917}
1918
1919unsafe impl<T: ComponentType> func::Lift for StreamReader<T> {
1921 fn linear_lift_from_flat(
1922 cx: &mut LiftContext<'_>,
1923 ty: InterfaceType,
1924 src: &Self::Lower,
1925 ) -> Result<Self> {
1926 let index = u32::linear_lift_from_flat(cx, InterfaceType::U32, src)?;
1927 Self::lift_from_index(cx, ty, index)
1928 }
1929
1930 fn linear_lift_from_memory(
1931 cx: &mut LiftContext<'_>,
1932 ty: InterfaceType,
1933 bytes: &[u8],
1934 ) -> Result<Self> {
1935 let index = u32::linear_lift_from_memory(cx, InterfaceType::U32, bytes)?;
1936 Self::lift_from_index(cx, ty, index)
1937 }
1938}
1939
1940pub struct GuardedStreamReader<T, A>
1948where
1949 A: AsAccessor,
1950{
1951 reader: Option<StreamReader<T>>,
1955 accessor: A,
1956}
1957
1958impl<T, A> GuardedStreamReader<T, A>
1959where
1960 A: AsAccessor,
1961{
1962 pub fn new(accessor: A, reader: StreamReader<T>) -> Self {
1971 assert!(
1972 accessor
1973 .as_accessor()
1974 .with(|a| a.as_context().0.concurrency_support())
1975 );
1976 Self {
1977 reader: Some(reader),
1978 accessor,
1979 }
1980 }
1981
1982 pub fn into_stream(self) -> StreamReader<T> {
1985 self.into()
1986 }
1987}
1988
1989impl<T, A> From<GuardedStreamReader<T, A>> for StreamReader<T>
1990where
1991 A: AsAccessor,
1992{
1993 fn from(mut guard: GuardedStreamReader<T, A>) -> Self {
1994 guard.reader.take().unwrap()
1995 }
1996}
1997
1998impl<T, A> Drop for GuardedStreamReader<T, A>
1999where
2000 A: AsAccessor,
2001{
2002 fn drop(&mut self) {
2003 if let Some(reader) = &mut self.reader {
2004 let result = reader.close_with(&self.accessor);
2007 debug_assert!(result.is_ok());
2008 }
2009 }
2010}
2011
2012pub struct ErrorContext {
2014 rep: u32,
2015}
2016
2017impl ErrorContext {
2018 pub(crate) fn new(rep: u32) -> Self {
2019 Self { rep }
2020 }
2021
2022 pub fn into_val(self) -> Val {
2024 Val::ErrorContext(ErrorContextAny(self.rep))
2025 }
2026
2027 pub fn from_val(_: impl AsContextMut, value: &Val) -> Result<Self> {
2029 let Val::ErrorContext(ErrorContextAny(rep)) = value else {
2030 bail!("expected `error-context`; got `{}`", value.desc());
2031 };
2032 Ok(Self::new(*rep))
2033 }
2034
2035 fn lift_from_index(cx: &mut LiftContext<'_>, ty: InterfaceType, index: u32) -> Result<Self> {
2036 match ty {
2037 InterfaceType::ErrorContext(src) => {
2038 let rep = cx
2039 .instance_mut()
2040 .table_for_error_context(src)
2041 .error_context_rep(index)?;
2042
2043 Ok(Self { rep })
2044 }
2045 _ => func::bad_type_info(),
2046 }
2047 }
2048}
2049
2050pub(crate) fn lower_error_context_to_index<U>(
2051 rep: u32,
2052 cx: &mut LowerContext<'_, U>,
2053 ty: InterfaceType,
2054) -> Result<u32> {
2055 match ty {
2056 InterfaceType::ErrorContext(dst) => {
2057 let tbl = cx.instance_mut().table_for_error_context(dst);
2058 tbl.error_context_insert(rep)
2059 }
2060 _ => func::bad_type_info(),
2061 }
2062}
2063unsafe impl func::ComponentType for ErrorContext {
2066 const ABI: CanonicalAbiInfo = CanonicalAbiInfo::SCALAR4;
2067
2068 type Lower = <u32 as func::ComponentType>::Lower;
2069
2070 fn typecheck(ty: &InterfaceType, _types: &InstanceType<'_>) -> Result<()> {
2071 match ty {
2072 InterfaceType::ErrorContext(_) => Ok(()),
2073 other => bail!("expected `error`, found `{}`", func::desc(other)),
2074 }
2075 }
2076}
2077
2078unsafe impl func::Lower for ErrorContext {
2080 fn linear_lower_to_flat<T>(
2081 &self,
2082 cx: &mut LowerContext<'_, T>,
2083 ty: InterfaceType,
2084 dst: &mut MaybeUninit<Self::Lower>,
2085 ) -> Result<()> {
2086 lower_error_context_to_index(self.rep, cx, ty)?.linear_lower_to_flat(
2087 cx,
2088 InterfaceType::U32,
2089 dst,
2090 )
2091 }
2092
2093 fn linear_lower_to_memory<T>(
2094 &self,
2095 cx: &mut LowerContext<'_, T>,
2096 ty: InterfaceType,
2097 offset: usize,
2098 ) -> Result<()> {
2099 lower_error_context_to_index(self.rep, cx, ty)?.linear_lower_to_memory(
2100 cx,
2101 InterfaceType::U32,
2102 offset,
2103 )
2104 }
2105}
2106
2107unsafe impl func::Lift for ErrorContext {
2109 fn linear_lift_from_flat(
2110 cx: &mut LiftContext<'_>,
2111 ty: InterfaceType,
2112 src: &Self::Lower,
2113 ) -> Result<Self> {
2114 let index = u32::linear_lift_from_flat(cx, InterfaceType::U32, src)?;
2115 Self::lift_from_index(cx, ty, index)
2116 }
2117
2118 fn linear_lift_from_memory(
2119 cx: &mut LiftContext<'_>,
2120 ty: InterfaceType,
2121 bytes: &[u8],
2122 ) -> Result<Self> {
2123 let index = u32::linear_lift_from_memory(cx, InterfaceType::U32, bytes)?;
2124 Self::lift_from_index(cx, ty, index)
2125 }
2126}
2127
2128pub(super) struct TransmitHandle {
2130 pub(super) common: WaitableCommon,
2131 state: TableId<TransmitState>,
2133}
2134
2135impl TransmitHandle {
2136 fn new(state: TableId<TransmitState>) -> Self {
2137 Self {
2138 common: WaitableCommon::default(),
2139 state,
2140 }
2141 }
2142}
2143
2144impl TableDebug for TransmitHandle {
2145 fn type_name() -> &'static str {
2146 "TransmitHandle"
2147 }
2148}
2149
2150struct TransmitState {
2152 write_handle: TableId<TransmitHandle>,
2154 read_handle: TableId<TransmitHandle>,
2156 write: WriteState,
2158 read: ReadState,
2160 done: bool,
2162 pub(super) origin: TransmitOrigin,
2165}
2166
2167#[derive(Copy, Clone)]
2168pub(super) enum TransmitOrigin {
2169 Host,
2170 GuestFuture(ComponentInstanceId, TypeFutureTableIndex),
2171 GuestStream(ComponentInstanceId, TypeStreamTableIndex),
2172}
2173
2174impl TransmitState {
2175 fn new(origin: TransmitOrigin) -> Self {
2176 Self {
2177 write_handle: TableId::new(u32::MAX),
2178 read_handle: TableId::new(u32::MAX),
2179 read: ReadState::Open,
2180 write: WriteState::Open,
2181 done: false,
2182 origin,
2183 }
2184 }
2185}
2186
2187impl TableDebug for TransmitState {
2188 fn type_name() -> &'static str {
2189 "TransmitState"
2190 }
2191}
2192
2193impl TransmitOrigin {
2194 fn guest(id: ComponentInstanceId, index: TransmitIndex) -> Self {
2195 match index {
2196 TransmitIndex::Future(ty) => TransmitOrigin::GuestFuture(id, ty),
2197 TransmitIndex::Stream(ty) => TransmitOrigin::GuestStream(id, ty),
2198 }
2199 }
2200}
2201
2202type PollStream = Box<
2203 dyn Fn() -> Pin<Box<dyn Future<Output = Result<StreamResult>> + Send + 'static>> + Send + Sync,
2204>;
2205
2206type TryInto = Box<dyn Fn(TypeId) -> Option<Box<dyn Any>> + Send + Sync>;
2207
2208enum WriteState {
2210 Open,
2212 GuestReady {
2214 instance: Instance,
2215 caller: RuntimeComponentInstanceIndex,
2216 ty: TransmitIndex,
2217 flat_abi: Option<FlatAbi>,
2218 options: OptionsIndex,
2219 address: usize,
2220 count: ItemCount,
2221 handle: u32,
2222 },
2223 HostReady {
2225 produce: PollStream,
2226 try_into: TryInto,
2227 guest_offset: ItemCount,
2228 cancel: bool,
2229 cancel_waker: Option<Waker>,
2230 },
2231 Dropped,
2233}
2234
2235impl fmt::Debug for WriteState {
2236 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
2237 match self {
2238 Self::Open => f.debug_tuple("Open").finish(),
2239 Self::GuestReady { .. } => f.debug_tuple("GuestReady").finish(),
2240 Self::HostReady { .. } => f.debug_tuple("HostReady").finish(),
2241 Self::Dropped => f.debug_tuple("Dropped").finish(),
2242 }
2243 }
2244}
2245
2246enum ReadState {
2248 Open,
2250 GuestReady {
2252 ty: TransmitIndex,
2253 caller_instance: RuntimeComponentInstanceIndex,
2254 caller_thread: QualifiedThreadId,
2255 flat_abi: Option<FlatAbi>,
2256 instance: Instance,
2257 options: OptionsIndex,
2258 address: usize,
2259 count: ItemCount,
2260 handle: u32,
2261 },
2262 HostReady {
2264 consume: PollStream,
2265 guest_offset: ItemCount,
2266 cancel: bool,
2267 cancel_waker: Option<Waker>,
2268 },
2269 HostToHost {
2271 accept: Box<
2272 dyn for<'a> Fn(
2273 &'a mut UntypedWriteBuffer<'a>,
2274 )
2275 -> Pin<Box<dyn Future<Output = Result<StreamResult>> + Send + 'a>>
2276 + Send
2277 + Sync,
2278 >,
2279 buffer: Vec<u8>,
2280 limit: usize,
2281 },
2282 Dropped,
2284}
2285
2286impl fmt::Debug for ReadState {
2287 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
2288 match self {
2289 Self::Open => f.debug_tuple("Open").finish(),
2290 Self::GuestReady { .. } => f.debug_tuple("GuestReady").finish(),
2291 Self::HostReady { .. } => f.debug_tuple("HostReady").finish(),
2292 Self::HostToHost { .. } => f.debug_tuple("HostToHost").finish(),
2293 Self::Dropped => f.debug_tuple("Dropped").finish(),
2294 }
2295 }
2296}
2297
2298fn return_code(kind: TransmitKind, state: StreamResult, count: ItemCount) -> Result<ReturnCode> {
2299 Ok(match state {
2300 StreamResult::Dropped => ReturnCode::Dropped(count),
2301 StreamResult::Completed => ReturnCode::completed(kind, count),
2302 StreamResult::Cancelled => ReturnCode::Cancelled(count),
2303 })
2304}
2305
2306fn settle_host_read(
2307 transmit: &mut TransmitState,
2308 kind: TransmitKind,
2309 state: StreamResult,
2310) -> Result<ReturnCode> {
2311 let ReadState::HostReady {
2312 consume,
2313 guest_offset,
2314 ..
2315 } = mem::replace(&mut transmit.read, ReadState::Open)
2316 else {
2317 bail_bug!("expected ReadState::HostReady")
2318 };
2319 let code = return_code(kind, state, guest_offset)?;
2320 transmit.read = match state {
2321 StreamResult::Dropped => ReadState::Dropped,
2322 StreamResult::Completed | StreamResult::Cancelled => ReadState::HostReady {
2323 consume,
2324 guest_offset: ItemCount::ZERO,
2325 cancel: false,
2326 cancel_waker: None,
2327 },
2328 };
2329 Ok(code)
2330}
2331
2332fn settle_host_write(
2333 transmit: &mut TransmitState,
2334 kind: TransmitKind,
2335 state: StreamResult,
2336) -> Result<ReturnCode> {
2337 let WriteState::HostReady {
2338 produce,
2339 try_into,
2340 guest_offset,
2341 ..
2342 } = mem::replace(&mut transmit.write, WriteState::Open)
2343 else {
2344 bail_bug!("expected WriteState::HostReady")
2345 };
2346 let code = return_code(kind, state, guest_offset)?;
2347 transmit.write = match state {
2348 StreamResult::Dropped => WriteState::Dropped,
2349 StreamResult::Completed | StreamResult::Cancelled => WriteState::HostReady {
2350 produce,
2351 try_into,
2352 guest_offset: ItemCount::ZERO,
2353 cancel: false,
2354 cancel_waker: None,
2355 },
2356 };
2357 Ok(code)
2358}
2359
2360impl StoreOpaque {
2361 fn pipe_from_guest(
2362 &mut self,
2363 kind: TransmitKind,
2364 id: TableId<TransmitState>,
2365 future: Pin<Box<dyn Future<Output = Result<StreamResult>> + Send + 'static>>,
2366 ) {
2367 let future = async move {
2368 let stream_state = future.await?;
2369 tls::get(|store| {
2370 let state = store.concurrent_state_mut()?;
2371 let transmit = state.get_mut(id)?;
2372 let code = settle_host_read(transmit, kind, stream_state)?;
2373 let WriteState::GuestReady { ty, handle, .. } =
2374 mem::replace(&mut transmit.write, WriteState::Open)
2375 else {
2376 bail_bug!("expected WriteState::GuestReady")
2377 };
2378 state.send_write_result(ty, id, handle, code)?;
2379 Ok(())
2380 })
2381 };
2382
2383 self.concurrent_state_mut_already_forced_current_thread()
2384 .push_future(future.boxed());
2385 }
2386
2387 fn pipe_to_guest(
2388 &mut self,
2389 kind: TransmitKind,
2390 id: TableId<TransmitState>,
2391 future: Pin<Box<dyn Future<Output = Result<StreamResult>> + Send + 'static>>,
2392 ) {
2393 let future = async move {
2394 let stream_state = future.await?;
2395 tls::get(|store| {
2396 let state = store.concurrent_state_mut()?;
2397 let transmit = state.get_mut(id)?;
2398 let code = settle_host_write(transmit, kind, stream_state)?;
2399 let ReadState::GuestReady { ty, handle, .. } =
2400 mem::replace(&mut transmit.read, ReadState::Open)
2401 else {
2402 bail_bug!("expected ReadState::GuestReady")
2403 };
2404 state.send_read_result(ty, id, handle, code)?;
2405 Ok(())
2406 })
2407 };
2408
2409 self.concurrent_state_mut_already_forced_current_thread()
2410 .push_future(future.boxed());
2411 }
2412
2413 fn host_drop_reader(&mut self, id: TableId<TransmitHandle>, kind: TransmitKind) -> Result<()> {
2415 let state = self.concurrent_state_mut()?;
2416 Waitable::Transmit(id).join(state, None)?;
2417 let transmit_id = state.get_mut(id)?.state;
2418 let transmit = state
2419 .get_mut(transmit_id)
2420 .with_context(|| format!("error closing reader {transmit_id:?}"))?;
2421 log::trace!(
2422 "host_drop_reader state {transmit_id:?}; read state {:?} write state {:?}",
2423 transmit.read,
2424 transmit.write
2425 );
2426
2427 transmit.read = ReadState::Dropped;
2428
2429 let new_state = if let WriteState::Dropped = &transmit.write {
2432 WriteState::Dropped
2433 } else {
2434 WriteState::Open
2435 };
2436
2437 let write_handle = transmit.write_handle;
2438
2439 match mem::replace(&mut transmit.write, new_state) {
2440 WriteState::GuestReady { ty, handle, .. } => {
2443 state.update_event(
2444 write_handle.rep(),
2445 match ty {
2446 TransmitIndex::Future(ty) => Event::FutureWrite {
2447 code: ReturnCode::Dropped(ItemCount::ZERO),
2448 pending: Some((ty, handle)),
2449 },
2450 TransmitIndex::Stream(ty) => Event::StreamWrite {
2451 code: ReturnCode::Dropped(ItemCount::ZERO),
2452 pending: Some((ty, handle)),
2453 },
2454 },
2455 )?;
2456 }
2457
2458 WriteState::Open => {
2459 state.update_event(
2460 write_handle.rep(),
2461 match kind {
2462 TransmitKind::Future => Event::FutureWrite {
2463 code: ReturnCode::Dropped(ItemCount::ZERO),
2464 pending: None,
2465 },
2466 TransmitKind::Stream => Event::StreamWrite {
2467 code: ReturnCode::Dropped(ItemCount::ZERO),
2468 pending: None,
2469 },
2470 },
2471 )?;
2472 }
2473
2474 WriteState::Dropped | WriteState::HostReady { .. } => {
2479 log::trace!("host_drop_reader delete {transmit_id:?}");
2480 state.delete_transmit(transmit_id)?;
2481 }
2482 }
2483 Ok(())
2484 }
2485
2486 fn host_drop_writer(
2488 &mut self,
2489 id: TableId<TransmitHandle>,
2490 on_drop_open: Option<fn() -> Result<()>>,
2491 ) -> Result<()> {
2492 let state = self.concurrent_state_mut()?;
2493 Waitable::Transmit(id).join(state, None)?;
2494 let transmit_id = state.get_mut(id)?.state;
2495 let transmit = state
2496 .get_mut(transmit_id)
2497 .with_context(|| format!("error closing writer {transmit_id:?}"))?;
2498 log::trace!(
2499 "host_drop_writer state {transmit_id:?}; read state {:?} writer state {:?}",
2500 transmit.read,
2501 transmit.write
2502 );
2503
2504 match &mut transmit.write {
2506 WriteState::GuestReady { .. } => {
2507 bail_bug!("can't call `host_drop_writer` on a guest-owned writer");
2508 }
2509 WriteState::HostReady { .. } => {}
2510 v @ WriteState::Open => {
2511 if let (Some(on_drop_open), false) = (on_drop_open, transmit.done) {
2512 on_drop_open()?;
2513 } else {
2514 *v = WriteState::Dropped;
2515 }
2516 }
2517 WriteState::Dropped => bail_bug!("write state is already dropped"),
2518 }
2519
2520 let transmit = self.concurrent_state_mut()?.get_mut(transmit_id)?;
2521
2522 let new_state = if let ReadState::Dropped = &transmit.read {
2528 ReadState::Dropped
2529 } else {
2530 ReadState::Open
2531 };
2532
2533 let read_handle = transmit.read_handle;
2534
2535 match mem::replace(&mut transmit.read, new_state) {
2537 ReadState::GuestReady { ty, handle, .. } => {
2541 self.concurrent_state_mut()?.update_event(
2543 read_handle.rep(),
2544 match ty {
2545 TransmitIndex::Future(ty) => Event::FutureRead {
2546 code: ReturnCode::Dropped(ItemCount::ZERO),
2547 pending: Some((ty, handle)),
2548 },
2549 TransmitIndex::Stream(ty) => Event::StreamRead {
2550 code: ReturnCode::Dropped(ItemCount::ZERO),
2551 pending: Some((ty, handle)),
2552 },
2553 },
2554 )?;
2555 }
2556
2557 ReadState::Open => {
2559 self.concurrent_state_mut()?.update_event(
2560 read_handle.rep(),
2561 match on_drop_open {
2562 Some(_) => Event::FutureRead {
2563 code: ReturnCode::Dropped(ItemCount::ZERO),
2564 pending: None,
2565 },
2566 None => Event::StreamRead {
2567 code: ReturnCode::Dropped(ItemCount::ZERO),
2568 pending: None,
2569 },
2570 },
2571 )?;
2572 }
2573
2574 ReadState::Dropped | ReadState::HostReady { .. } | ReadState::HostToHost { .. } => {
2581 log::trace!("host_drop_writer delete {transmit_id:?}");
2582 self.concurrent_state_mut()?.delete_transmit(transmit_id)?;
2583 }
2584 }
2585 Ok(())
2586 }
2587
2588 pub(super) fn transmit_origin(
2589 &mut self,
2590 id: TableId<TransmitHandle>,
2591 ) -> Result<TransmitOrigin> {
2592 let state = self.concurrent_state_mut()?;
2593 let state_id = state.get_mut(id)?.state;
2594 Ok(state.get_mut(state_id)?.origin)
2595 }
2596}
2597
2598impl<T> StoreContextMut<'_, T> {
2599 fn new_transmit<P: StreamProducer<T>>(
2600 mut self,
2601 kind: TransmitKind,
2602 producer: P,
2603 ) -> Result<TableId<TransmitHandle>>
2604 where
2605 P::Item: func::Lower,
2606 {
2607 let token = StoreToken::new(self.as_context_mut());
2608 let state = self.0.concurrent_state_mut()?;
2609 let (_, read) = state.new_transmit(TransmitOrigin::Host)?;
2610 let producer = Arc::new(LockedState::new((Box::pin(producer), P::Buffer::default())));
2611 let id = state.get_mut(read)?.state;
2612 let mut dropped = false;
2613 let produce = Box::new({
2614 let producer = producer.clone();
2615 move || {
2616 let producer = producer.clone();
2617 async move {
2618 let mut state = producer.take()?;
2619 let (mine, buffer) = &mut *state;
2620
2621 let (result, cancelled) = if buffer.remaining().is_empty() {
2622 future::poll_fn(|cx| {
2623 tls::get(|store| {
2624 let transmit = store.concurrent_state_mut()?.get_mut(id)?;
2625
2626 let &WriteState::HostReady { cancel, .. } = &transmit.write else {
2627 bail_bug!("expected WriteState::HostReady")
2628 };
2629
2630 let mut host_written = 0;
2631 let mut host_buffer =
2632 if let ReadState::HostToHost { buffer, .. } = &mut transmit.read {
2633 Some(mem::take(buffer))
2634 } else {
2635 None
2636 };
2637
2638 let poll = mine.as_mut().poll_produce(
2639 cx,
2640 token.as_context_mut(store),
2641 Destination {
2642 id,
2643 buffer,
2644 host_buffer: host_buffer.as_mut().map(|b| {
2645 HostBuffer {
2646 dst: b,
2647 marked_written: &mut host_written,
2648 }
2649 }),
2650 _phantom: PhantomData,
2651 },
2652 cancel,
2653 );
2654
2655 let transmit = store.concurrent_state_mut()?.get_mut(id)?;
2656
2657 let host_offset = if let (
2658 Some(host_buffer),
2659 ReadState::HostToHost { buffer, limit, .. },
2660 ) = (host_buffer, &mut transmit.read)
2661 {
2662 *limit = host_written;
2663 *buffer = host_buffer;
2664 *limit
2665 } else {
2666 0
2667 };
2668
2669 {
2670 let WriteState::HostReady {
2671 guest_offset,
2672 cancel,
2673 cancel_waker,
2674 ..
2675 } = &mut transmit.write
2676 else {
2677 bail_bug!("expected WriteState::HostReady")
2678 };
2679
2680 if poll.is_pending() {
2681 if !buffer.remaining().is_empty()
2682 || *guest_offset > 0
2683 || host_offset > 0
2684 {
2685 bail!(
2686 "StreamProducer::poll_produce returned Poll::Pending \
2687 after producing at least one item"
2688 )
2689 }
2690 *cancel_waker = Some(cx.waker().clone());
2691 } else {
2692 *cancel_waker = None;
2693 *cancel = false;
2694 }
2695 }
2696
2697 Ok(poll.map(|v| v.map(|result| (result, cancel))))
2698 })?
2699 })
2700 .await?
2701 } else {
2702 (StreamResult::Completed, false)
2703 };
2704
2705 let (guest_offset, host_offset, count) = tls::get(|store| {
2706 let transmit = store.concurrent_state_mut()?.get_mut(id)?;
2707 let (count, host_offset) = match &transmit.read {
2708 &ReadState::GuestReady { count, .. } => (count.as_u32(), 0),
2709 &ReadState::HostToHost { limit, .. } => (1, limit),
2710 _ => bail_bug!("invalid read state"),
2711 };
2712 let guest_offset = match &transmit.write {
2713 &WriteState::HostReady { guest_offset, .. } => guest_offset,
2714 _ => bail_bug!("invalid write state"),
2715 };
2716 Ok((guest_offset, host_offset, count))
2717 })?;
2718
2719 match result {
2720 StreamResult::Completed => {
2721 if count > 1
2722 && buffer.remaining().is_empty()
2723 && guest_offset == 0
2724 && host_offset == 0
2725 {
2726 bail!(
2727 "StreamProducer::poll_produce returned StreamResult::Completed \
2728 without producing any items"
2729 );
2730 }
2731 }
2732 StreamResult::Cancelled => {
2733 if !cancelled {
2734 bail!(
2735 "StreamProducer::poll_produce returned StreamResult::Cancelled \
2736 without being given a `finish` parameter value of true"
2737 );
2738 }
2739 }
2740 StreamResult::Dropped => {
2741 dropped = true;
2742 }
2743 }
2744
2745 let write_buffer = !buffer.remaining().is_empty() || host_offset > 0;
2746
2747 drop(state);
2748
2749 if write_buffer {
2750 write(token, id, producer.clone(), kind).await?;
2751 }
2752
2753 Ok(if dropped {
2754 if producer.with(|p| p.1.remaining().is_empty())? {
2755 StreamResult::Dropped
2756 } else {
2757 StreamResult::Completed
2758 }
2759 } else {
2760 result
2761 })
2762 }
2763 .boxed()
2764 }
2765 });
2766 let try_into = Box::new(move |ty| {
2767 let (mine, buffer) = producer.try_lock().ok()?.take()?;
2768 match P::try_into(mine, ty) {
2769 Ok(value) => Some(value),
2770 Err(mine) => {
2771 *producer.try_lock().ok()? = Some((mine, buffer));
2772 None
2773 }
2774 }
2775 });
2776 state.get_mut(id)?.write = WriteState::HostReady {
2777 produce,
2778 try_into,
2779 guest_offset: ItemCount::ZERO,
2780 cancel: false,
2781 cancel_waker: None,
2782 };
2783 Ok(read)
2784 }
2785
2786 fn set_consumer<C: StreamConsumer<T>>(
2787 mut self,
2788 id: TableId<TransmitHandle>,
2789 kind: TransmitKind,
2790 consumer: C,
2791 ) -> Result<()> {
2792 let token = StoreToken::new(self.as_context_mut());
2793 let state = self.0.concurrent_state_mut()?;
2794 let id = state.get_mut(id)?.state;
2795 let transmit = state.get_mut(id)?;
2796 let consumer = Arc::new(LockedState::new(Box::pin(consumer)));
2797 let consume_with_buffer = {
2798 let consumer = consumer.clone();
2799 async move |mut host_buffer: Option<&mut dyn WriteBuffer<C::Item>>| {
2800 let mut mine = consumer.take()?;
2801
2802 let host_buffer_remaining_before =
2803 host_buffer.as_deref_mut().map(|v| v.remaining().len());
2804
2805 let (result, cancelled) = future::poll_fn(|cx| {
2806 tls::get(|store| {
2807 let cancel = match &store.concurrent_state_mut()?.get_mut(id)?.read {
2808 &ReadState::HostReady { cancel, .. } => cancel,
2809 ReadState::Open => false,
2810 _ => bail_bug!("unexpected read state"),
2811 };
2812
2813 let poll = mine.as_mut().poll_consume(
2814 cx,
2815 token.as_context_mut(store),
2816 Source {
2817 id,
2818 host_buffer: host_buffer.as_deref_mut(),
2819 },
2820 cancel,
2821 );
2822
2823 if let ReadState::HostReady {
2824 cancel_waker,
2825 cancel,
2826 ..
2827 } = &mut store.concurrent_state_mut()?.get_mut(id)?.read
2828 {
2829 if poll.is_pending() {
2830 *cancel_waker = Some(cx.waker().clone());
2831 } else {
2832 *cancel_waker = None;
2833 *cancel = false;
2834 }
2835 }
2836
2837 Ok(poll.map(|v| v.map(|result| (result, cancel))))
2838 })?
2839 })
2840 .await?;
2841
2842 let (guest_offset, count) = tls::get(|store| {
2843 let transmit = store.concurrent_state_mut()?.get_mut(id)?;
2844 Ok((
2845 match &transmit.read {
2846 &ReadState::HostReady { guest_offset, .. } => guest_offset,
2847 ReadState::Open => ItemCount::ZERO,
2848 _ => bail_bug!("invalid read state"),
2849 },
2850 match &transmit.write {
2851 WriteState::GuestReady { count, .. } => count.as_usize(),
2852 WriteState::HostReady { .. } => match host_buffer_remaining_before {
2853 Some(n) => n,
2854 None => bail_bug!("host_buffer_remaining_before should be set"),
2855 },
2856 _ => bail_bug!("invalid write state"),
2857 },
2858 ))
2859 })?;
2860
2861 match result {
2862 StreamResult::Completed => {
2863 if count > 0
2864 && guest_offset == 0
2865 && host_buffer_remaining_before
2866 .zip(host_buffer.map(|v| v.remaining().len()))
2867 .map(|(before, after)| before == after)
2868 .unwrap_or(false)
2869 {
2870 bail!(
2871 "StreamConsumer::poll_consume returned StreamResult::Completed \
2872 without consuming any items"
2873 );
2874 }
2875
2876 if let TransmitKind::Future = kind {
2877 tls::get(|store| {
2878 store.concurrent_state_mut()?.get_mut(id)?.done = true;
2879 crate::error::Ok(())
2880 })?;
2881 }
2882 }
2883 StreamResult::Cancelled => {
2884 if !cancelled {
2885 bail!(
2886 "StreamConsumer::poll_consume returned StreamResult::Cancelled \
2887 without being given a `finish` parameter value of true"
2888 );
2889 }
2890 }
2891 StreamResult::Dropped => {}
2892 }
2893
2894 Ok(result)
2895 }
2896 };
2897 let consume = {
2898 let consume = consume_with_buffer.clone();
2899 Box::new(move || {
2900 let consume = consume.clone();
2901 async move { consume(None).await }.boxed()
2902 })
2903 };
2904
2905 match &transmit.write {
2906 WriteState::Open => {
2907 transmit.read = ReadState::HostReady {
2908 consume,
2909 guest_offset: ItemCount::ZERO,
2910 cancel: false,
2911 cancel_waker: None,
2912 };
2913 }
2914 &WriteState::GuestReady { .. } => {
2915 let future = consume();
2916 transmit.read = ReadState::HostReady {
2917 consume,
2918 guest_offset: ItemCount::ZERO,
2919 cancel: false,
2920 cancel_waker: None,
2921 };
2922 self.0.pipe_from_guest(kind, id, future);
2923 }
2924 WriteState::HostReady { .. } => {
2925 let WriteState::HostReady { produce, .. } = mem::replace(
2926 &mut transmit.write,
2927 WriteState::HostReady {
2928 produce: Box::new(|| {
2929 Box::pin(async { bail_bug!("unexpected invocation of `produce`") })
2930 }),
2931 try_into: Box::new(|_| None),
2932 guest_offset: ItemCount::ZERO,
2933 cancel: false,
2934 cancel_waker: None,
2935 },
2936 ) else {
2937 bail_bug!("expected WriteState::HostReady")
2938 };
2939
2940 transmit.read = ReadState::HostToHost {
2941 accept: Box::new(move |input| {
2942 let consume = consume_with_buffer.clone();
2943 async move { consume(Some(input.get_mut::<C::Item>())).await }.boxed()
2944 }),
2945 buffer: Vec::new(),
2946 limit: 0,
2947 };
2948
2949 let future = async move {
2950 loop {
2951 if tls::get(|store| {
2952 crate::error::Ok(matches!(
2953 store.concurrent_state_mut()?.get_mut(id)?.read,
2954 ReadState::Dropped
2955 ))
2956 })? {
2957 break Ok(());
2958 }
2959
2960 match produce().await? {
2961 StreamResult::Completed | StreamResult::Cancelled => {}
2962 StreamResult::Dropped => break Ok(()),
2963 }
2964
2965 if let TransmitKind::Future = kind {
2966 break Ok(());
2967 }
2968 }
2969 }
2970 .map(move |result| {
2971 tls::get(|store| store.concurrent_state_mut()?.delete_transmit(id))?;
2972 result
2973 });
2974
2975 state.push_future(Box::pin(future));
2976 }
2977 WriteState::Dropped => {
2978 let reader = transmit.read_handle;
2979 self.0.host_drop_reader(reader, kind)?;
2980 }
2981 }
2982 Ok(())
2983 }
2984}
2985
2986async fn write<D: 'static, P: Send + 'static, T: func::Lower + 'static, B: WriteBuffer<T>>(
2987 token: StoreToken<D>,
2988 id: TableId<TransmitState>,
2989 pair: Arc<LockedState<(P, B)>>,
2990 kind: TransmitKind,
2991) -> Result<()> {
2992 let (read, guest_offset) = tls::get(|store| {
2993 let transmit = store.concurrent_state_mut()?.get_mut(id)?;
2994
2995 let guest_offset = if let &WriteState::HostReady { guest_offset, .. } = &transmit.write {
2996 Some(guest_offset)
2997 } else {
2998 None
2999 };
3000
3001 crate::error::Ok((
3002 mem::replace(&mut transmit.read, ReadState::Open),
3003 guest_offset,
3004 ))
3005 })?;
3006
3007 match read {
3008 ReadState::GuestReady {
3009 ty,
3010 flat_abi,
3011 options,
3012 address,
3013 count,
3014 handle,
3015 instance,
3016 caller_instance,
3017 caller_thread,
3018 } => {
3019 let guest_offset = match guest_offset {
3020 Some(i) => i,
3021 None => bail_bug!("guest_offset should be present if ready"),
3022 };
3023
3024 if let TransmitKind::Future = kind {
3025 tls::get(|store| {
3026 store.concurrent_state_mut()?.get_mut(id)?.done = true;
3027 crate::error::Ok(())
3028 })?;
3029 }
3030
3031 let old_remaining = pair.with(|p| p.1.remaining().len())?;
3032 let accept = {
3033 let pair = pair.clone();
3034 move |mut store: StoreContextMut<D>| {
3035 let mut state = pair.take()?;
3036 lower::<T, B, D>(
3037 store.as_context_mut(),
3038 instance,
3039 caller_thread,
3040 options,
3041 ty,
3042 address + (T::SIZE32 * guest_offset.as_usize()),
3043 count.as_usize() - guest_offset.as_usize(),
3044 &mut state.1,
3045 )?;
3046 crate::error::Ok(())
3047 }
3048 };
3049
3050 if guest_offset < count {
3051 if T::MAY_REQUIRE_REALLOC {
3052 let (tx, rx) = oneshot::channel();
3057 tls::get(move |store| {
3058 store
3059 .concurrent_state_mut()?
3060 .push_high_priority(WorkItem::WorkerFunction(AlwaysMut::new(
3061 Box::new(move |store| {
3062 _ = tx.send(accept(token.as_context_mut(store))?);
3063 Ok(())
3064 }),
3065 )));
3066 crate::error::Ok(())
3067 })?;
3068 match rx.await {
3069 Ok(r) => r,
3070 Err(oneshot::Canceled) => bail_bug!("work cancelled"),
3071 }
3072 } else {
3073 tls::get(|store| accept(token.as_context_mut(store)))?
3078 }
3079 }
3080
3081 tls::get(|store| {
3082 let count = old_remaining - pair.with(|p| p.1.remaining().len())?;
3083
3084 let transmit = store.concurrent_state_mut()?.get_mut(id)?;
3085
3086 let WriteState::HostReady { guest_offset, .. } = &mut transmit.write else {
3087 bail_bug!("expected WriteState::HostReady")
3088 };
3089
3090 guest_offset.inc(count)?;
3091
3092 transmit.read = ReadState::GuestReady {
3093 ty,
3094 flat_abi,
3095 options,
3096 address,
3097 count: ItemCount::new_usize(count)?,
3098 handle,
3099 instance,
3100 caller_instance,
3101 caller_thread,
3102 };
3103
3104 crate::error::Ok(())
3105 })?;
3106
3107 Ok(())
3108 }
3109
3110 ReadState::HostToHost {
3111 accept,
3112 mut buffer,
3113 limit,
3114 } => {
3115 let mut state = StreamResult::Completed;
3116 let mut position = 0;
3117
3118 while !matches!(state, StreamResult::Dropped) && position < limit {
3119 let mut slice_buffer = SliceBuffer::new(buffer, position, limit);
3120 state = accept(&mut UntypedWriteBuffer::new(&mut slice_buffer)).await?;
3121 (buffer, position, _) = slice_buffer.into_parts();
3122 }
3123
3124 {
3125 let mut pair = pair.take()?;
3126 let (_, buffer) = &mut *pair;
3127
3128 while !(matches!(state, StreamResult::Dropped) || buffer.remaining().is_empty()) {
3129 state = accept(&mut UntypedWriteBuffer::new(buffer)).await?;
3130 }
3131 }
3132
3133 tls::get(|store| {
3134 store.concurrent_state_mut()?.get_mut(id)?.read = match state {
3135 StreamResult::Dropped => ReadState::Dropped,
3136 StreamResult::Completed | StreamResult::Cancelled => ReadState::HostToHost {
3137 accept,
3138 buffer,
3139 limit: 0,
3140 },
3141 };
3142
3143 crate::error::Ok(())
3144 })?;
3145 Ok(())
3146 }
3147
3148 _ => bail_bug!("unexpected read state"),
3149 }
3150}
3151
3152impl Instance {
3153 fn consume(
3156 self,
3157 store: &mut dyn VMStore,
3158 kind: TransmitKind,
3159 transmit_id: TableId<TransmitState>,
3160 consume: PollStream,
3161 guest_offset: ItemCount,
3162 cancel: bool,
3163 ) -> Result<ReturnCode> {
3164 let mut future = consume();
3165 store.concurrent_state_mut()?.get_mut(transmit_id)?.read = ReadState::HostReady {
3166 consume,
3167 guest_offset,
3168 cancel,
3169 cancel_waker: None,
3170 };
3171 let poll = tls::set(store, || {
3172 future
3173 .as_mut()
3174 .poll(&mut Context::from_waker(&Waker::noop()))
3175 });
3176
3177 Ok(match poll {
3178 Poll::Ready(state) => {
3179 let transmit = store.concurrent_state_mut()?.get_mut(transmit_id)?;
3180 let code = settle_host_read(transmit, kind, state?)?;
3181 transmit.write = WriteState::Open;
3182 code
3183 }
3184 Poll::Pending => {
3185 store.pipe_from_guest(kind, transmit_id, future);
3186 ReturnCode::Blocked
3187 }
3188 })
3189 }
3190
3191 fn produce(
3194 self,
3195 store: &mut dyn VMStore,
3196 kind: TransmitKind,
3197 transmit_id: TableId<TransmitState>,
3198 produce: PollStream,
3199 try_into: TryInto,
3200 guest_offset: ItemCount,
3201 cancel: bool,
3202 ) -> Result<ReturnCode> {
3203 let mut future = produce();
3204 store.concurrent_state_mut()?.get_mut(transmit_id)?.write = WriteState::HostReady {
3205 produce,
3206 try_into,
3207 guest_offset,
3208 cancel,
3209 cancel_waker: None,
3210 };
3211 let poll = tls::set(store, || {
3212 future
3213 .as_mut()
3214 .poll(&mut Context::from_waker(&Waker::noop()))
3215 });
3216
3217 Ok(match poll {
3218 Poll::Ready(state) => {
3219 let transmit = store.concurrent_state_mut()?.get_mut(transmit_id)?;
3220 let code = settle_host_write(transmit, kind, state?)?;
3221 transmit.read = ReadState::Open;
3222 code
3223 }
3224 Poll::Pending => {
3225 store.pipe_to_guest(kind, transmit_id, future);
3226 ReturnCode::Blocked
3227 }
3228 })
3229 }
3230
3231 pub(super) fn guest_drop_writable(
3233 self,
3234 store: &mut StoreOpaque,
3235 ty: TransmitIndex,
3236 writer: u32,
3237 ) -> Result<()> {
3238 let table = self.id().get_mut(store).table_for_transmit(ty);
3239 let (transmit_rep, is_done) = match ty {
3240 TransmitIndex::Future(ty) => table.future_remove_writable(ty, writer)?,
3241 TransmitIndex::Stream(ty) => (table.stream_remove_writable(ty, writer)?, false),
3242 };
3243
3244 let id = TableId::<TransmitHandle>::new(transmit_rep);
3245 log::trace!("guest_drop_writable: drop writer {id:?}");
3246 match ty {
3247 TransmitIndex::Stream(_) => store.host_drop_writer(id, None),
3248 TransmitIndex::Future(_) => store.host_drop_writer(
3249 id,
3250 if is_done {
3251 None
3252 } else {
3253 Some(|| {
3254 Err(format_err!(
3255 "cannot drop future write end without first writing a value"
3256 ))
3257 })
3258 },
3259 ),
3260 }
3261 }
3262
3263 fn copy<T: 'static>(
3266 store: StoreContextMut<T>,
3267 flat_abi: Option<FlatAbi>,
3268 write_instance: Instance,
3269 write_caller_instance: RuntimeComponentInstanceIndex,
3270 write_ty: TransmitIndex,
3271 write_options: OptionsIndex,
3272 write_address: usize,
3273 read_instance: Instance,
3274 read_caller_instance: RuntimeComponentInstanceIndex,
3275 read_caller_thread: QualifiedThreadId,
3276 read_ty: TransmitIndex,
3277 read_options: OptionsIndex,
3278 read_address: usize,
3279 count: ItemCount,
3280 rep: u32,
3281 ) -> Result<()> {
3282 let (write_component, store) = write_instance.component_and_store_mut(store.0);
3283 let (read_component, mut store) = read_instance.component_and_store_mut(store);
3284 let write_types = write_component.types();
3285 let read_types = read_component.types();
3286 let count = count.as_usize();
3287
3288 let write_payload_ty = write_ty.payload(write_types);
3291 let write_abi = match write_payload_ty {
3292 Some(ty) => write_types.canonical_abi(ty),
3293 None => &CanonicalAbiInfo::ZERO,
3294 };
3295 let write_length_in_bytes = match flat_abi {
3296 Some(abi) => usize::try_from(abi.size)? * count,
3297 None => usize::try_from(write_abi.size32)? * count,
3298 };
3299 if write_length_in_bytes > 0 {
3300 if write_address % usize::try_from(write_abi.align32)? != 0 {
3301 bail!("write pointer not aligned");
3302 }
3303 write_instance
3304 .options_memory(store, write_options)
3305 .get(write_address..)
3306 .and_then(|b| b.get(..write_length_in_bytes))
3307 .ok_or_else(|| crate::format_err!("write pointer out of bounds"))?;
3308 }
3309
3310 let read_payload_ty = read_ty.payload(read_types);
3311 let read_abi = match read_payload_ty {
3312 Some(ty) => read_types.canonical_abi(ty),
3313 None => &CanonicalAbiInfo::ZERO,
3314 };
3315 let read_length_in_bytes = match flat_abi {
3316 Some(abi) => usize::try_from(abi.size)? * count,
3317 None => usize::try_from(read_abi.size32)? * count,
3318 };
3319 if read_length_in_bytes > 0 {
3320 if read_address % usize::try_from(read_abi.align32)? != 0 {
3321 bail!("read pointer not aligned");
3322 }
3323 read_instance
3324 .options_memory(store, read_options)
3325 .get(read_address..)
3326 .and_then(|b| b.get(..read_length_in_bytes))
3327 .ok_or_else(|| crate::format_err!("read pointer out of bounds"))?;
3328 }
3329
3330 if write_caller_instance == read_caller_instance
3331 && !allow_intra_component_read_write(write_payload_ty)
3332 {
3333 bail!(
3334 "cannot read from and write to intra-component future/stream with non-numeric payload"
3335 )
3336 }
3337
3338 match (write_ty, read_ty) {
3339 (TransmitIndex::Future(_), TransmitIndex::Future(_)) => {
3340 if count != 1 {
3341 bail_bug!("futures can only send 1 item");
3342 }
3343
3344 let val = write_payload_ty
3345 .map(|ty| {
3346 let lift = &mut LiftContext::new(store, write_options, write_instance)?;
3347 let bytes = &lift.memory()[write_address..][..write_length_in_bytes];
3348 Val::load(lift, *ty, bytes)
3349 })
3350 .transpose()?;
3351
3352 if let Some(val) = val {
3353 let old_thread = store.set_thread(read_caller_thread)?;
3357 let lower =
3358 &mut LowerContext::new(store.as_context_mut(), read_options, read_instance);
3359 let ptr = func::validate_inbounds_dynamic(
3360 read_abi,
3361 lower.as_slice_mut(),
3362 &ValRaw::u32(read_address.try_into()?),
3363 )?;
3364 let ty = match read_payload_ty {
3365 Some(ty) => ty,
3366 None => bail_bug!("expected read payload type to be present"),
3367 };
3368 val.store(lower, *ty, ptr)?;
3369 store.set_thread(old_thread)?;
3370 }
3371 }
3372 (TransmitIndex::Stream(_), TransmitIndex::Stream(_)) => {
3373 if write_length_in_bytes == 0 {
3374 return Ok(());
3375 }
3376 let write_payload_ty = match write_payload_ty {
3377 Some(ty) => ty,
3378 None => bail_bug!("expected write payload type to be present"),
3379 };
3380 let read_payload_ty = match read_payload_ty {
3381 Some(ty) => ty,
3382 None => bail_bug!("expected read payload type to be present"),
3383 };
3384 if flat_abi.is_some() {
3385 let store_opaque = store.store_opaque_mut();
3387
3388 assert_eq!(read_length_in_bytes, write_length_in_bytes);
3389
3390 if read_instance
3391 .options_memory(store_opaque, read_options)
3392 .as_ptr()
3393 == write_instance
3394 .options_memory(store_opaque, write_options)
3395 .as_ptr()
3396 {
3397 let memory = read_instance.options_memory_mut(store_opaque, read_options);
3398 memory.copy_within(
3399 write_address..write_address + write_length_in_bytes,
3400 read_address,
3401 );
3402 } else {
3403 let src = write_instance.options_memory(store_opaque, write_options)
3404 [write_address..][..write_length_in_bytes]
3405 .as_ptr();
3406 let dst = read_instance.options_memory_mut(store_opaque, read_options)
3407 [read_address..][..read_length_in_bytes]
3408 .as_mut_ptr();
3409
3410 unsafe {
3420 src.copy_to_nonoverlapping(dst, write_length_in_bytes);
3421 }
3422 }
3423 } else {
3424 let store_opaque = store.store_opaque_mut();
3425 let lift = &mut LiftContext::new(store_opaque, write_options, write_instance)?;
3426 let bytes = &lift.memory()[write_address..][..write_length_in_bytes];
3427 lift.consume_fuel_array(count, size_of::<Val>())?;
3428
3429 let values = (0..count)
3430 .map(|index| {
3431 let size = usize::try_from(write_abi.size32)?;
3432 Val::load(lift, *write_payload_ty, &bytes[(index * size)..][..size])
3433 })
3434 .collect::<Result<Vec<_>>>()?;
3435
3436 let id = TableId::<TransmitHandle>::new(rep);
3437 log::trace!("copy values {values:?} for {id:?}");
3438
3439 let old_thread = store.set_thread(read_caller_thread)?;
3443 let lower =
3444 &mut LowerContext::new(store.as_context_mut(), read_options, read_instance);
3445 let mut ptr = read_address;
3446 for value in values {
3447 value.store(lower, *read_payload_ty, ptr)?;
3448 ptr += usize::try_from(read_abi.size32)?;
3449 }
3450 store.set_thread(old_thread)?;
3451 }
3452 }
3453 _ => bail_bug!("mismatched transmit types in copy"),
3454 }
3455
3456 Ok(())
3457 }
3458
3459 fn check_bounds(
3460 self,
3461 store: &StoreOpaque,
3462 options: OptionsIndex,
3463 ty: TransmitIndex,
3464 address: usize,
3465 count: usize,
3466 ) -> Result<()> {
3467 let types = self.id().get(store).component().types();
3468 let size = usize::try_from(
3469 match ty {
3470 TransmitIndex::Future(ty) => types[types[ty].ty]
3471 .payload
3472 .map(|ty| types.canonical_abi(&ty).size32),
3473 TransmitIndex::Stream(ty) => types[types[ty].ty]
3474 .payload
3475 .map(|ty| types.canonical_abi(&ty).size32),
3476 }
3477 .unwrap_or(0),
3478 )?;
3479
3480 if count > 0 && size > 0 {
3481 self.options_memory(store, options)
3482 .get(address..)
3483 .and_then(|b| b.get(..size.checked_mul(count)?))
3484 .map(drop)
3485 .ok_or_else(|| crate::format_err!("read pointer out of bounds of memory"))
3486 } else {
3487 Ok(())
3488 }
3489 }
3490
3491 pub(super) fn guest_write<T: 'static>(
3493 self,
3494 mut store: StoreContextMut<T>,
3495 caller: RuntimeComponentInstanceIndex,
3496 ty: TransmitIndex,
3497 options: OptionsIndex,
3498 flat_abi: Option<FlatAbi>,
3499 handle: u32,
3500 address: u32,
3501 count: u32,
3502 ) -> Result<ReturnCode> {
3503 let count = ItemCount::new(count)?;
3504
3505 if !self.options(store.0, options).async_ {
3506 store.0.check_blocking()?;
3510 }
3511
3512 let address = usize::try_from(address)?;
3513 self.check_bounds(store.0, options, ty, address, count.as_usize())?;
3514 let (rep, state) = self.id().get_mut(store.0).get_mut_by_index(ty, handle)?;
3515 let TransmitLocalState::Write { done } = *state else {
3516 bail!(Trap::ConcurrentFutureStreamOp);
3517 };
3518
3519 if done {
3520 bail!("cannot write after being notified that the readable end dropped");
3521 }
3522
3523 *state = TransmitLocalState::Busy;
3524 let transmit_handle = TableId::<TransmitHandle>::new(rep);
3525 let concurrent_state = store.0.concurrent_state_mut()?;
3526 let transmit_id = concurrent_state.get_mut(transmit_handle)?.state;
3527 let transmit = concurrent_state.get_mut(transmit_id)?;
3528 log::trace!(
3529 "guest_write {count} to {transmit_handle:?} (handle {handle}; state {transmit_id:?}); {:?}",
3530 transmit.read
3531 );
3532
3533 if transmit.done {
3534 bail!("cannot write to future after previous write succeeded or readable end dropped");
3535 }
3536
3537 let new_state = if let ReadState::Dropped = &transmit.read {
3538 ReadState::Dropped
3539 } else {
3540 ReadState::Open
3541 };
3542
3543 let set_guest_ready = |me: &mut ConcurrentState| {
3544 let transmit = me.get_mut(transmit_id)?;
3545 if !matches!(&transmit.write, WriteState::Open) {
3546 bail_bug!("expected `WriteState::Open`; got `{:?}`", transmit.write);
3547 }
3548 transmit.write = WriteState::GuestReady {
3549 instance: self,
3550 caller,
3551 ty,
3552 flat_abi,
3553 options,
3554 address,
3555 count,
3556 handle,
3557 };
3558 Ok::<_, crate::Error>(())
3559 };
3560
3561 let mut result = match mem::replace(&mut transmit.read, new_state) {
3562 ReadState::GuestReady {
3563 ty: read_ty,
3564 flat_abi: read_flat_abi,
3565 options: read_options,
3566 address: read_address,
3567 count: read_count,
3568 handle: read_handle,
3569 instance: read_instance,
3570 caller_instance: read_caller_instance,
3571 caller_thread: read_caller_thread,
3572 } => {
3573 if flat_abi != read_flat_abi {
3574 bail_bug!("expected flat ABI calculations to be the same");
3575 }
3576
3577 if let TransmitIndex::Future(_) = ty {
3578 transmit.done = true;
3579 }
3580
3581 let write_complete = count == 0 || read_count > 0;
3603 let read_complete = count > 0;
3604 let read_buffer_remaining = count < read_count;
3605
3606 let read_handle_rep = transmit.read_handle.rep();
3607
3608 let count = count.min(read_count);
3609
3610 Instance::copy(
3611 store.as_context_mut(),
3612 flat_abi,
3613 self,
3614 caller,
3615 ty,
3616 options,
3617 address,
3618 read_instance,
3619 read_caller_instance,
3620 read_caller_thread,
3621 read_ty,
3622 read_options,
3623 read_address,
3624 count,
3625 rep,
3626 )?;
3627
3628 let instance = read_instance.id().get(store.0);
3629 let types = instance.component().types();
3630 let item_size = match read_ty.payload(types) {
3631 Some(ty) => usize::try_from(types.canonical_abi(ty).size32)?,
3632 None => 0,
3633 };
3634 let concurrent_state = store.0.concurrent_state_mut()?;
3635 if read_complete {
3636 let total = if let Some(Event::StreamRead {
3637 code: ReturnCode::Completed(old_total),
3638 ..
3639 }) = concurrent_state.take_event(read_handle_rep)?
3640 {
3641 count.add(old_total)?
3642 } else {
3643 count
3644 };
3645
3646 let code = ReturnCode::completed(ty.kind(), total);
3647
3648 concurrent_state.send_read_result(read_ty, transmit_id, read_handle, code)?;
3649 }
3650
3651 if read_buffer_remaining || (count == 0 && read_count == 0) {
3658 let transmit = concurrent_state.get_mut(transmit_id)?;
3659 transmit.read = ReadState::GuestReady {
3660 ty: read_ty,
3661 flat_abi: read_flat_abi,
3662 options: read_options,
3663 address: read_address + (count.as_usize() * item_size),
3664 count: read_count.sub(count)?,
3665 handle: read_handle,
3666 instance: read_instance,
3667 caller_instance: read_caller_instance,
3668 caller_thread: read_caller_thread,
3669 };
3670 }
3671
3672 if write_complete {
3673 ReturnCode::completed(ty.kind(), count)
3674 } else {
3675 set_guest_ready(concurrent_state)?;
3676 ReturnCode::Blocked
3677 }
3678 }
3679
3680 ReadState::HostReady {
3681 consume,
3682 guest_offset,
3683 cancel,
3684 cancel_waker,
3685 } => {
3686 if cancel_waker.is_some() {
3687 bail_bug!("expected cancel_waker to be none");
3688 }
3689 if cancel {
3690 bail_bug!("expected cancel to be false");
3691 }
3692 if guest_offset != 0 {
3693 bail_bug!("expected guest_offset to be 0");
3694 }
3695
3696 if let TransmitIndex::Future(_) = ty {
3697 transmit.done = true;
3698 }
3699
3700 set_guest_ready(concurrent_state)?;
3701 self.consume(
3702 store.0,
3703 ty.kind(),
3704 transmit_id,
3705 consume,
3706 ItemCount::ZERO,
3707 false,
3708 )?
3709 }
3710
3711 ReadState::HostToHost { .. } => bail_bug!("unexpected HostToHost"),
3712
3713 ReadState::Open => {
3714 set_guest_ready(concurrent_state)?;
3715 ReturnCode::Blocked
3716 }
3717
3718 ReadState::Dropped => {
3719 if let TransmitIndex::Future(_) = ty {
3720 transmit.done = true;
3721 }
3722
3723 ReturnCode::Dropped(ItemCount::ZERO)
3724 }
3725 };
3726
3727 if result == ReturnCode::Blocked && !self.options(store.0, options).async_ {
3728 result = self.wait_for_write(store.0, transmit_handle)?;
3729 }
3730
3731 if result != ReturnCode::Blocked {
3732 *self.id().get_mut(store.0).get_mut_by_index(ty, handle)?.1 =
3733 TransmitLocalState::Write {
3734 done: matches!(result, ReturnCode::Dropped(_)),
3735 };
3736 }
3737
3738 log::trace!(
3739 "guest_write result for {transmit_handle:?} (handle {handle}; state {transmit_id:?}): {result:?}",
3740 );
3741
3742 Ok(result)
3743 }
3744
3745 pub(super) fn guest_read<T: 'static>(
3747 self,
3748 mut store: StoreContextMut<T>,
3749 caller_instance: RuntimeComponentInstanceIndex,
3750 ty: TransmitIndex,
3751 options: OptionsIndex,
3752 flat_abi: Option<FlatAbi>,
3753 handle: u32,
3754 address: u32,
3755 count: u32,
3756 ) -> Result<ReturnCode> {
3757 let count = ItemCount::new(count)?;
3758
3759 if !self.options(store.0, options).async_ {
3760 store.0.check_blocking()?;
3764 }
3765
3766 let address = usize::try_from(address)?;
3767 self.check_bounds(store.0, options, ty, address, count.as_usize())?;
3768 let (rep, state) = self.id().get_mut(store.0).get_mut_by_index(ty, handle)?;
3769 let TransmitLocalState::Read { done } = *state else {
3770 bail!(Trap::ConcurrentFutureStreamOp);
3771 };
3772
3773 if done {
3774 bail!("cannot read after being notified that the writable end dropped");
3775 }
3776
3777 *state = TransmitLocalState::Busy;
3778 let transmit_handle = TableId::<TransmitHandle>::new(rep);
3779 let caller_thread = store.0.current_guest_thread()?;
3780 let concurrent_state = store.0.concurrent_state_mut()?;
3781 let transmit_id = concurrent_state.get_mut(transmit_handle)?.state;
3782 let transmit = concurrent_state.get_mut(transmit_id)?;
3783 log::trace!(
3784 "guest_read {count} from {transmit_handle:?} (handle {handle}; state {transmit_id:?}); {:?}",
3785 transmit.write
3786 );
3787
3788 if transmit.done {
3789 bail!("cannot read from future after previous read succeeded");
3790 }
3791
3792 let new_state = if let WriteState::Dropped = &transmit.write {
3793 WriteState::Dropped
3794 } else {
3795 WriteState::Open
3796 };
3797
3798 let set_guest_ready = |me: &mut ConcurrentState| {
3799 let transmit = me.get_mut(transmit_id)?;
3800 if !matches!(&transmit.read, ReadState::Open) {
3801 bail_bug!("expected `ReadState::Open`; got `{:?}`", transmit.read);
3802 }
3803 transmit.read = ReadState::GuestReady {
3804 ty,
3805 flat_abi,
3806 options,
3807 address,
3808 count,
3809 handle,
3810 instance: self,
3811 caller_instance,
3812 caller_thread,
3813 };
3814 Ok::<_, crate::Error>(())
3815 };
3816
3817 let mut result = match mem::replace(&mut transmit.write, new_state) {
3818 WriteState::GuestReady {
3819 instance: write_instance,
3820 ty: write_ty,
3821 flat_abi: write_flat_abi,
3822 options: write_options,
3823 address: write_address,
3824 count: write_count,
3825 handle: write_handle,
3826 caller: write_caller,
3827 } => {
3828 if flat_abi != write_flat_abi {
3829 bail_bug!("expected flat ABI calculations to be the same");
3830 }
3831
3832 if let TransmitIndex::Future(_) = ty {
3833 transmit.done = true;
3834 }
3835
3836 let write_handle_rep = transmit.write_handle.rep();
3837
3838 let write_complete = write_count == 0 || count > 0;
3843 let read_complete = write_count > 0;
3844 let write_buffer_remaining = count < write_count;
3845
3846 let count = count.min(write_count);
3847
3848 Instance::copy(
3849 store.as_context_mut(),
3850 flat_abi,
3851 write_instance,
3852 write_caller,
3853 write_ty,
3854 write_options,
3855 write_address,
3856 self,
3857 caller_instance,
3858 caller_thread,
3859 ty,
3860 options,
3861 address,
3862 count,
3863 rep,
3864 )?;
3865
3866 let instance = write_instance.id().get(store.0);
3867 let types = instance.component().types();
3868 let item_size = match write_ty.payload(types) {
3869 Some(ty) => usize::try_from(types.canonical_abi(ty).size32)?,
3870 None => 0,
3871 };
3872 let concurrent_state = store.0.concurrent_state_mut()?;
3873
3874 if write_complete {
3875 let total = if let Some(Event::StreamWrite {
3876 code: ReturnCode::Completed(old_total),
3877 ..
3878 }) = concurrent_state.take_event(write_handle_rep)?
3879 {
3880 count.add(old_total)?
3881 } else {
3882 count
3883 };
3884
3885 let code = ReturnCode::completed(ty.kind(), total);
3886
3887 concurrent_state.send_write_result(
3888 write_ty,
3889 transmit_id,
3890 write_handle,
3891 code,
3892 )?;
3893 }
3894
3895 if write_buffer_remaining {
3896 let transmit = concurrent_state.get_mut(transmit_id)?;
3897 transmit.write = WriteState::GuestReady {
3898 instance: write_instance,
3899 caller: write_caller,
3900 ty: write_ty,
3901 flat_abi: write_flat_abi,
3902 options: write_options,
3903 address: write_address + (count.as_usize() * item_size),
3904 count: write_count.sub(count)?,
3905 handle: write_handle,
3906 };
3907 }
3908
3909 if read_complete {
3910 ReturnCode::completed(ty.kind(), count)
3911 } else {
3912 set_guest_ready(concurrent_state)?;
3913 ReturnCode::Blocked
3914 }
3915 }
3916
3917 WriteState::HostReady {
3918 produce,
3919 try_into,
3920 guest_offset,
3921 cancel,
3922 cancel_waker,
3923 } => {
3924 if cancel_waker.is_some() {
3925 bail_bug!("expected cancel_waker to be none");
3926 }
3927 if cancel {
3928 bail_bug!("expected cancel to be false");
3929 }
3930 if guest_offset != 0 {
3931 bail_bug!("expected guest_offset to be 0");
3932 }
3933
3934 set_guest_ready(concurrent_state)?;
3935
3936 let code = self.produce(
3937 store.0,
3938 ty.kind(),
3939 transmit_id,
3940 produce,
3941 try_into,
3942 ItemCount::ZERO,
3943 false,
3944 )?;
3945
3946 if let (TransmitIndex::Future(_), ReturnCode::Completed(_)) = (ty, code) {
3947 store.0.concurrent_state_mut()?.get_mut(transmit_id)?.done = true;
3948 }
3949
3950 code
3951 }
3952
3953 WriteState::Open => {
3954 set_guest_ready(concurrent_state)?;
3955 ReturnCode::Blocked
3956 }
3957
3958 WriteState::Dropped => ReturnCode::Dropped(ItemCount::ZERO),
3959 };
3960
3961 if result == ReturnCode::Blocked && !self.options(store.0, options).async_ {
3962 result = self.wait_for_read(store.0, transmit_handle)?;
3963 }
3964
3965 if result != ReturnCode::Blocked {
3966 *self.id().get_mut(store.0).get_mut_by_index(ty, handle)?.1 =
3967 TransmitLocalState::Read {
3968 done: matches!(
3969 (result, ty),
3970 (ReturnCode::Dropped(_), TransmitIndex::Stream(_))
3971 ),
3972 };
3973 }
3974
3975 log::trace!(
3976 "guest_read result for {transmit_handle:?} (handle {handle}; state {transmit_id:?}): {result:?}",
3977 );
3978
3979 Ok(result)
3980 }
3981
3982 fn wait_for_write(
3983 self,
3984 store: &mut StoreOpaque,
3985 handle: TableId<TransmitHandle>,
3986 ) -> Result<ReturnCode> {
3987 let waitable = Waitable::Transmit(handle);
3988 store.wait_for_event(waitable)?;
3989 let event = waitable.take_event(store.concurrent_state_mut()?)?;
3990 if let Some(event @ (Event::StreamWrite { code, .. } | Event::FutureWrite { code, .. })) =
3991 event
3992 {
3993 waitable.on_delivery(store, self, event)?;
3994 Ok(code)
3995 } else {
3996 bail_bug!("expected either a stream or future write event")
3997 }
3998 }
3999
4000 fn cancel_write(
4002 self,
4003 store: &mut StoreOpaque,
4004 transmit_id: TableId<TransmitState>,
4005 async_: bool,
4006 ) -> Result<ReturnCode> {
4007 let state = store.concurrent_state_mut()?;
4008 let transmit = state.get_mut(transmit_id)?;
4009 log::trace!(
4010 "host_cancel_write state {transmit_id:?}; write state {:?} read state {:?}",
4011 transmit.read,
4012 transmit.write
4013 );
4014 let waitable = Waitable::Transmit(transmit.write_handle);
4015
4016 if !async_ {
4017 waitable.trap_if_in_waitable_set(state)?;
4018 }
4019
4020 let code = if let Some(event) = waitable.take_event(state)? {
4021 let (Event::FutureWrite { code, .. } | Event::StreamWrite { code, .. }) = event else {
4022 bail_bug!("expected either a stream or future write event")
4023 };
4024 waitable.on_delivery(store, self, event)?;
4025 match (code, event) {
4026 (ReturnCode::Completed(count), Event::StreamWrite { .. }) => {
4027 ReturnCode::Cancelled(count)
4028 }
4029 (ReturnCode::Dropped(_) | ReturnCode::Completed(_), _) => code,
4030 _ => bail_bug!("unexpected code/event combo"),
4031 }
4032 } else if let ReadState::HostReady {
4033 cancel,
4034 cancel_waker,
4035 ..
4036 } = &mut state.get_mut(transmit_id)?.read
4037 {
4038 *cancel = true;
4039 if let Some(waker) = cancel_waker.take() {
4040 waker.wake();
4041 }
4042
4043 if async_ {
4044 ReturnCode::Blocked
4045 } else {
4046 let handle = store
4047 .concurrent_state_mut()?
4048 .get_mut(transmit_id)?
4049 .write_handle;
4050 self.wait_for_write(store, handle)?
4051 }
4052 } else {
4053 ReturnCode::Cancelled(ItemCount::ZERO)
4054 };
4055
4056 if !matches!(code, ReturnCode::Blocked) {
4057 let transmit = store.concurrent_state_mut()?.get_mut(transmit_id)?;
4058
4059 match &transmit.write {
4060 WriteState::GuestReady { .. } => {
4061 transmit.write = WriteState::Open;
4062 }
4063 WriteState::HostReady { .. } => bail_bug!("support host write cancellation"),
4064 WriteState::Open | WriteState::Dropped => {}
4065 }
4066 }
4067
4068 log::trace!("cancelled write {transmit_id:?}: {code:?}");
4069
4070 Ok(code)
4071 }
4072
4073 fn wait_for_read(
4074 self,
4075 store: &mut StoreOpaque,
4076 handle: TableId<TransmitHandle>,
4077 ) -> Result<ReturnCode> {
4078 let waitable = Waitable::Transmit(handle);
4079 store.wait_for_event(waitable)?;
4080 let event = waitable.take_event(store.concurrent_state_mut()?)?;
4081 if let Some(event @ (Event::StreamRead { code, .. } | Event::FutureRead { code, .. })) =
4082 event
4083 {
4084 waitable.on_delivery(store, self, event)?;
4085 Ok(code)
4086 } else {
4087 bail_bug!("expected either a stream or future read event")
4088 }
4089 }
4090
4091 fn cancel_read(
4093 self,
4094 store: &mut StoreOpaque,
4095 transmit_id: TableId<TransmitState>,
4096 async_: bool,
4097 ) -> Result<ReturnCode> {
4098 let state = store.concurrent_state_mut()?;
4099 let transmit = state.get_mut(transmit_id)?;
4100 log::trace!(
4101 "host_cancel_read state {transmit_id:?}; read state {:?} write state {:?}",
4102 transmit.read,
4103 transmit.write
4104 );
4105
4106 let waitable = Waitable::Transmit(transmit.read_handle);
4107
4108 if !async_ {
4109 waitable.trap_if_in_waitable_set(state)?;
4110 }
4111
4112 let code = if let Some(event) = waitable.take_event(state)? {
4113 let (Event::FutureRead { code, .. } | Event::StreamRead { code, .. }) = event else {
4114 bail_bug!("expected either a stream or future read event")
4115 };
4116 waitable.on_delivery(store, self, event)?;
4117 match (code, event) {
4118 (ReturnCode::Completed(count), Event::StreamRead { .. }) => {
4119 ReturnCode::Cancelled(count)
4120 }
4121 (ReturnCode::Dropped(_) | ReturnCode::Completed(_), _) => code,
4122 _ => bail_bug!("unexpected code/event combo"),
4123 }
4124 } else if let WriteState::HostReady {
4125 cancel,
4126 cancel_waker,
4127 ..
4128 } = &mut state.get_mut(transmit_id)?.write
4129 {
4130 *cancel = true;
4131 if let Some(waker) = cancel_waker.take() {
4132 waker.wake();
4133 }
4134
4135 if async_ {
4136 ReturnCode::Blocked
4137 } else {
4138 let handle = store
4139 .concurrent_state_mut()?
4140 .get_mut(transmit_id)?
4141 .read_handle;
4142 self.wait_for_read(store, handle)?
4143 }
4144 } else {
4145 ReturnCode::Cancelled(ItemCount::ZERO)
4146 };
4147
4148 if !matches!(code, ReturnCode::Blocked) {
4149 let transmit = store.concurrent_state_mut()?.get_mut(transmit_id)?;
4150
4151 match &transmit.read {
4152 ReadState::GuestReady { .. } => {
4153 transmit.read = ReadState::Open;
4154 }
4155 ReadState::HostReady { .. } | ReadState::HostToHost { .. } => {
4156 bail_bug!("support host read cancellation")
4157 }
4158 ReadState::Open | ReadState::Dropped => {}
4159 }
4160 }
4161
4162 log::trace!("cancelled read {transmit_id:?}: {code:?}");
4163
4164 Ok(code)
4165 }
4166
4167 fn guest_cancel_write(
4169 self,
4170 store: &mut StoreOpaque,
4171 ty: TransmitIndex,
4172 async_: bool,
4173 writer: u32,
4174 ) -> Result<ReturnCode> {
4175 if !async_ {
4176 store.check_blocking()?;
4180 }
4181
4182 let (rep, state) =
4183 get_mut_by_index_from(self.id().get_mut(store).table_for_transmit(ty), ty, writer)?;
4184 let id = TableId::<TransmitHandle>::new(rep);
4185 log::trace!("guest cancel write {id:?} (handle {writer})");
4186 match state {
4187 TransmitLocalState::Write { .. } => {
4188 bail!("stream or future write cancelled when no write is pending")
4189 }
4190 TransmitLocalState::Read { .. } => {
4191 bail!("passed read end to `{{stream|future}}.cancel-write`")
4192 }
4193 TransmitLocalState::Busy => {}
4194 }
4195 let transmit_id = store.concurrent_state_mut()?.get_mut(id)?.state;
4196 let code = self.cancel_write(store, transmit_id, async_)?;
4197 if !matches!(code, ReturnCode::Blocked) {
4198 let state =
4199 get_mut_by_index_from(self.id().get_mut(store).table_for_transmit(ty), ty, writer)?
4200 .1;
4201 if let TransmitLocalState::Busy = state {
4202 *state = TransmitLocalState::Write { done: false };
4203 }
4204 }
4205 Ok(code)
4206 }
4207
4208 fn guest_cancel_read(
4210 self,
4211 store: &mut StoreOpaque,
4212 ty: TransmitIndex,
4213 async_: bool,
4214 reader: u32,
4215 ) -> Result<ReturnCode> {
4216 if !async_ {
4217 store.check_blocking()?;
4221 }
4222
4223 let (rep, state) =
4224 get_mut_by_index_from(self.id().get_mut(store).table_for_transmit(ty), ty, reader)?;
4225 let id = TableId::<TransmitHandle>::new(rep);
4226 log::trace!("guest cancel read {id:?} (handle {reader})");
4227 match state {
4228 TransmitLocalState::Read { .. } => {
4229 bail!("stream or future read cancelled when no read is pending")
4230 }
4231 TransmitLocalState::Write { .. } => {
4232 bail!("passed write end to `{{stream|future}}.cancel-read`")
4233 }
4234 TransmitLocalState::Busy => {}
4235 }
4236 let transmit_id = store.concurrent_state_mut()?.get_mut(id)?.state;
4237 let code = self.cancel_read(store, transmit_id, async_)?;
4238 if !matches!(code, ReturnCode::Blocked) {
4239 let state =
4240 get_mut_by_index_from(self.id().get_mut(store).table_for_transmit(ty), ty, reader)?
4241 .1;
4242 if let TransmitLocalState::Busy = state {
4243 *state = TransmitLocalState::Read { done: false };
4244 }
4245 }
4246 Ok(code)
4247 }
4248
4249 fn guest_drop_readable(
4251 self,
4252 store: &mut StoreOpaque,
4253 ty: TransmitIndex,
4254 reader: u32,
4255 ) -> Result<()> {
4256 let table = self.id().get_mut(store).table_for_transmit(ty);
4257 let (rep, _is_done) = match ty {
4258 TransmitIndex::Stream(ty) => table.stream_remove_readable(ty, reader)?,
4259 TransmitIndex::Future(ty) => table.future_remove_readable(ty, reader)?,
4260 };
4261 let kind = match ty {
4262 TransmitIndex::Stream(_) => TransmitKind::Stream,
4263 TransmitIndex::Future(_) => TransmitKind::Future,
4264 };
4265 let id = TableId::<TransmitHandle>::new(rep);
4266 log::trace!("guest_drop_readable: drop reader {id:?}");
4267 store.host_drop_reader(id, kind)
4268 }
4269
4270 pub(crate) fn error_context_new(
4272 self,
4273 store: &mut StoreOpaque,
4274 ty: TypeComponentLocalErrorContextTableIndex,
4275 options: OptionsIndex,
4276 debug_msg_address: u32,
4277 debug_msg_len: u32,
4278 ) -> Result<u32> {
4279 let lift_ctx = &mut LiftContext::new(store, options, self)?;
4280 let debug_msg = String::linear_lift_from_flat(
4281 lift_ctx,
4282 InterfaceType::String,
4283 &[ValRaw::u32(debug_msg_address), ValRaw::u32(debug_msg_len)],
4284 )?;
4285
4286 let err_ctx = ErrorContextState { debug_msg };
4288 let state = store.concurrent_state_mut()?;
4289 let table_id = state.push(err_ctx)?;
4290 let global_ref_count_idx =
4291 TypeComponentGlobalErrorContextTableIndex::from_u32(table_id.rep());
4292
4293 let _ = state
4295 .global_error_context_ref_counts
4296 .insert(global_ref_count_idx, GlobalErrorContextRefCount(1));
4297
4298 let local_idx = self
4305 .id()
4306 .get_mut(store)
4307 .table_for_error_context(ty)
4308 .error_context_insert(table_id.rep())?;
4309
4310 Ok(local_idx)
4311 }
4312
4313 pub(super) fn error_context_debug_message<T>(
4315 self,
4316 store: StoreContextMut<T>,
4317 ty: TypeComponentLocalErrorContextTableIndex,
4318 options: OptionsIndex,
4319 err_ctx_handle: u32,
4320 debug_msg_address: u32,
4321 ) -> Result<()> {
4322 let handle_table_id_rep = self
4324 .id()
4325 .get_mut(store.0)
4326 .table_for_error_context(ty)
4327 .error_context_rep(err_ctx_handle)?;
4328
4329 let state = store.0.concurrent_state_mut()?;
4330 let ErrorContextState { debug_msg } =
4332 state.get_mut(TableId::<ErrorContextState>::new(handle_table_id_rep))?;
4333 let debug_msg = debug_msg.clone();
4334
4335 let lower_cx = &mut LowerContext::new(store, options, self);
4336 let debug_msg_address = usize::try_from(debug_msg_address)?;
4337 let offset = lower_cx
4343 .as_slice_mut()
4344 .get(debug_msg_address..)
4345 .and_then(|b| b.get(..8))
4346 .map(|_| debug_msg_address)
4347 .ok_or_else(|| crate::format_err!("invalid debug message pointer: out of bounds"))?;
4348 debug_msg
4349 .as_str()
4350 .linear_lower_to_memory(lower_cx, InterfaceType::String, offset)?;
4351
4352 Ok(())
4353 }
4354
4355 pub(crate) fn future_cancel_read(
4357 self,
4358 store: &mut StoreOpaque,
4359 ty: TypeFutureTableIndex,
4360 async_: bool,
4361 reader: u32,
4362 ) -> Result<u32> {
4363 self.guest_cancel_read(store, TransmitIndex::Future(ty), async_, reader)
4364 .map(|v| v.encode())
4365 }
4366
4367 pub(crate) fn future_cancel_write(
4369 self,
4370 store: &mut StoreOpaque,
4371 ty: TypeFutureTableIndex,
4372 async_: bool,
4373 writer: u32,
4374 ) -> Result<u32> {
4375 self.guest_cancel_write(store, TransmitIndex::Future(ty), async_, writer)
4376 .map(|v| v.encode())
4377 }
4378
4379 pub(crate) fn stream_cancel_read(
4381 self,
4382 store: &mut StoreOpaque,
4383 ty: TypeStreamTableIndex,
4384 async_: bool,
4385 reader: u32,
4386 ) -> Result<u32> {
4387 self.guest_cancel_read(store, TransmitIndex::Stream(ty), async_, reader)
4388 .map(|v| v.encode())
4389 }
4390
4391 pub(crate) fn stream_cancel_write(
4393 self,
4394 store: &mut StoreOpaque,
4395 ty: TypeStreamTableIndex,
4396 async_: bool,
4397 writer: u32,
4398 ) -> Result<u32> {
4399 self.guest_cancel_write(store, TransmitIndex::Stream(ty), async_, writer)
4400 .map(|v| v.encode())
4401 }
4402
4403 pub(crate) fn future_drop_readable(
4405 self,
4406 store: &mut StoreOpaque,
4407 ty: TypeFutureTableIndex,
4408 reader: u32,
4409 ) -> Result<()> {
4410 self.guest_drop_readable(store, TransmitIndex::Future(ty), reader)
4411 }
4412
4413 pub(crate) fn stream_drop_readable(
4415 self,
4416 store: &mut StoreOpaque,
4417 ty: TypeStreamTableIndex,
4418 reader: u32,
4419 ) -> Result<()> {
4420 self.guest_drop_readable(store, TransmitIndex::Stream(ty), reader)
4421 }
4422
4423 fn guest_new(self, store: &mut StoreOpaque, ty: TransmitIndex) -> Result<ResourcePair> {
4427 let (write, read) = store
4428 .concurrent_state_mut()?
4429 .new_transmit(TransmitOrigin::guest(self.id().instance(), ty))?;
4430
4431 let table = self.id().get_mut(store).table_for_transmit(ty);
4432 let (read_handle, write_handle) = match ty {
4433 TransmitIndex::Future(ty) => (
4434 table.future_insert_read(ty, read.rep())?,
4435 table.future_insert_write(ty, write.rep())?,
4436 ),
4437 TransmitIndex::Stream(ty) => (
4438 table.stream_insert_read(ty, read.rep())?,
4439 table.stream_insert_write(ty, write.rep())?,
4440 ),
4441 };
4442
4443 let state = store.concurrent_state_mut()?;
4444 state.get_mut(read)?.common.handle = Some(read_handle);
4445 state.get_mut(write)?.common.handle = Some(write_handle);
4446
4447 Ok(ResourcePair {
4448 write: write_handle,
4449 read: read_handle,
4450 })
4451 }
4452
4453 pub(crate) fn error_context_drop(
4455 self,
4456 store: &mut StoreOpaque,
4457 ty: TypeComponentLocalErrorContextTableIndex,
4458 error_context: u32,
4459 ) -> Result<()> {
4460 let instance = self.id().get_mut(store);
4461
4462 let local_handle_table = instance.table_for_error_context(ty);
4463
4464 let rep = local_handle_table.error_context_drop(error_context)?;
4465
4466 let global_ref_count_idx = TypeComponentGlobalErrorContextTableIndex::from_u32(rep);
4467
4468 let state = store.concurrent_state_mut()?;
4469 let Some(GlobalErrorContextRefCount(global_ref_count)) = state
4470 .global_error_context_ref_counts
4471 .get_mut(&global_ref_count_idx)
4472 else {
4473 bail_bug!("retrieve concurrent state for error context during drop")
4474 };
4475
4476 if *global_ref_count < 1 {
4478 bail_bug!("ref count unexpectedly zero");
4479 }
4480 *global_ref_count -= 1;
4481 if *global_ref_count == 0 {
4482 state
4483 .global_error_context_ref_counts
4484 .remove(&global_ref_count_idx);
4485
4486 state
4487 .delete(TableId::<ErrorContextState>::new(rep))
4488 .context("deleting component-global error context data")?;
4489 }
4490
4491 Ok(())
4492 }
4493
4494 fn guest_transfer(
4497 self,
4498 store: &mut StoreOpaque,
4499 src_idx: u32,
4500 src: TransmitIndex,
4501 dst: TransmitIndex,
4502 ) -> Result<u32> {
4503 let id = self.lift_index_to_transmit(store, src, src_idx)?;
4504 self.lower_transmit_to_index(store, dst, id)
4505 }
4506
4507 fn lift_index_to_transmit(
4508 self,
4509 store: &mut StoreOpaque,
4510 ty: TransmitIndex,
4511 src_idx: u32,
4512 ) -> Result<TableId<TransmitHandle>> {
4513 let (state, _, _, instance) = store.lift_context_parts(self);
4514 lift_index_to_transmit(instance, state.concurrent_state_mut(), ty, src_idx)
4515 }
4516
4517 fn lower_transmit_to_index(
4518 self,
4519 store: &mut StoreOpaque,
4520 ty: TransmitIndex,
4521 id: TableId<TransmitHandle>,
4522 ) -> Result<u32> {
4523 let (state, _, _, instance) = store.lift_context_parts(self);
4524 lower_transmit_to_index(instance, state.concurrent_state_mut(), ty, id)
4525 }
4526
4527 pub(crate) fn future_new(
4529 self,
4530 store: &mut StoreOpaque,
4531 ty: TypeFutureTableIndex,
4532 ) -> Result<ResourcePair> {
4533 self.guest_new(store, TransmitIndex::Future(ty))
4534 }
4535
4536 pub(crate) fn stream_new(
4538 self,
4539 store: &mut StoreOpaque,
4540 ty: TypeStreamTableIndex,
4541 ) -> Result<ResourcePair> {
4542 self.guest_new(store, TransmitIndex::Stream(ty))
4543 }
4544
4545 pub(crate) fn future_transfer(
4548 self,
4549 store: &mut StoreOpaque,
4550 src_idx: u32,
4551 src: TypeFutureTableIndex,
4552 dst: TypeFutureTableIndex,
4553 ) -> Result<u32> {
4554 self.guest_transfer(
4555 store,
4556 src_idx,
4557 TransmitIndex::Future(src),
4558 TransmitIndex::Future(dst),
4559 )
4560 }
4561
4562 pub(crate) fn stream_transfer(
4565 self,
4566 store: &mut StoreOpaque,
4567 src_idx: u32,
4568 src: TypeStreamTableIndex,
4569 dst: TypeStreamTableIndex,
4570 ) -> Result<u32> {
4571 self.guest_transfer(
4572 store,
4573 src_idx,
4574 TransmitIndex::Stream(src),
4575 TransmitIndex::Stream(dst),
4576 )
4577 }
4578
4579 pub(crate) fn error_context_transfer(
4581 self,
4582 store: &mut StoreOpaque,
4583 src_idx: u32,
4584 src: TypeComponentLocalErrorContextTableIndex,
4585 dst: TypeComponentLocalErrorContextTableIndex,
4586 ) -> Result<u32> {
4587 let mut instance = self.id().get_mut(store);
4588 let rep = instance
4589 .as_mut()
4590 .table_for_error_context(src)
4591 .error_context_rep(src_idx)?;
4592 let dst_idx = instance
4593 .table_for_error_context(dst)
4594 .error_context_insert(rep)?;
4595
4596 let global_ref_count = store
4600 .concurrent_state_mut()?
4601 .global_error_context_ref_counts
4602 .get_mut(&TypeComponentGlobalErrorContextTableIndex::from_u32(rep))
4603 .context("global ref count present for existing (sub)component error context")?;
4604
4605 global_ref_count.0 = global_ref_count
4606 .0
4607 .checked_add(1)
4608 .ok_or_else(|| format_err!(Trap::ReferenceCountOverflow))?;
4609
4610 Ok(dst_idx)
4611 }
4612}
4613
4614fn lift_index_to_transmit(
4620 instance: Pin<&mut ComponentInstance>,
4621 concurrent_state: &mut ConcurrentState,
4622 ty: TransmitIndex,
4623 src_idx: u32,
4624) -> Result<TableId<TransmitHandle>> {
4625 let handle_table = instance.table_for_transmit(ty);
4626 let (rep, is_done) = match ty {
4627 TransmitIndex::Future(idx) => handle_table.future_remove_readable(idx, src_idx)?,
4628 TransmitIndex::Stream(idx) => handle_table.stream_remove_readable(idx, src_idx)?,
4629 };
4630 let desc = match ty {
4631 TransmitIndex::Future(_) => "future",
4632 TransmitIndex::Stream(_) => "stream",
4633 };
4634 if is_done {
4635 bail!("cannot lift {desc} after being notified that the writable end dropped");
4636 }
4637 let id = TableId::<TransmitHandle>::new(rep);
4638 let future = concurrent_state.get_mut(id)?;
4639 if future.common.set.is_some() {
4640 bail!("cannot lift {desc} while it's in a waitable set");
4641 }
4642 future.common.handle = None;
4643
4644 let state = future.state;
4645 if concurrent_state.get_mut(state)?.done {
4646 bail!("cannot lift {desc} after previous read succeeded");
4647 }
4648
4649 Ok(id)
4650}
4651
4652fn lower_transmit_to_index(
4655 instance: Pin<&mut ComponentInstance>,
4656 concurrent_state: &mut ConcurrentState,
4657 ty: TransmitIndex,
4658 id: TableId<TransmitHandle>,
4659) -> Result<u32> {
4660 let state = concurrent_state.get_mut(id)?.state;
4661 debug_assert_eq!(concurrent_state.get_mut(state)?.read_handle, id);
4662 let handle_table = instance.table_for_transmit(ty);
4663 let handle = match ty {
4664 TransmitIndex::Future(idx) => handle_table.future_insert_read(idx, id.rep()),
4665 TransmitIndex::Stream(idx) => handle_table.stream_insert_read(idx, id.rep()),
4666 }?;
4667 concurrent_state.get_mut(id)?.common.handle = Some(handle);
4668 Ok(handle)
4669}
4670
4671impl ComponentInstance {
4672 fn table_for_transmit(self: Pin<&mut Self>, ty: TransmitIndex) -> &mut HandleTable {
4673 let (states, types) = self.instance_states();
4674 let runtime_instance = match ty {
4675 TransmitIndex::Stream(ty) => types[ty].instance,
4676 TransmitIndex::Future(ty) => types[ty].instance,
4677 };
4678 states[runtime_instance].handle_table()
4679 }
4680
4681 fn table_for_error_context(
4682 self: Pin<&mut Self>,
4683 ty: TypeComponentLocalErrorContextTableIndex,
4684 ) -> &mut HandleTable {
4685 let (states, types) = self.instance_states();
4686 let runtime_instance = types[ty].instance;
4687 states[runtime_instance].handle_table()
4688 }
4689
4690 fn get_mut_by_index(
4691 self: Pin<&mut Self>,
4692 ty: TransmitIndex,
4693 index: u32,
4694 ) -> Result<(u32, &mut TransmitLocalState)> {
4695 get_mut_by_index_from(self.table_for_transmit(ty), ty, index)
4696 }
4697}
4698
4699impl ConcurrentState {
4700 fn send_write_result(
4701 &mut self,
4702 ty: TransmitIndex,
4703 id: TableId<TransmitState>,
4704 handle: u32,
4705 code: ReturnCode,
4706 ) -> Result<()> {
4707 let write_handle = self.get_mut(id)?.write_handle.rep();
4708 self.set_event(
4709 write_handle,
4710 match ty {
4711 TransmitIndex::Future(ty) => Event::FutureWrite {
4712 code,
4713 pending: Some((ty, handle)),
4714 },
4715 TransmitIndex::Stream(ty) => Event::StreamWrite {
4716 code,
4717 pending: Some((ty, handle)),
4718 },
4719 },
4720 )
4721 }
4722
4723 fn send_read_result(
4724 &mut self,
4725 ty: TransmitIndex,
4726 id: TableId<TransmitState>,
4727 handle: u32,
4728 code: ReturnCode,
4729 ) -> Result<()> {
4730 let read_handle = self.get_mut(id)?.read_handle.rep();
4731 self.set_event(
4732 read_handle,
4733 match ty {
4734 TransmitIndex::Future(ty) => Event::FutureRead {
4735 code,
4736 pending: Some((ty, handle)),
4737 },
4738 TransmitIndex::Stream(ty) => Event::StreamRead {
4739 code,
4740 pending: Some((ty, handle)),
4741 },
4742 },
4743 )
4744 }
4745
4746 fn take_event(&mut self, waitable: u32) -> Result<Option<Event>> {
4747 Waitable::Transmit(TableId::<TransmitHandle>::new(waitable)).take_event(self)
4748 }
4749
4750 fn set_event(&mut self, waitable: u32, event: Event) -> Result<()> {
4751 Waitable::Transmit(TableId::<TransmitHandle>::new(waitable)).set_event(self, Some(event))
4752 }
4753
4754 fn update_event(&mut self, waitable: u32, event: Event) -> Result<()> {
4765 let waitable = Waitable::Transmit(TableId::<TransmitHandle>::new(waitable));
4766
4767 fn update_code(old: ReturnCode, new: ReturnCode) -> Result<ReturnCode> {
4768 let (ReturnCode::Completed(count)
4769 | ReturnCode::Dropped(count)
4770 | ReturnCode::Cancelled(count)) = old
4771 else {
4772 bail_bug!("unexpected old return code")
4773 };
4774
4775 Ok(match new {
4776 ReturnCode::Dropped(ItemCount::ZERO) => ReturnCode::Dropped(count),
4777 ReturnCode::Cancelled(ItemCount::ZERO) => ReturnCode::Cancelled(count),
4778 _ => bail_bug!("unexpected new return code"),
4779 })
4780 }
4781
4782 let event = match (waitable.take_event(self)?, event) {
4783 (None, _) => event,
4784 (Some(old @ Event::FutureWrite { .. }), Event::FutureWrite { .. }) => old,
4785 (Some(old @ Event::FutureRead { .. }), Event::FutureRead { .. }) => old,
4786 (
4787 Some(Event::StreamWrite {
4788 code: old_code,
4789 pending: old_pending,
4790 }),
4791 Event::StreamWrite { code, pending },
4792 ) => Event::StreamWrite {
4793 code: update_code(old_code, code)?,
4794 pending: old_pending.or(pending),
4795 },
4796 (
4797 Some(Event::StreamRead {
4798 code: old_code,
4799 pending: old_pending,
4800 }),
4801 Event::StreamRead { code, pending },
4802 ) => Event::StreamRead {
4803 code: update_code(old_code, code)?,
4804 pending: old_pending.or(pending),
4805 },
4806 _ => bail_bug!("unexpected event combination"),
4807 };
4808
4809 waitable.set_event(self, Some(event))
4810 }
4811
4812 fn new_transmit(
4815 &mut self,
4816 origin: TransmitOrigin,
4817 ) -> Result<(TableId<TransmitHandle>, TableId<TransmitHandle>)> {
4818 let state_id = self.push(TransmitState::new(origin))?;
4819
4820 let write = self.push(TransmitHandle::new(state_id))?;
4821 let read = self.push(TransmitHandle::new(state_id))?;
4822
4823 let state = self.get_mut(state_id)?;
4824 state.write_handle = write;
4825 state.read_handle = read;
4826
4827 log::trace!("new transmit: state {state_id:?}; write {write:?}; read {read:?}",);
4828
4829 Ok((write, read))
4830 }
4831
4832 fn delete_transmit(&mut self, state_id: TableId<TransmitState>) -> Result<()> {
4834 let state = self.delete(state_id)?;
4835 self.delete(state.write_handle)?;
4836 self.delete(state.read_handle)?;
4837
4838 log::trace!(
4839 "delete transmit: state {state_id:?}; write {:?}; read {:?}",
4840 state.write_handle,
4841 state.read_handle,
4842 );
4843
4844 Ok(())
4845 }
4846}
4847
4848pub(crate) struct ResourcePair {
4849 pub(crate) write: u32,
4850 pub(crate) read: u32,
4851}
4852
4853impl Waitable {
4854 pub(super) fn on_delivery(
4857 &self,
4858 store: &mut StoreOpaque,
4859 instance: Instance,
4860 event: Event,
4861 ) -> Result<()> {
4862 let instance = instance.id().get_mut(store);
4863 let (rep, state, code) = match event {
4864 Event::FutureRead {
4865 pending: Some((ty, handle)),
4866 code,
4867 }
4868 | Event::FutureWrite {
4869 pending: Some((ty, handle)),
4870 code,
4871 } => {
4872 let runtime_instance = instance.component().types()[ty].instance;
4873 let (rep, state) = instance.instance_states().0[runtime_instance]
4874 .handle_table()
4875 .future_rep(ty, handle)?;
4876 (rep, state, code)
4877 }
4878 Event::StreamRead {
4879 pending: Some((ty, handle)),
4880 code,
4881 }
4882 | Event::StreamWrite {
4883 pending: Some((ty, handle)),
4884 code,
4885 } => {
4886 let runtime_instance = instance.component().types()[ty].instance;
4887 let (rep, state) = instance.instance_states().0[runtime_instance]
4888 .handle_table()
4889 .stream_rep(ty, handle)?;
4890 (rep, state, code)
4891 }
4892 _ => return Ok(()),
4893 };
4894 if rep != self.rep() {
4895 bail_bug!("unexpected rep mismatch");
4896 }
4897 if *state != TransmitLocalState::Busy {
4898 bail_bug!("expected state to be busy");
4899 }
4900 let done = matches!(code, ReturnCode::Dropped(_));
4901 *state = match event {
4902 Event::FutureRead { .. } | Event::StreamRead { .. } => {
4903 TransmitLocalState::Read { done }
4904 }
4905 Event::FutureWrite { .. } | Event::StreamWrite { .. } => {
4906 TransmitLocalState::Write { done }
4907 }
4908 _ => bail_bug!("unexpected event for stream"),
4909 };
4910
4911 let transmit_handle = TableId::<TransmitHandle>::new(rep);
4912 let state = store.concurrent_state_mut()?;
4913 let transmit_id = state.get_mut(transmit_handle)?.state;
4914 let transmit = state.get_mut(transmit_id)?;
4915
4916 match event {
4917 Event::StreamRead { .. } => {
4918 transmit.read = ReadState::Open;
4919 }
4920 Event::StreamWrite { .. } => transmit.write = WriteState::Open,
4921 _ => {}
4922 }
4923 Ok(())
4924 }
4925}
4926
4927fn allow_intra_component_read_write(ty: Option<&InterfaceType>) -> bool {
4931 matches!(
4932 ty,
4933 None | Some(
4934 InterfaceType::S8
4935 | InterfaceType::U8
4936 | InterfaceType::S16
4937 | InterfaceType::U16
4938 | InterfaceType::S32
4939 | InterfaceType::U32
4940 | InterfaceType::S64
4941 | InterfaceType::U64
4942 | InterfaceType::Float32
4943 | InterfaceType::Float64
4944 )
4945 )
4946}
4947
4948struct LockedState<T> {
4952 inner: TryMutex<Option<T>>,
4953}
4954
4955impl<T> LockedState<T> {
4956 fn new(value: T) -> Self {
4958 Self {
4959 inner: TryMutex::new(Some(value)),
4960 }
4961 }
4962
4963 fn try_lock(&self) -> Result<TryMutexGuard<'_, Option<T>>> {
4972 match self.inner.try_lock() {
4973 Some(lock) => Ok(lock),
4974 None => bail_bug!("should not have contention on state lock"),
4975 }
4976 }
4977
4978 fn take(&self) -> Result<LockedStateGuard<'_, T>> {
4985 let result = self.try_lock()?.take();
4986 match result {
4987 Some(result) => Ok(LockedStateGuard {
4988 value: ManuallyDrop::new(result),
4989 state: self,
4990 }),
4991 None => bail_bug!("lock value unexpectedly missing"),
4992 }
4993 }
4994
4995 fn with<R>(&self, f: impl FnOnce(&mut T) -> R) -> Result<R> {
5004 let mut inner = self.try_lock()?;
5005 match &mut *inner {
5006 Some(state) => Ok(f(state)),
5007 None => bail_bug!("lock value unexpectedly missing"),
5008 }
5009 }
5010}
5011
5012struct LockedStateGuard<'a, T> {
5015 value: ManuallyDrop<T>,
5016 state: &'a LockedState<T>,
5017}
5018
5019impl<T> Deref for LockedStateGuard<'_, T> {
5020 type Target = T;
5021
5022 fn deref(&self) -> &T {
5023 &self.value
5024 }
5025}
5026
5027impl<T> DerefMut for LockedStateGuard<'_, T> {
5028 fn deref_mut(&mut self) -> &mut T {
5029 &mut self.value
5030 }
5031}
5032
5033impl<T> Drop for LockedStateGuard<'_, T> {
5034 fn drop(&mut self) {
5035 let value = unsafe { ManuallyDrop::take(&mut self.value) };
5040
5041 if let Ok(mut lock) = self.state.try_lock() {
5045 *lock = Some(value);
5046 }
5047 }
5048}
5049
5050#[cfg(test)]
5051mod tests {
5052 use super::*;
5053 use crate::{Engine, Store};
5054 use core::future::pending;
5055 use core::pin::pin;
5056 use std::sync::LazyLock;
5057
5058 static ENGINE: LazyLock<Engine> = LazyLock::new(Engine::default);
5059
5060 fn poll_future_producer<T>(rx: Pin<&mut T>, finish: bool) -> Poll<Result<Option<T::Item>>>
5061 where
5062 T: FutureProducer<()>,
5063 {
5064 rx.poll_produce(
5065 &mut Context::from_waker(Waker::noop()),
5066 Store::new(&ENGINE, ()).as_context_mut(),
5067 finish,
5068 )
5069 }
5070
5071 #[test]
5072 fn future_producer() {
5073 let mut fut = pin!(async { crate::error::Ok(()) });
5074 assert!(matches!(
5075 poll_future_producer(fut.as_mut(), false),
5076 Poll::Ready(Ok(Some(()))),
5077 ));
5078
5079 let mut fut = pin!(async { crate::error::Ok(()) });
5080 assert!(matches!(
5081 poll_future_producer(fut.as_mut(), true),
5082 Poll::Ready(Ok(Some(()))),
5083 ));
5084
5085 let mut fut = pin!(pending::<Result<()>>());
5086 assert!(matches!(
5087 poll_future_producer(fut.as_mut(), false),
5088 Poll::Pending,
5089 ));
5090 assert!(matches!(
5091 poll_future_producer(fut.as_mut(), true),
5092 Poll::Ready(Ok(None)),
5093 ));
5094
5095 let (tx, rx) = oneshot::channel();
5096 let mut rx = pin!(rx);
5097 assert!(matches!(
5098 poll_future_producer(rx.as_mut(), false),
5099 Poll::Pending,
5100 ));
5101 assert!(matches!(
5102 poll_future_producer(rx.as_mut(), true),
5103 Poll::Ready(Ok(None)),
5104 ));
5105 tx.send(()).unwrap();
5106 assert!(matches!(
5107 poll_future_producer(rx.as_mut(), true),
5108 Poll::Ready(Ok(Some(()))),
5109 ));
5110
5111 let (tx, rx) = oneshot::channel();
5112 let mut rx = pin!(rx);
5113 tx.send(()).unwrap();
5114 assert!(matches!(
5115 poll_future_producer(rx.as_mut(), false),
5116 Poll::Ready(Ok(Some(()))),
5117 ));
5118
5119 let (tx, rx) = oneshot::channel::<()>();
5120 let mut rx = pin!(rx);
5121 drop(tx);
5122 assert!(matches!(
5123 poll_future_producer(rx.as_mut(), false),
5124 Poll::Ready(Err(..)),
5125 ));
5126
5127 let (tx, rx) = oneshot::channel::<()>();
5128 let mut rx = pin!(rx);
5129 drop(tx);
5130 assert!(matches!(
5131 poll_future_producer(rx.as_mut(), true),
5132 Poll::Ready(Err(..)),
5133 ));
5134 }
5135}