1use super::table::{TableDebug, TableId};
2use super::{Event, GlobalErrorContextRefCount, Waitable, WaitableCommon};
3use crate::component::concurrent::{ConcurrentState, QualifiedThreadId, WaitReason, 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,
10 RuntimeInstance, StreamAny, 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_runtime_instance: RuntimeInstance,
3269 write_ty: TransmitIndex,
3270 write_options: OptionsIndex,
3271 write_address: usize,
3272 read_runtime_instance: RuntimeInstance,
3273 read_caller_thread: QualifiedThreadId,
3274 read_ty: TransmitIndex,
3275 read_options: OptionsIndex,
3276 read_address: usize,
3277 count: ItemCount,
3278 rep: u32,
3279 ) -> Result<()> {
3280 let write_instance = Instance::from_runtime_instance(store.0, write_runtime_instance);
3281 let read_instance = Instance::from_runtime_instance(store.0, read_runtime_instance);
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_runtime_instance == read_runtime_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 let address = usize::try_from(address)?;
3506 self.check_bounds(store.0, options, ty, address, count.as_usize())?;
3507 let (rep, state) = self.id().get_mut(store.0).get_mut_by_index(ty, handle)?;
3508 let TransmitLocalState::Write { done } = *state else {
3509 bail!(Trap::ConcurrentFutureStreamOp);
3510 };
3511
3512 if done {
3513 bail!("cannot write after being notified that the readable end dropped");
3514 }
3515
3516 *state = TransmitLocalState::Busy;
3517 let transmit_handle = TableId::<TransmitHandle>::new(rep);
3518 let concurrent_state = store.0.concurrent_state_mut()?;
3519 let transmit_id = concurrent_state.get_mut(transmit_handle)?.state;
3520 let transmit = concurrent_state.get_mut(transmit_id)?;
3521 log::trace!(
3522 "guest_write {count} to {transmit_handle:?} (handle {handle}; state {transmit_id:?}); {:?}",
3523 transmit.read
3524 );
3525
3526 if transmit.done {
3527 bail!("cannot write to future after previous write succeeded or readable end dropped");
3528 }
3529
3530 let new_state = if let ReadState::Dropped = &transmit.read {
3531 ReadState::Dropped
3532 } else {
3533 ReadState::Open
3534 };
3535
3536 let set_guest_ready = |me: &mut ConcurrentState| {
3537 let transmit = me.get_mut(transmit_id)?;
3538 if !matches!(&transmit.write, WriteState::Open) {
3539 bail_bug!("expected `WriteState::Open`; got `{:?}`", transmit.write);
3540 }
3541 transmit.write = WriteState::GuestReady {
3542 instance: self,
3543 caller,
3544 ty,
3545 flat_abi,
3546 options,
3547 address,
3548 count,
3549 handle,
3550 };
3551 Ok::<_, crate::Error>(())
3552 };
3553
3554 let mut result = match mem::replace(&mut transmit.read, new_state) {
3555 ReadState::GuestReady {
3556 ty: read_ty,
3557 flat_abi: read_flat_abi,
3558 options: read_options,
3559 address: read_address,
3560 count: read_count,
3561 handle: read_handle,
3562 instance: read_instance,
3563 caller_instance: read_caller_instance,
3564 caller_thread: read_caller_thread,
3565 } => {
3566 if flat_abi != read_flat_abi {
3567 bail_bug!("expected flat ABI calculations to be the same");
3568 }
3569
3570 if let TransmitIndex::Future(_) = ty {
3571 transmit.done = true;
3572 }
3573
3574 let write_count = count;
3596 let write_complete = count == 0 || read_count > 0;
3597 let read_complete = count > 0;
3598 let read_buffer_remaining = count < read_count;
3599
3600 let read_handle_rep = transmit.read_handle.rep();
3601
3602 let count = count.min(read_count);
3603
3604 Instance::copy(
3605 store.as_context_mut(),
3606 flat_abi,
3607 self.runtime_instance(caller),
3608 ty,
3609 options,
3610 address,
3611 read_instance.runtime_instance(read_caller_instance),
3612 read_caller_thread,
3613 read_ty,
3614 read_options,
3615 read_address,
3616 count,
3617 rep,
3618 )?;
3619
3620 let instance = read_instance.id().get(store.0);
3621 let types = instance.component().types();
3622 let item_size = match read_ty.payload(types) {
3623 Some(ty) => usize::try_from(types.canonical_abi(ty).size32)?,
3624 None => 0,
3625 };
3626 let concurrent_state = store.0.concurrent_state_mut()?;
3627 if read_complete {
3628 let total = if let Some(Event::StreamRead {
3629 code: ReturnCode::Completed(old_total),
3630 ..
3631 }) = concurrent_state.take_event(read_handle_rep)?
3632 {
3633 count.add(old_total)?
3634 } else {
3635 count
3636 };
3637
3638 let code = ReturnCode::completed(ty.kind(), total);
3639
3640 concurrent_state.send_read_result(read_ty, transmit_id, read_handle, code)?;
3641 }
3642
3643 if read_buffer_remaining || (write_count == 0 && read_count == 0) {
3650 let transmit = concurrent_state.get_mut(transmit_id)?;
3651 transmit.read = ReadState::GuestReady {
3652 ty: read_ty,
3653 flat_abi: read_flat_abi,
3654 options: read_options,
3655 address: read_address + (count.as_usize() * item_size),
3656 count: read_count.sub(count)?,
3657 handle: read_handle,
3658 instance: read_instance,
3659 caller_instance: read_caller_instance,
3660 caller_thread: read_caller_thread,
3661 };
3662 }
3663
3664 if write_complete {
3665 ReturnCode::completed(ty.kind(), count)
3666 } else {
3667 set_guest_ready(concurrent_state)?;
3668 ReturnCode::Blocked
3669 }
3670 }
3671
3672 ReadState::HostReady {
3673 consume,
3674 guest_offset,
3675 cancel,
3676 cancel_waker,
3677 } => {
3678 if cancel_waker.is_some() {
3679 bail_bug!("expected cancel_waker to be none");
3680 }
3681 if cancel {
3682 bail_bug!("expected cancel to be false");
3683 }
3684 if guest_offset != 0 {
3685 bail_bug!("expected guest_offset to be 0");
3686 }
3687
3688 if let TransmitIndex::Future(_) = ty {
3689 transmit.done = true;
3690 }
3691
3692 set_guest_ready(concurrent_state)?;
3693 self.consume(
3694 store.0,
3695 ty.kind(),
3696 transmit_id,
3697 consume,
3698 ItemCount::ZERO,
3699 false,
3700 )?
3701 }
3702
3703 ReadState::HostToHost { .. } => bail_bug!("unexpected HostToHost"),
3704
3705 ReadState::Open => {
3706 set_guest_ready(concurrent_state)?;
3707 ReturnCode::Blocked
3708 }
3709
3710 ReadState::Dropped => {
3711 if let TransmitIndex::Future(_) = ty {
3712 transmit.done = true;
3713 }
3714
3715 match Waitable::Transmit(transmit_handle).take_event(concurrent_state)? {
3716 Some(
3717 Event::StreamWrite {
3718 code: ReturnCode::Dropped(ItemCount::ZERO),
3719 ..
3720 }
3721 | Event::FutureWrite {
3722 code: ReturnCode::Dropped(ItemCount::ZERO),
3723 ..
3724 },
3725 ) => {}
3726 None => bail!(match ty {
3727 TransmitIndex::Future(_) => Trap::WriteToDroppedFuture,
3728 TransmitIndex::Stream(_) => Trap::WriteToDroppedStream,
3729 }),
3730 event => bail_bug!("expected pending dropped event for writer; got {event:?}"),
3731 }
3732
3733 ReturnCode::Dropped(ItemCount::ZERO)
3734 }
3735 };
3736
3737 if result == ReturnCode::Blocked && !self.options(store.0, options).async_ {
3738 result = self.wait_for_write(store.0, caller, transmit_handle)?;
3739 }
3740
3741 if result != ReturnCode::Blocked {
3742 *self.id().get_mut(store.0).get_mut_by_index(ty, handle)?.1 =
3743 TransmitLocalState::Write {
3744 done: matches!(result, ReturnCode::Dropped(_)),
3745 };
3746 }
3747
3748 log::trace!(
3749 "guest_write result for {transmit_handle:?} (handle {handle}; state {transmit_id:?}): {result:?}",
3750 );
3751
3752 Ok(result)
3753 }
3754
3755 pub(super) fn guest_read<T: 'static>(
3757 self,
3758 mut store: StoreContextMut<T>,
3759 caller_instance: RuntimeComponentInstanceIndex,
3760 ty: TransmitIndex,
3761 options: OptionsIndex,
3762 flat_abi: Option<FlatAbi>,
3763 handle: u32,
3764 address: u32,
3765 count: u32,
3766 ) -> Result<ReturnCode> {
3767 let count = ItemCount::new(count)?;
3768
3769 let address = usize::try_from(address)?;
3770 self.check_bounds(store.0, options, ty, address, count.as_usize())?;
3771 let (rep, state) = self.id().get_mut(store.0).get_mut_by_index(ty, handle)?;
3772 let TransmitLocalState::Read { done } = *state else {
3773 bail!(Trap::ConcurrentFutureStreamOp);
3774 };
3775
3776 if done {
3777 bail!("cannot read after being notified that the writable end dropped");
3778 }
3779
3780 *state = TransmitLocalState::Busy;
3781 let transmit_handle = TableId::<TransmitHandle>::new(rep);
3782 let caller_thread = store.0.current_guest_thread()?;
3783 let concurrent_state = store.0.concurrent_state_mut()?;
3784 let transmit_id = concurrent_state.get_mut(transmit_handle)?.state;
3785 let transmit = concurrent_state.get_mut(transmit_id)?;
3786 log::trace!(
3787 "guest_read {count} from {transmit_handle:?} (handle {handle}; state {transmit_id:?}); {:?}",
3788 transmit.write
3789 );
3790
3791 if transmit.done {
3792 bail!("cannot read from future after previous read succeeded");
3793 }
3794
3795 let new_state = if let WriteState::Dropped = &transmit.write {
3796 WriteState::Dropped
3797 } else {
3798 WriteState::Open
3799 };
3800
3801 let set_guest_ready = |me: &mut ConcurrentState| {
3802 let transmit = me.get_mut(transmit_id)?;
3803 if !matches!(&transmit.read, ReadState::Open) {
3804 bail_bug!("expected `ReadState::Open`; got `{:?}`", transmit.read);
3805 }
3806 transmit.read = ReadState::GuestReady {
3807 ty,
3808 flat_abi,
3809 options,
3810 address,
3811 count,
3812 handle,
3813 instance: self,
3814 caller_instance,
3815 caller_thread,
3816 };
3817 Ok::<_, crate::Error>(())
3818 };
3819
3820 let mut result = match mem::replace(&mut transmit.write, new_state) {
3821 WriteState::GuestReady {
3822 instance: write_instance,
3823 ty: write_ty,
3824 flat_abi: write_flat_abi,
3825 options: write_options,
3826 address: write_address,
3827 count: write_count,
3828 handle: write_handle,
3829 caller: write_caller,
3830 } => {
3831 if flat_abi != write_flat_abi {
3832 bail_bug!("expected flat ABI calculations to be the same");
3833 }
3834
3835 if let TransmitIndex::Future(_) = ty {
3836 transmit.done = true;
3837 }
3838
3839 let write_handle_rep = transmit.write_handle.rep();
3840
3841 let write_complete = write_count == 0 || count > 0;
3846 let read_complete = write_count > 0;
3847 let write_buffer_remaining = count < write_count;
3848
3849 let count = count.min(write_count);
3850
3851 Instance::copy(
3852 store.as_context_mut(),
3853 flat_abi,
3854 write_instance.runtime_instance(write_caller),
3855 write_ty,
3856 write_options,
3857 write_address,
3858 self.runtime_instance(caller_instance),
3859 caller_thread,
3860 ty,
3861 options,
3862 address,
3863 count,
3864 rep,
3865 )?;
3866
3867 let instance = write_instance.id().get(store.0);
3868 let types = instance.component().types();
3869 let item_size = match write_ty.payload(types) {
3870 Some(ty) => usize::try_from(types.canonical_abi(ty).size32)?,
3871 None => 0,
3872 };
3873 let concurrent_state = store.0.concurrent_state_mut()?;
3874
3875 if write_complete {
3876 let total = if let Some(Event::StreamWrite {
3877 code: ReturnCode::Completed(old_total),
3878 ..
3879 }) = concurrent_state.take_event(write_handle_rep)?
3880 {
3881 count.add(old_total)?
3882 } else {
3883 count
3884 };
3885
3886 let code = ReturnCode::completed(ty.kind(), total);
3887
3888 concurrent_state.send_write_result(
3889 write_ty,
3890 transmit_id,
3891 write_handle,
3892 code,
3893 )?;
3894 }
3895
3896 if write_buffer_remaining {
3897 let transmit = concurrent_state.get_mut(transmit_id)?;
3898 transmit.write = WriteState::GuestReady {
3899 instance: write_instance,
3900 caller: write_caller,
3901 ty: write_ty,
3902 flat_abi: write_flat_abi,
3903 options: write_options,
3904 address: write_address + (count.as_usize() * item_size),
3905 count: write_count.sub(count)?,
3906 handle: write_handle,
3907 };
3908 }
3909
3910 if read_complete {
3911 ReturnCode::completed(ty.kind(), count)
3912 } else {
3913 set_guest_ready(concurrent_state)?;
3914 ReturnCode::Blocked
3915 }
3916 }
3917
3918 WriteState::HostReady {
3919 produce,
3920 try_into,
3921 guest_offset,
3922 cancel,
3923 cancel_waker,
3924 } => {
3925 if cancel_waker.is_some() {
3926 bail_bug!("expected cancel_waker to be none");
3927 }
3928 if cancel {
3929 bail_bug!("expected cancel to be false");
3930 }
3931 if guest_offset != 0 {
3932 bail_bug!("expected guest_offset to be 0");
3933 }
3934
3935 set_guest_ready(concurrent_state)?;
3936
3937 let code = self.produce(
3938 store.0,
3939 ty.kind(),
3940 transmit_id,
3941 produce,
3942 try_into,
3943 ItemCount::ZERO,
3944 false,
3945 )?;
3946
3947 if let (TransmitIndex::Future(_), ReturnCode::Completed(_)) = (ty, code) {
3948 store.0.concurrent_state_mut()?.get_mut(transmit_id)?.done = true;
3949 }
3950
3951 code
3952 }
3953
3954 WriteState::Open => {
3955 set_guest_ready(concurrent_state)?;
3956 ReturnCode::Blocked
3957 }
3958
3959 WriteState::Dropped => {
3960 if let TransmitIndex::Future(_) = ty {
3961 bail_bug!(
3962 "should not be possible to read from a future whose write end was dropped"
3963 );
3964 }
3965
3966 match Waitable::Transmit(transmit_handle).take_event(concurrent_state)? {
3967 Some(Event::StreamRead {
3968 code: ReturnCode::Dropped(ItemCount::ZERO),
3969 ..
3970 }) => {}
3971 None => bail!(Trap::ReadFromDroppedStream),
3972 event => bail_bug!("expected pending dropped event for reader; got {event:?}"),
3973 }
3974
3975 ReturnCode::Dropped(ItemCount::ZERO)
3976 }
3977 };
3978
3979 if result == ReturnCode::Blocked && !self.options(store.0, options).async_ {
3980 result = self.wait_for_read(store.0, caller_instance, transmit_handle)?;
3981 }
3982
3983 if result != ReturnCode::Blocked {
3984 *self.id().get_mut(store.0).get_mut_by_index(ty, handle)?.1 =
3985 TransmitLocalState::Read {
3986 done: matches!(
3987 (result, ty),
3988 (ReturnCode::Dropped(_), TransmitIndex::Stream(_))
3989 ),
3990 };
3991 }
3992
3993 log::trace!(
3994 "guest_read result for {transmit_handle:?} (handle {handle}; state {transmit_id:?}): {result:?}",
3995 );
3996
3997 Ok(result)
3998 }
3999
4000 fn wait_for_write(
4001 self,
4002 store: &mut StoreOpaque,
4003 caller: RuntimeComponentInstanceIndex,
4004 handle: TableId<TransmitHandle>,
4005 ) -> Result<ReturnCode> {
4006 let waitable = Waitable::Transmit(handle);
4007 store.wait_for_event(self.runtime_instance(caller), waitable, WaitReason::Other)?;
4008 let event = waitable.take_event(store.concurrent_state_mut()?)?;
4009 if let Some(event @ (Event::StreamWrite { code, .. } | Event::FutureWrite { code, .. })) =
4010 event
4011 {
4012 waitable.on_delivery(store, self, event)?;
4013 Ok(code)
4014 } else {
4015 bail_bug!("expected either a stream or future write event")
4016 }
4017 }
4018
4019 fn cancel_write(
4021 self,
4022 store: &mut StoreOpaque,
4023 caller: RuntimeComponentInstanceIndex,
4024 transmit_id: TableId<TransmitState>,
4025 async_: bool,
4026 ) -> Result<ReturnCode> {
4027 let state = store.concurrent_state_mut()?;
4028 let transmit = state.get_mut(transmit_id)?;
4029 log::trace!(
4030 "host_cancel_write state {transmit_id:?}; write state {:?} read state {:?}",
4031 transmit.read,
4032 transmit.write
4033 );
4034 let waitable = Waitable::Transmit(transmit.write_handle);
4035
4036 if !async_ {
4037 waitable.trap_if_in_waitable_set(state)?;
4038 }
4039
4040 let code = if let Some(event) = waitable.take_event(state)? {
4041 let (Event::FutureWrite { code, .. } | Event::StreamWrite { code, .. }) = event else {
4042 bail_bug!("expected either a stream or future write event")
4043 };
4044 waitable.on_delivery(store, self, event)?;
4045 match (code, event) {
4046 (ReturnCode::Completed(count), Event::StreamWrite { .. }) => {
4047 ReturnCode::Cancelled(count)
4048 }
4049 (ReturnCode::Dropped(_) | ReturnCode::Completed(_), _) => code,
4050 _ => bail_bug!("unexpected code/event combo"),
4051 }
4052 } else if let ReadState::HostReady {
4053 cancel,
4054 cancel_waker,
4055 ..
4056 } = &mut state.get_mut(transmit_id)?.read
4057 {
4058 *cancel = true;
4059 if let Some(waker) = cancel_waker.take() {
4060 waker.wake();
4061 }
4062
4063 if async_ {
4064 ReturnCode::Blocked
4065 } else {
4066 let handle = store
4067 .concurrent_state_mut()?
4068 .get_mut(transmit_id)?
4069 .write_handle;
4070 self.wait_for_write(store, caller, handle)?
4071 }
4072 } else {
4073 ReturnCode::Cancelled(ItemCount::ZERO)
4074 };
4075
4076 if !matches!(code, ReturnCode::Blocked) {
4077 let transmit = store.concurrent_state_mut()?.get_mut(transmit_id)?;
4078
4079 match &transmit.write {
4080 WriteState::GuestReady { .. } => {
4081 transmit.write = WriteState::Open;
4082 }
4083 WriteState::HostReady { .. } => bail_bug!("support host write cancellation"),
4084 WriteState::Open | WriteState::Dropped => {}
4085 }
4086 }
4087
4088 log::trace!("cancelled write {transmit_id:?}: {code:?}");
4089
4090 Ok(code)
4091 }
4092
4093 fn wait_for_read(
4094 self,
4095 store: &mut StoreOpaque,
4096 caller: RuntimeComponentInstanceIndex,
4097 handle: TableId<TransmitHandle>,
4098 ) -> Result<ReturnCode> {
4099 let waitable = Waitable::Transmit(handle);
4100 store.wait_for_event(self.runtime_instance(caller), waitable, WaitReason::Other)?;
4101 let event = waitable.take_event(store.concurrent_state_mut()?)?;
4102 if let Some(event @ (Event::StreamRead { code, .. } | Event::FutureRead { code, .. })) =
4103 event
4104 {
4105 waitable.on_delivery(store, self, event)?;
4106 Ok(code)
4107 } else {
4108 bail_bug!("expected either a stream or future read event")
4109 }
4110 }
4111
4112 fn cancel_read(
4114 self,
4115 store: &mut StoreOpaque,
4116 caller: RuntimeComponentInstanceIndex,
4117 transmit_id: TableId<TransmitState>,
4118 async_: bool,
4119 ) -> Result<ReturnCode> {
4120 let state = store.concurrent_state_mut()?;
4121 let transmit = state.get_mut(transmit_id)?;
4122 log::trace!(
4123 "host_cancel_read state {transmit_id:?}; read state {:?} write state {:?}",
4124 transmit.read,
4125 transmit.write
4126 );
4127
4128 let waitable = Waitable::Transmit(transmit.read_handle);
4129
4130 if !async_ {
4131 waitable.trap_if_in_waitable_set(state)?;
4132 }
4133
4134 let code = if let Some(event) = waitable.take_event(state)? {
4135 let (Event::FutureRead { code, .. } | Event::StreamRead { code, .. }) = event else {
4136 bail_bug!("expected either a stream or future read event")
4137 };
4138 waitable.on_delivery(store, self, event)?;
4139 match (code, event) {
4140 (ReturnCode::Completed(count), Event::StreamRead { .. }) => {
4141 ReturnCode::Cancelled(count)
4142 }
4143 (ReturnCode::Dropped(_) | ReturnCode::Completed(_), _) => code,
4144 _ => bail_bug!("unexpected code/event combo"),
4145 }
4146 } else if let WriteState::HostReady {
4147 cancel,
4148 cancel_waker,
4149 ..
4150 } = &mut state.get_mut(transmit_id)?.write
4151 {
4152 *cancel = true;
4153 if let Some(waker) = cancel_waker.take() {
4154 waker.wake();
4155 }
4156
4157 if async_ {
4158 ReturnCode::Blocked
4159 } else {
4160 let handle = store
4161 .concurrent_state_mut()?
4162 .get_mut(transmit_id)?
4163 .read_handle;
4164 self.wait_for_read(store, caller, handle)?
4165 }
4166 } else {
4167 ReturnCode::Cancelled(ItemCount::ZERO)
4168 };
4169
4170 if !matches!(code, ReturnCode::Blocked) {
4171 let transmit = store.concurrent_state_mut()?.get_mut(transmit_id)?;
4172
4173 match &transmit.read {
4174 ReadState::GuestReady { .. } => {
4175 transmit.read = ReadState::Open;
4176 }
4177 ReadState::HostReady { .. } | ReadState::HostToHost { .. } => {
4178 bail_bug!("support host read cancellation")
4179 }
4180 ReadState::Open | ReadState::Dropped => {}
4181 }
4182 }
4183
4184 log::trace!("cancelled read {transmit_id:?}: {code:?}");
4185
4186 Ok(code)
4187 }
4188
4189 fn guest_cancel_write(
4191 self,
4192 store: &mut StoreOpaque,
4193 caller: RuntimeComponentInstanceIndex,
4194 ty: TransmitIndex,
4195 async_: bool,
4196 writer: u32,
4197 ) -> Result<ReturnCode> {
4198 let (rep, state) =
4199 get_mut_by_index_from(self.id().get_mut(store).table_for_transmit(ty), ty, writer)?;
4200 let id = TableId::<TransmitHandle>::new(rep);
4201 log::trace!("guest cancel write {id:?} (handle {writer})");
4202 match state {
4203 TransmitLocalState::Write { .. } => {
4204 bail!("stream or future write cancelled when no write is pending")
4205 }
4206 TransmitLocalState::Read { .. } => {
4207 bail!("passed read end to `{{stream|future}}.cancel-write`")
4208 }
4209 TransmitLocalState::Busy => {}
4210 }
4211 let transmit_id = store.concurrent_state_mut()?.get_mut(id)?.state;
4212 let code = self.cancel_write(store, caller, transmit_id, async_)?;
4213 if !matches!(code, ReturnCode::Blocked) {
4214 let state =
4215 get_mut_by_index_from(self.id().get_mut(store).table_for_transmit(ty), ty, writer)?
4216 .1;
4217 if let TransmitLocalState::Busy = state {
4218 *state = TransmitLocalState::Write { done: false };
4219 }
4220 }
4221 Ok(code)
4222 }
4223
4224 fn guest_cancel_read(
4226 self,
4227 store: &mut StoreOpaque,
4228 caller: RuntimeComponentInstanceIndex,
4229 ty: TransmitIndex,
4230 async_: bool,
4231 reader: u32,
4232 ) -> Result<ReturnCode> {
4233 let (rep, state) =
4234 get_mut_by_index_from(self.id().get_mut(store).table_for_transmit(ty), ty, reader)?;
4235 let id = TableId::<TransmitHandle>::new(rep);
4236 log::trace!("guest cancel read {id:?} (handle {reader})");
4237 match state {
4238 TransmitLocalState::Read { .. } => {
4239 bail!("stream or future read cancelled when no read is pending")
4240 }
4241 TransmitLocalState::Write { .. } => {
4242 bail!("passed write end to `{{stream|future}}.cancel-read`")
4243 }
4244 TransmitLocalState::Busy => {}
4245 }
4246 let transmit_id = store.concurrent_state_mut()?.get_mut(id)?.state;
4247 let code = self.cancel_read(store, caller, transmit_id, async_)?;
4248 if !matches!(code, ReturnCode::Blocked) {
4249 let state =
4250 get_mut_by_index_from(self.id().get_mut(store).table_for_transmit(ty), ty, reader)?
4251 .1;
4252 if let TransmitLocalState::Busy = state {
4253 *state = TransmitLocalState::Read { done: false };
4254 }
4255 }
4256 Ok(code)
4257 }
4258
4259 fn guest_drop_readable(
4261 self,
4262 store: &mut StoreOpaque,
4263 ty: TransmitIndex,
4264 reader: u32,
4265 ) -> Result<()> {
4266 let table = self.id().get_mut(store).table_for_transmit(ty);
4267 let (rep, _is_done) = match ty {
4268 TransmitIndex::Stream(ty) => table.stream_remove_readable(ty, reader)?,
4269 TransmitIndex::Future(ty) => table.future_remove_readable(ty, reader)?,
4270 };
4271 let kind = match ty {
4272 TransmitIndex::Stream(_) => TransmitKind::Stream,
4273 TransmitIndex::Future(_) => TransmitKind::Future,
4274 };
4275 let id = TableId::<TransmitHandle>::new(rep);
4276 log::trace!("guest_drop_readable: drop reader {id:?}");
4277 store.host_drop_reader(id, kind)
4278 }
4279
4280 pub(crate) fn error_context_new(
4282 self,
4283 store: &mut StoreOpaque,
4284 ty: TypeComponentLocalErrorContextTableIndex,
4285 options: OptionsIndex,
4286 debug_msg_address: u32,
4287 debug_msg_len: u32,
4288 ) -> Result<u32> {
4289 let lift_ctx = &mut LiftContext::new(store, options, self)?;
4290 let debug_msg = String::linear_lift_from_flat(
4291 lift_ctx,
4292 InterfaceType::String,
4293 &[ValRaw::u32(debug_msg_address), ValRaw::u32(debug_msg_len)],
4294 )?;
4295
4296 let err_ctx = ErrorContextState { debug_msg };
4298 let state = store.concurrent_state_mut()?;
4299 let table_id = state.push(err_ctx)?;
4300 let global_ref_count_idx =
4301 TypeComponentGlobalErrorContextTableIndex::from_u32(table_id.rep());
4302
4303 let _ = state
4305 .global_error_context_ref_counts
4306 .insert(global_ref_count_idx, GlobalErrorContextRefCount(1));
4307
4308 let local_idx = self
4315 .id()
4316 .get_mut(store)
4317 .table_for_error_context(ty)
4318 .error_context_insert(table_id.rep())?;
4319
4320 Ok(local_idx)
4321 }
4322
4323 pub(super) fn error_context_debug_message<T>(
4325 self,
4326 store: StoreContextMut<T>,
4327 ty: TypeComponentLocalErrorContextTableIndex,
4328 options: OptionsIndex,
4329 err_ctx_handle: u32,
4330 debug_msg_address: u32,
4331 ) -> Result<()> {
4332 let handle_table_id_rep = self
4334 .id()
4335 .get_mut(store.0)
4336 .table_for_error_context(ty)
4337 .error_context_rep(err_ctx_handle)?;
4338
4339 let state = store.0.concurrent_state_mut()?;
4340 let ErrorContextState { debug_msg } =
4342 state.get_mut(TableId::<ErrorContextState>::new(handle_table_id_rep))?;
4343 let debug_msg = debug_msg.clone();
4344
4345 let lower_cx = &mut LowerContext::new(store, options, self);
4346 let debug_msg_address = usize::try_from(debug_msg_address)?;
4347 let offset = lower_cx
4353 .as_slice_mut()
4354 .get(debug_msg_address..)
4355 .and_then(|b| b.get(..8))
4356 .map(|_| debug_msg_address)
4357 .ok_or_else(|| crate::format_err!("invalid debug message pointer: out of bounds"))?;
4358 debug_msg
4359 .as_str()
4360 .linear_lower_to_memory(lower_cx, InterfaceType::String, offset)?;
4361
4362 Ok(())
4363 }
4364
4365 pub(crate) fn future_cancel_read(
4367 self,
4368 store: &mut StoreOpaque,
4369 caller: RuntimeComponentInstanceIndex,
4370 ty: TypeFutureTableIndex,
4371 async_: bool,
4372 reader: u32,
4373 ) -> Result<u32> {
4374 self.guest_cancel_read(store, caller, TransmitIndex::Future(ty), async_, reader)
4375 .map(|v| v.encode())
4376 }
4377
4378 pub(crate) fn future_cancel_write(
4380 self,
4381 store: &mut StoreOpaque,
4382 caller: RuntimeComponentInstanceIndex,
4383 ty: TypeFutureTableIndex,
4384 async_: bool,
4385 writer: u32,
4386 ) -> Result<u32> {
4387 self.guest_cancel_write(store, caller, TransmitIndex::Future(ty), async_, writer)
4388 .map(|v| v.encode())
4389 }
4390
4391 pub(crate) fn stream_cancel_read(
4393 self,
4394 store: &mut StoreOpaque,
4395 caller: RuntimeComponentInstanceIndex,
4396 ty: TypeStreamTableIndex,
4397 async_: bool,
4398 reader: u32,
4399 ) -> Result<u32> {
4400 self.guest_cancel_read(store, caller, TransmitIndex::Stream(ty), async_, reader)
4401 .map(|v| v.encode())
4402 }
4403
4404 pub(crate) fn stream_cancel_write(
4406 self,
4407 store: &mut StoreOpaque,
4408 caller: RuntimeComponentInstanceIndex,
4409 ty: TypeStreamTableIndex,
4410 async_: bool,
4411 writer: u32,
4412 ) -> Result<u32> {
4413 self.guest_cancel_write(store, caller, TransmitIndex::Stream(ty), async_, writer)
4414 .map(|v| v.encode())
4415 }
4416
4417 pub(crate) fn future_drop_readable(
4419 self,
4420 store: &mut StoreOpaque,
4421 ty: TypeFutureTableIndex,
4422 reader: u32,
4423 ) -> Result<()> {
4424 self.guest_drop_readable(store, TransmitIndex::Future(ty), reader)
4425 }
4426
4427 pub(crate) fn stream_drop_readable(
4429 self,
4430 store: &mut StoreOpaque,
4431 ty: TypeStreamTableIndex,
4432 reader: u32,
4433 ) -> Result<()> {
4434 self.guest_drop_readable(store, TransmitIndex::Stream(ty), reader)
4435 }
4436
4437 fn guest_new(self, store: &mut StoreOpaque, ty: TransmitIndex) -> Result<ResourcePair> {
4441 let (write, read) = store
4442 .concurrent_state_mut()?
4443 .new_transmit(TransmitOrigin::guest(self.id().instance(), ty))?;
4444
4445 let table = self.id().get_mut(store).table_for_transmit(ty);
4446 let (read_handle, write_handle) = match ty {
4447 TransmitIndex::Future(ty) => (
4448 table.future_insert_read(ty, read.rep())?,
4449 table.future_insert_write(ty, write.rep())?,
4450 ),
4451 TransmitIndex::Stream(ty) => (
4452 table.stream_insert_read(ty, read.rep())?,
4453 table.stream_insert_write(ty, write.rep())?,
4454 ),
4455 };
4456
4457 let state = store.concurrent_state_mut()?;
4458 state.get_mut(read)?.common.handle = Some(read_handle);
4459 state.get_mut(write)?.common.handle = Some(write_handle);
4460
4461 Ok(ResourcePair {
4462 write: write_handle,
4463 read: read_handle,
4464 })
4465 }
4466
4467 pub(crate) fn error_context_drop(
4469 self,
4470 store: &mut StoreOpaque,
4471 ty: TypeComponentLocalErrorContextTableIndex,
4472 error_context: u32,
4473 ) -> Result<()> {
4474 let instance = self.id().get_mut(store);
4475
4476 let local_handle_table = instance.table_for_error_context(ty);
4477
4478 let rep = local_handle_table.error_context_drop(error_context)?;
4479
4480 let global_ref_count_idx = TypeComponentGlobalErrorContextTableIndex::from_u32(rep);
4481
4482 let state = store.concurrent_state_mut()?;
4483 let Some(GlobalErrorContextRefCount(global_ref_count)) = state
4484 .global_error_context_ref_counts
4485 .get_mut(&global_ref_count_idx)
4486 else {
4487 bail_bug!("retrieve concurrent state for error context during drop")
4488 };
4489
4490 if *global_ref_count < 1 {
4492 bail_bug!("ref count unexpectedly zero");
4493 }
4494 *global_ref_count -= 1;
4495 if *global_ref_count == 0 {
4496 state
4497 .global_error_context_ref_counts
4498 .remove(&global_ref_count_idx);
4499
4500 state
4501 .delete(TableId::<ErrorContextState>::new(rep))
4502 .context("deleting component-global error context data")?;
4503 }
4504
4505 Ok(())
4506 }
4507
4508 fn guest_transfer(
4511 self,
4512 store: &mut StoreOpaque,
4513 src_idx: u32,
4514 src: TransmitIndex,
4515 dst: TransmitIndex,
4516 ) -> Result<u32> {
4517 let id = self.lift_index_to_transmit(store, src, src_idx)?;
4518 self.lower_transmit_to_index(store, dst, id)
4519 }
4520
4521 fn lift_index_to_transmit(
4522 self,
4523 store: &mut StoreOpaque,
4524 ty: TransmitIndex,
4525 src_idx: u32,
4526 ) -> Result<TableId<TransmitHandle>> {
4527 let (state, _, _, instance) = store.lift_context_parts(self);
4528 lift_index_to_transmit(instance, state.concurrent_state_mut(), ty, src_idx)
4529 }
4530
4531 fn lower_transmit_to_index(
4532 self,
4533 store: &mut StoreOpaque,
4534 ty: TransmitIndex,
4535 id: TableId<TransmitHandle>,
4536 ) -> Result<u32> {
4537 let (state, _, _, instance) = store.lift_context_parts(self);
4538 lower_transmit_to_index(instance, state.concurrent_state_mut(), ty, id)
4539 }
4540
4541 pub(crate) fn future_new(
4543 self,
4544 store: &mut StoreOpaque,
4545 ty: TypeFutureTableIndex,
4546 ) -> Result<ResourcePair> {
4547 self.guest_new(store, TransmitIndex::Future(ty))
4548 }
4549
4550 pub(crate) fn stream_new(
4552 self,
4553 store: &mut StoreOpaque,
4554 ty: TypeStreamTableIndex,
4555 ) -> Result<ResourcePair> {
4556 self.guest_new(store, TransmitIndex::Stream(ty))
4557 }
4558
4559 pub(crate) fn future_transfer(
4562 self,
4563 store: &mut StoreOpaque,
4564 src_idx: u32,
4565 src: TypeFutureTableIndex,
4566 dst: TypeFutureTableIndex,
4567 ) -> Result<u32> {
4568 self.guest_transfer(
4569 store,
4570 src_idx,
4571 TransmitIndex::Future(src),
4572 TransmitIndex::Future(dst),
4573 )
4574 }
4575
4576 pub(crate) fn stream_transfer(
4579 self,
4580 store: &mut StoreOpaque,
4581 src_idx: u32,
4582 src: TypeStreamTableIndex,
4583 dst: TypeStreamTableIndex,
4584 ) -> Result<u32> {
4585 self.guest_transfer(
4586 store,
4587 src_idx,
4588 TransmitIndex::Stream(src),
4589 TransmitIndex::Stream(dst),
4590 )
4591 }
4592
4593 pub(crate) fn error_context_transfer(
4595 self,
4596 store: &mut StoreOpaque,
4597 src_idx: u32,
4598 src: TypeComponentLocalErrorContextTableIndex,
4599 dst: TypeComponentLocalErrorContextTableIndex,
4600 ) -> Result<u32> {
4601 let mut instance = self.id().get_mut(store);
4602 let rep = instance
4603 .as_mut()
4604 .table_for_error_context(src)
4605 .error_context_rep(src_idx)?;
4606 let dst_idx = instance
4607 .table_for_error_context(dst)
4608 .error_context_insert(rep)?;
4609
4610 let global_ref_count = store
4614 .concurrent_state_mut()?
4615 .global_error_context_ref_counts
4616 .get_mut(&TypeComponentGlobalErrorContextTableIndex::from_u32(rep))
4617 .context("global ref count present for existing (sub)component error context")?;
4618
4619 global_ref_count.0 = global_ref_count
4620 .0
4621 .checked_add(1)
4622 .ok_or_else(|| format_err!(Trap::ReferenceCountOverflow))?;
4623
4624 Ok(dst_idx)
4625 }
4626}
4627
4628fn lift_index_to_transmit(
4634 instance: Pin<&mut ComponentInstance>,
4635 concurrent_state: &mut ConcurrentState,
4636 ty: TransmitIndex,
4637 src_idx: u32,
4638) -> Result<TableId<TransmitHandle>> {
4639 let handle_table = instance.table_for_transmit(ty);
4640 let (rep, is_done) = match ty {
4641 TransmitIndex::Future(idx) => handle_table.future_remove_readable(idx, src_idx)?,
4642 TransmitIndex::Stream(idx) => handle_table.stream_remove_readable(idx, src_idx)?,
4643 };
4644 let desc = match ty {
4645 TransmitIndex::Future(_) => "future",
4646 TransmitIndex::Stream(_) => "stream",
4647 };
4648 if is_done {
4649 bail!("cannot lift {desc} after being notified that the writable end dropped");
4650 }
4651 let id = TableId::<TransmitHandle>::new(rep);
4652 let future = concurrent_state.get_mut(id)?;
4653 if future.common.set.is_some() {
4654 bail!("cannot lift {desc} while it's in a waitable set");
4655 }
4656 future.common.handle = None;
4657
4658 let state = future.state;
4659 if concurrent_state.get_mut(state)?.done {
4660 match ty {
4661 TransmitIndex::Future(_) => bail!("cannot lift {desc} after previous read succeeded"),
4662 TransmitIndex::Stream(_) => bail!(Trap::LiftDroppedStream),
4663 };
4664 }
4665
4666 Ok(id)
4667}
4668
4669fn lower_transmit_to_index(
4672 instance: Pin<&mut ComponentInstance>,
4673 concurrent_state: &mut ConcurrentState,
4674 ty: TransmitIndex,
4675 id: TableId<TransmitHandle>,
4676) -> Result<u32> {
4677 let state = concurrent_state.get_mut(id)?.state;
4678 debug_assert_eq!(concurrent_state.get_mut(state)?.read_handle, id);
4679 let handle_table = instance.table_for_transmit(ty);
4680 let handle = match ty {
4681 TransmitIndex::Future(idx) => handle_table.future_insert_read(idx, id.rep()),
4682 TransmitIndex::Stream(idx) => handle_table.stream_insert_read(idx, id.rep()),
4683 }?;
4684 concurrent_state.get_mut(id)?.common.handle = Some(handle);
4685 Ok(handle)
4686}
4687
4688impl ComponentInstance {
4689 fn table_for_transmit(self: Pin<&mut Self>, ty: TransmitIndex) -> &mut HandleTable {
4690 let (states, types) = self.instance_states();
4691 let runtime_instance = match ty {
4692 TransmitIndex::Stream(ty) => types[ty].instance,
4693 TransmitIndex::Future(ty) => types[ty].instance,
4694 };
4695 states[runtime_instance].handle_table()
4696 }
4697
4698 fn table_for_error_context(
4699 self: Pin<&mut Self>,
4700 ty: TypeComponentLocalErrorContextTableIndex,
4701 ) -> &mut HandleTable {
4702 let (states, types) = self.instance_states();
4703 let runtime_instance = types[ty].instance;
4704 states[runtime_instance].handle_table()
4705 }
4706
4707 fn get_mut_by_index(
4708 self: Pin<&mut Self>,
4709 ty: TransmitIndex,
4710 index: u32,
4711 ) -> Result<(u32, &mut TransmitLocalState)> {
4712 get_mut_by_index_from(self.table_for_transmit(ty), ty, index)
4713 }
4714}
4715
4716impl ConcurrentState {
4717 fn send_write_result(
4718 &mut self,
4719 ty: TransmitIndex,
4720 id: TableId<TransmitState>,
4721 handle: u32,
4722 code: ReturnCode,
4723 ) -> Result<()> {
4724 let write_handle = self.get_mut(id)?.write_handle.rep();
4725 self.set_event(
4726 write_handle,
4727 match ty {
4728 TransmitIndex::Future(ty) => Event::FutureWrite {
4729 code,
4730 pending: Some((ty, handle)),
4731 },
4732 TransmitIndex::Stream(ty) => Event::StreamWrite {
4733 code,
4734 pending: Some((ty, handle)),
4735 },
4736 },
4737 )
4738 }
4739
4740 fn send_read_result(
4741 &mut self,
4742 ty: TransmitIndex,
4743 id: TableId<TransmitState>,
4744 handle: u32,
4745 code: ReturnCode,
4746 ) -> Result<()> {
4747 let read_handle = self.get_mut(id)?.read_handle.rep();
4748 self.set_event(
4749 read_handle,
4750 match ty {
4751 TransmitIndex::Future(ty) => Event::FutureRead {
4752 code,
4753 pending: Some((ty, handle)),
4754 },
4755 TransmitIndex::Stream(ty) => Event::StreamRead {
4756 code,
4757 pending: Some((ty, handle)),
4758 },
4759 },
4760 )
4761 }
4762
4763 fn take_event(&mut self, waitable: u32) -> Result<Option<Event>> {
4764 Waitable::Transmit(TableId::<TransmitHandle>::new(waitable)).take_event(self)
4765 }
4766
4767 fn set_event(&mut self, waitable: u32, event: Event) -> Result<()> {
4768 Waitable::Transmit(TableId::<TransmitHandle>::new(waitable)).set_event(self, Some(event))
4769 }
4770
4771 fn update_event(&mut self, waitable: u32, event: Event) -> Result<()> {
4782 let waitable = Waitable::Transmit(TableId::<TransmitHandle>::new(waitable));
4783
4784 fn update_code(old: ReturnCode, new: ReturnCode) -> Result<ReturnCode> {
4785 let (ReturnCode::Completed(count)
4786 | ReturnCode::Dropped(count)
4787 | ReturnCode::Cancelled(count)) = old
4788 else {
4789 bail_bug!("unexpected old return code")
4790 };
4791
4792 Ok(match new {
4793 ReturnCode::Dropped(ItemCount::ZERO) => ReturnCode::Dropped(count),
4794 ReturnCode::Cancelled(ItemCount::ZERO) => ReturnCode::Cancelled(count),
4795 _ => bail_bug!("unexpected new return code"),
4796 })
4797 }
4798
4799 let event = match (waitable.take_event(self)?, event) {
4800 (None, _) => event,
4801 (Some(old @ Event::FutureWrite { .. }), Event::FutureWrite { .. }) => old,
4802 (Some(old @ Event::FutureRead { .. }), Event::FutureRead { .. }) => old,
4803 (
4804 Some(Event::StreamWrite {
4805 code: old_code,
4806 pending: old_pending,
4807 }),
4808 Event::StreamWrite { code, pending },
4809 ) => Event::StreamWrite {
4810 code: update_code(old_code, code)?,
4811 pending: old_pending.or(pending),
4812 },
4813 (
4814 Some(Event::StreamRead {
4815 code: old_code,
4816 pending: old_pending,
4817 }),
4818 Event::StreamRead { code, pending },
4819 ) => Event::StreamRead {
4820 code: update_code(old_code, code)?,
4821 pending: old_pending.or(pending),
4822 },
4823 _ => bail_bug!("unexpected event combination"),
4824 };
4825
4826 waitable.set_event(self, Some(event))
4827 }
4828
4829 fn new_transmit(
4832 &mut self,
4833 origin: TransmitOrigin,
4834 ) -> Result<(TableId<TransmitHandle>, TableId<TransmitHandle>)> {
4835 let state_id = self.push(TransmitState::new(origin))?;
4836
4837 let write = self.push(TransmitHandle::new(state_id))?;
4838 let read = self.push(TransmitHandle::new(state_id))?;
4839
4840 let state = self.get_mut(state_id)?;
4841 state.write_handle = write;
4842 state.read_handle = read;
4843
4844 log::trace!("new transmit: state {state_id:?}; write {write:?}; read {read:?}",);
4845
4846 Ok((write, read))
4847 }
4848
4849 fn delete_transmit(&mut self, state_id: TableId<TransmitState>) -> Result<()> {
4851 let state = self.delete(state_id)?;
4852 self.delete(state.write_handle)?;
4853 self.delete(state.read_handle)?;
4854
4855 log::trace!(
4856 "delete transmit: state {state_id:?}; write {:?}; read {:?}",
4857 state.write_handle,
4858 state.read_handle,
4859 );
4860
4861 Ok(())
4862 }
4863}
4864
4865pub(crate) struct ResourcePair {
4866 pub(crate) write: u32,
4867 pub(crate) read: u32,
4868}
4869
4870impl Waitable {
4871 pub(super) fn on_delivery(
4874 &self,
4875 store: &mut StoreOpaque,
4876 instance: Instance,
4877 event: Event,
4878 ) -> Result<()> {
4879 if let Event::FutureRead {
4880 code: ReturnCode::Dropped(_),
4881 ..
4882 }
4883 | Event::FutureWrite {
4884 code: ReturnCode::Dropped(_),
4885 ..
4886 }
4887 | Event::StreamRead {
4888 code: ReturnCode::Dropped(_),
4889 ..
4890 }
4891 | Event::StreamWrite {
4892 code: ReturnCode::Dropped(_),
4893 ..
4894 } = event
4895 {
4896 let Waitable::Transmit(transmit_handle) = self else {
4897 bail_bug!("unexpected `{event:?}` for `{self:?}`");
4898 };
4899 let state = store.concurrent_state_mut()?;
4900 let transmit_id = state.get_mut(*transmit_handle)?.state;
4901 state.get_mut(transmit_id)?.done = true;
4902 }
4903
4904 let instance = instance.id().get_mut(store);
4905 let (rep, state, code) = match event {
4906 Event::FutureRead {
4907 pending: Some((ty, handle)),
4908 code,
4909 }
4910 | Event::FutureWrite {
4911 pending: Some((ty, handle)),
4912 code,
4913 } => {
4914 let runtime_instance = instance.component().types()[ty].instance;
4915 let (rep, state) = instance.instance_states().0[runtime_instance]
4916 .handle_table()
4917 .future_rep(ty, handle)?;
4918 (rep, state, code)
4919 }
4920 Event::StreamRead {
4921 pending: Some((ty, handle)),
4922 code,
4923 }
4924 | Event::StreamWrite {
4925 pending: Some((ty, handle)),
4926 code,
4927 } => {
4928 let runtime_instance = instance.component().types()[ty].instance;
4929 let (rep, state) = instance.instance_states().0[runtime_instance]
4930 .handle_table()
4931 .stream_rep(ty, handle)?;
4932 (rep, state, code)
4933 }
4934 _ => return Ok(()),
4935 };
4936 if rep != self.rep() {
4937 bail_bug!("unexpected rep mismatch");
4938 }
4939 if *state != TransmitLocalState::Busy {
4940 bail_bug!("expected state to be busy");
4941 }
4942 let done = matches!(code, ReturnCode::Dropped(_));
4943 *state = match event {
4944 Event::FutureRead { .. } | Event::StreamRead { .. } => {
4945 TransmitLocalState::Read { done }
4946 }
4947 Event::FutureWrite { .. } | Event::StreamWrite { .. } => {
4948 TransmitLocalState::Write { done }
4949 }
4950 _ => bail_bug!("unexpected event for stream"),
4951 };
4952
4953 let transmit_handle = TableId::<TransmitHandle>::new(rep);
4954 let state = store.concurrent_state_mut()?;
4955 let transmit_id = state.get_mut(transmit_handle)?.state;
4956 let transmit = state.get_mut(transmit_id)?;
4957
4958 match event {
4959 Event::StreamRead { .. } => {
4960 transmit.read = ReadState::Open;
4961 }
4962 Event::StreamWrite { .. } => transmit.write = WriteState::Open,
4963 _ => {}
4964 }
4965 Ok(())
4966 }
4967}
4968
4969fn allow_intra_component_read_write(ty: Option<&InterfaceType>) -> bool {
4973 matches!(
4974 ty,
4975 None | Some(
4976 InterfaceType::S8
4977 | InterfaceType::U8
4978 | InterfaceType::S16
4979 | InterfaceType::U16
4980 | InterfaceType::S32
4981 | InterfaceType::U32
4982 | InterfaceType::S64
4983 | InterfaceType::U64
4984 | InterfaceType::Float32
4985 | InterfaceType::Float64
4986 )
4987 )
4988}
4989
4990struct LockedState<T> {
4994 inner: TryMutex<Option<T>>,
4995}
4996
4997impl<T> LockedState<T> {
4998 fn new(value: T) -> Self {
5000 Self {
5001 inner: TryMutex::new(Some(value)),
5002 }
5003 }
5004
5005 fn try_lock(&self) -> Result<TryMutexGuard<'_, Option<T>>> {
5014 match self.inner.try_lock() {
5015 Some(lock) => Ok(lock),
5016 None => bail_bug!("should not have contention on state lock"),
5017 }
5018 }
5019
5020 fn take(&self) -> Result<LockedStateGuard<'_, T>> {
5027 let result = self.try_lock()?.take();
5028 match result {
5029 Some(result) => Ok(LockedStateGuard {
5030 value: ManuallyDrop::new(result),
5031 state: self,
5032 }),
5033 None => bail_bug!("lock value unexpectedly missing"),
5034 }
5035 }
5036
5037 fn with<R>(&self, f: impl FnOnce(&mut T) -> R) -> Result<R> {
5046 let mut inner = self.try_lock()?;
5047 match &mut *inner {
5048 Some(state) => Ok(f(state)),
5049 None => bail_bug!("lock value unexpectedly missing"),
5050 }
5051 }
5052}
5053
5054struct LockedStateGuard<'a, T> {
5057 value: ManuallyDrop<T>,
5058 state: &'a LockedState<T>,
5059}
5060
5061impl<T> Deref for LockedStateGuard<'_, T> {
5062 type Target = T;
5063
5064 fn deref(&self) -> &T {
5065 &self.value
5066 }
5067}
5068
5069impl<T> DerefMut for LockedStateGuard<'_, T> {
5070 fn deref_mut(&mut self) -> &mut T {
5071 &mut self.value
5072 }
5073}
5074
5075impl<T> Drop for LockedStateGuard<'_, T> {
5076 fn drop(&mut self) {
5077 let value = unsafe { ManuallyDrop::take(&mut self.value) };
5082
5083 if let Ok(mut lock) = self.state.try_lock() {
5087 *lock = Some(value);
5088 }
5089 }
5090}
5091
5092#[cfg(test)]
5093mod tests {
5094 use super::*;
5095 use crate::{Engine, Store};
5096 use core::future::pending;
5097 use core::pin::pin;
5098 use std::sync::LazyLock;
5099
5100 static ENGINE: LazyLock<Engine> = LazyLock::new(Engine::default);
5101
5102 fn poll_future_producer<T>(rx: Pin<&mut T>, finish: bool) -> Poll<Result<Option<T::Item>>>
5103 where
5104 T: FutureProducer<()>,
5105 {
5106 rx.poll_produce(
5107 &mut Context::from_waker(Waker::noop()),
5108 Store::new(&ENGINE, ()).as_context_mut(),
5109 finish,
5110 )
5111 }
5112
5113 #[test]
5114 fn future_producer() {
5115 let mut fut = pin!(async { crate::error::Ok(()) });
5116 assert!(matches!(
5117 poll_future_producer(fut.as_mut(), false),
5118 Poll::Ready(Ok(Some(()))),
5119 ));
5120
5121 let mut fut = pin!(async { crate::error::Ok(()) });
5122 assert!(matches!(
5123 poll_future_producer(fut.as_mut(), true),
5124 Poll::Ready(Ok(Some(()))),
5125 ));
5126
5127 let mut fut = pin!(pending::<Result<()>>());
5128 assert!(matches!(
5129 poll_future_producer(fut.as_mut(), false),
5130 Poll::Pending,
5131 ));
5132 assert!(matches!(
5133 poll_future_producer(fut.as_mut(), true),
5134 Poll::Ready(Ok(None)),
5135 ));
5136
5137 let (tx, rx) = oneshot::channel();
5138 let mut rx = pin!(rx);
5139 assert!(matches!(
5140 poll_future_producer(rx.as_mut(), false),
5141 Poll::Pending,
5142 ));
5143 assert!(matches!(
5144 poll_future_producer(rx.as_mut(), true),
5145 Poll::Ready(Ok(None)),
5146 ));
5147 tx.send(()).unwrap();
5148 assert!(matches!(
5149 poll_future_producer(rx.as_mut(), true),
5150 Poll::Ready(Ok(Some(()))),
5151 ));
5152
5153 let (tx, rx) = oneshot::channel();
5154 let mut rx = pin!(rx);
5155 tx.send(()).unwrap();
5156 assert!(matches!(
5157 poll_future_producer(rx.as_mut(), false),
5158 Poll::Ready(Ok(Some(()))),
5159 ));
5160
5161 let (tx, rx) = oneshot::channel::<()>();
5162 let mut rx = pin!(rx);
5163 drop(tx);
5164 assert!(matches!(
5165 poll_future_producer(rx.as_mut(), false),
5166 Poll::Ready(Err(..)),
5167 ));
5168
5169 let (tx, rx) = oneshot::channel::<()>();
5170 let mut rx = pin!(rx);
5171 drop(tx);
5172 assert!(matches!(
5173 poll_future_producer(rx.as_mut(), true),
5174 Poll::Ready(Err(..)),
5175 ));
5176 }
5177}