Skip to content

Commit e079d70

Browse files
test(moonbit): add Rust async fixture pairings
1 parent b183f26 commit e079d70

2 files changed

Lines changed: 263 additions & 0 deletions

File tree

  • tests/runtime/moonbit
Lines changed: 150 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,150 @@
1+
include!(env!("BINDINGS"));
2+
3+
use crate::exports::test::moonbit_nested_future_stream::nested::Guest;
4+
use std::sync::Mutex;
5+
use std::sync::atomic::{AtomicBool, Ordering};
6+
use std::task::{Poll, Waker};
7+
use wit_bindgen::{FutureReader, StreamReader, StreamResult};
8+
9+
struct Component;
10+
11+
export!(Component);
12+
13+
static SHARED: Mutex<Option<u32>> = Mutex::new(None);
14+
static SHARED_WAKER: Mutex<Option<Waker>> = Mutex::new(None);
15+
static CANCELLATION_OBSERVED: AtomicBool = AtomicBool::new(false);
16+
static CANCELLATION_WAKER: Mutex<Option<Waker>> = Mutex::new(None);
17+
18+
struct CancellationGuard;
19+
20+
impl Drop for CancellationGuard {
21+
fn drop(&mut self) {
22+
CANCELLATION_OBSERVED.store(true, Ordering::SeqCst);
23+
if let Some(waker) = CANCELLATION_WAKER.lock().unwrap().take() {
24+
waker.wake();
25+
}
26+
}
27+
}
28+
29+
impl Guest for Component {
30+
async fn relay(
31+
value: FutureReader<FutureReader<StreamReader<u8>>>,
32+
) -> FutureReader<FutureReader<StreamReader<u8>>> {
33+
let (outer_writer, outer_reader) = wit_future::new(|| unreachable!());
34+
wit_bindgen::spawn_local(async move {
35+
let input_inner = value.await;
36+
let (inner_writer, inner_reader) = wit_future::new(|| unreachable!());
37+
let outer_open = outer_writer.write(inner_reader).await.is_ok();
38+
39+
// Continue consuming the input chain after an outer rejection so
40+
// that its producer observes the rejection at the innermost stream.
41+
let mut input_stream = input_inner.await;
42+
let (mut output_writer, output_reader) = wit_stream::new();
43+
let inner_open = inner_writer.write(output_reader).await.is_ok();
44+
if !outer_open || !inner_open {
45+
return;
46+
}
47+
48+
loop {
49+
let (result, values) = input_stream.read(Vec::with_capacity(16)).await;
50+
if !values.is_empty() && !output_writer.write_all(values).await.is_empty() {
51+
return;
52+
}
53+
match result {
54+
StreamResult::Complete(_) => {}
55+
StreamResult::Dropped => return,
56+
StreamResult::Cancelled => unreachable!(),
57+
}
58+
}
59+
});
60+
outer_reader
61+
}
62+
63+
async fn relay_stream(value: StreamReader<FutureReader<u8>>) -> StreamReader<FutureReader<u8>> {
64+
let (mut output_writer, output_reader) = wit_stream::new();
65+
wit_bindgen::spawn_local(async move {
66+
let mut input = value;
67+
loop {
68+
let (result, values) = input.read(Vec::with_capacity(1)).await;
69+
for input_value in values {
70+
let (value_writer, value_reader) = wit_future::new(|| unreachable!());
71+
if output_writer.write_one(value_reader).await.is_some() {
72+
let _ = input_value.await;
73+
return;
74+
}
75+
let value = input_value.await;
76+
let _ = value_writer.write(value).await;
77+
}
78+
match result {
79+
StreamResult::Complete(_) => {}
80+
StreamResult::Dropped => return,
81+
StreamResult::Cancelled => unreachable!(),
82+
}
83+
}
84+
});
85+
output_reader
86+
}
87+
88+
async fn concurrent_writes() -> StreamReader<u8> {
89+
let (mut writer, reader) = wit_stream::new();
90+
wit_bindgen::spawn_local(async move {
91+
assert!(writer.write_all(vec![1, 2]).await.is_empty());
92+
});
93+
reader
94+
}
95+
96+
async fn post_return_lazy() -> StreamReader<u8> {
97+
let (mut writer, reader) = wit_stream::new();
98+
wit_bindgen::spawn_local(async move {
99+
assert!(writer.write_one(42).await.is_none());
100+
});
101+
reader
102+
}
103+
104+
async fn wait_shared() -> u32 {
105+
std::future::poll_fn(|cx| {
106+
if let Some(value) = SHARED.lock().unwrap().take() {
107+
return Poll::Ready(value);
108+
}
109+
110+
*SHARED_WAKER.lock().unwrap() = Some(cx.waker().clone());
111+
match SHARED.lock().unwrap().take() {
112+
Some(value) => Poll::Ready(value),
113+
None => Poll::Pending,
114+
}
115+
})
116+
.await
117+
}
118+
119+
async fn resolve_shared(value: u32) {
120+
assert!(SHARED.lock().unwrap().replace(value).is_none());
121+
if let Some(waker) = SHARED_WAKER.lock().unwrap().take() {
122+
waker.wake();
123+
}
124+
}
125+
126+
async fn wait_cancelled() {
127+
let _guard = CancellationGuard;
128+
std::future::pending::<()>().await;
129+
}
130+
131+
async fn wait_cancellation_observed() {
132+
std::future::poll_fn(|cx| {
133+
if CANCELLATION_OBSERVED.load(Ordering::SeqCst) {
134+
return Poll::Ready(());
135+
}
136+
137+
*CANCELLATION_WAKER.lock().unwrap() = Some(cx.waker().clone());
138+
if CANCELLATION_OBSERVED.load(Ordering::SeqCst) {
139+
Poll::Ready(())
140+
} else {
141+
Poll::Pending
142+
}
143+
})
144+
.await
145+
}
146+
147+
fn cancellation_observed() -> bool {
148+
CANCELLATION_OBSERVED.load(Ordering::SeqCst)
149+
}
150+
}
Lines changed: 113 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,113 @@
1+
include!(env!("BINDINGS"));
2+
3+
use crate::exports::test::moonbit_stream_write_cancel::controller::Guest;
4+
use crate::test::moonbit_stream_write_cancel::holder;
5+
use std::sync::Mutex;
6+
use std::sync::atomic::{AtomicBool, Ordering};
7+
use std::task::{Poll, Waker};
8+
9+
struct Component;
10+
11+
export!(Component);
12+
13+
static CANCELLATION_OBSERVED: AtomicBool = AtomicBool::new(false);
14+
static CANCELLATION_WAKER: Mutex<Option<Waker>> = Mutex::new(None);
15+
static SECOND_WRITE_STARTED: AtomicBool = AtomicBool::new(false);
16+
static PRODUCER_FINISHED: AtomicBool = AtomicBool::new(false);
17+
static PRODUCER_WAKER: Mutex<Option<Waker>> = Mutex::new(None);
18+
19+
struct CancellationGuard;
20+
21+
impl Drop for CancellationGuard {
22+
fn drop(&mut self) {
23+
CANCELLATION_OBSERVED.store(true, Ordering::SeqCst);
24+
if let Some(waker) = CANCELLATION_WAKER.lock().unwrap().take() {
25+
waker.wake();
26+
}
27+
}
28+
}
29+
30+
struct ProducerGuard;
31+
32+
impl Drop for ProducerGuard {
33+
fn drop(&mut self) {
34+
PRODUCER_FINISHED.store(true, Ordering::SeqCst);
35+
if let Some(waker) = PRODUCER_WAKER.lock().unwrap().take() {
36+
waker.wake();
37+
}
38+
}
39+
}
40+
41+
impl Guest for Component {
42+
async fn run_until_cancelled() {
43+
let _guard = CancellationGuard;
44+
let (mut writer, reader) = wit_stream::new::<holder::Leaf>();
45+
holder::hold(reader).await;
46+
holder::mark_write_started();
47+
let _ = writer.write_one(holder::Leaf::new()).await;
48+
std::future::pending::<()>().await;
49+
}
50+
51+
fn cancellation_observed() -> bool {
52+
CANCELLATION_OBSERVED.load(Ordering::SeqCst)
53+
}
54+
55+
async fn wait_cancellation_observed() {
56+
std::future::poll_fn(|cx| {
57+
if CANCELLATION_OBSERVED.load(Ordering::SeqCst) {
58+
return Poll::Ready(());
59+
}
60+
61+
*CANCELLATION_WAKER.lock().unwrap() = Some(cx.waker().clone());
62+
if CANCELLATION_OBSERVED.load(Ordering::SeqCst) {
63+
Poll::Ready(())
64+
} else {
65+
Poll::Pending
66+
}
67+
})
68+
.await
69+
}
70+
71+
async fn start_peer_drop() {
72+
// The MoonBit implementation primes its language-level queue locally.
73+
// A CM stream cannot be read and written by the same component when
74+
// its payload contains resources, so account for that local item here.
75+
drop(holder::Leaf::new());
76+
77+
let (mut writer, reader) = wit_stream::new::<holder::Leaf>();
78+
wit_bindgen::spawn_local(async move {
79+
let _guard = ProducerGuard;
80+
81+
assert!(writer.write_one(holder::Leaf::new()).await.is_none());
82+
SECOND_WRITE_STARTED.store(true, Ordering::SeqCst);
83+
assert!(writer.write_one(holder::Leaf::new()).await.is_some());
84+
assert!(writer.write_one(holder::Leaf::new()).await.is_some());
85+
});
86+
87+
holder::hold(reader).await;
88+
}
89+
90+
fn second_write_started() -> bool {
91+
SECOND_WRITE_STARTED.load(Ordering::SeqCst)
92+
}
93+
94+
fn producer_finished() -> bool {
95+
PRODUCER_FINISHED.load(Ordering::SeqCst)
96+
}
97+
98+
async fn wait_producer_finished() {
99+
std::future::poll_fn(|cx| {
100+
if PRODUCER_FINISHED.load(Ordering::SeqCst) {
101+
return Poll::Ready(());
102+
}
103+
104+
*PRODUCER_WAKER.lock().unwrap() = Some(cx.waker().clone());
105+
if PRODUCER_FINISHED.load(Ordering::SeqCst) {
106+
Poll::Ready(())
107+
} else {
108+
Poll::Pending
109+
}
110+
})
111+
.await
112+
}
113+
}

0 commit comments

Comments
 (0)