Skip to content

Commit 8147838

Browse files
fix(moonbit): serialize stream writer teardown
1 parent 037b12a commit 8147838

7 files changed

Lines changed: 52 additions & 18 deletions

File tree

crates/moonbit/src/async_support.rs

Lines changed: 13 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1984,7 +1984,11 @@ fn wasm{symbol_name}StreamCommit(handle: Int) -> Unit {{
19841984
wasmImport{symbol_name}DropWritable(writer)
19851985
}}
19861986
}}
1987-
defer close_writer()
1987+
let close_writer_serialized = async fn() -> Unit {{
1988+
writer_lock.acquire()
1989+
defer writer_lock.release()
1990+
close_writer()
1991+
}}
19881992
let cleanup_value = fn(value : {result}) -> Unit {{
19891993
let ptr = wasm{symbol_name}Malloc(1)
19901994
wasm{symbol_name}Lower(value, ptr)
@@ -2054,11 +2058,7 @@ fn wasm{symbol_name}StreamCommit(handle: Int) -> Unit {{
20542058
// consumed the rest of this staging window.
20552059
data_len
20562060
}},
2057-
() => {{
2058-
writer_lock.acquire()
2059-
defer writer_lock.release()
2060-
close_writer()
2061-
}},
2061+
() => close_writer_serialized(),
20622062
() => !writer_closed.val,
20632063
Some(cleanup_value),
20642064
)
@@ -2096,6 +2096,13 @@ fn wasm{symbol_name}StreamCommit(handle: Int) -> Unit {{
20962096
}}
20972097
}}
20982098
}}
2099+
// A retained Sink may still own an in-flight canonical write after its
2100+
// producer returns. Do not drop the writable endpoint until that write
2101+
// has reached a terminal event and released its staging buffer.
2102+
{ffi}protect_from_cancel(
2103+
() => close_writer_serialized(),
2104+
resume_on_cancel=true,
2105+
)
20992106
}})
21002107
}}
21012108

crates/moonbit/src/lib.rs

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3585,8 +3585,13 @@ mod tests {
35853585
assert!(ffi.contains("suspend_for_future_write_terminal"));
35863586
assert!(ffi.contains("let data_len = if data.length() < 1"));
35873587
assert!(ffi.contains("let writer_lock = @async-core.Mutex()"));
3588+
assert!(ffi.contains("let close_writer_serialized = async fn()"));
35883589
assert!(ffi.contains("writer_lock.acquire()"));
35893590
assert!(ffi.contains("defer writer_lock.release()"));
3591+
assert!(
3592+
ffi.contains("() => close_writer_serialized(),\n resume_on_cancel=true,")
3593+
);
3594+
assert!(!ffi.contains("defer close_writer()"));
35903595
assert!(ffi.contains("read_cleanup : @async-core.CondVar"));
35913596
assert!(ffi.contains("self.read_cleanup.broadcast()"));
35923597
assert!(ffi.contains("@async-core.cancel_future_read("));

tests/runtime/moonbit/stream-write-cancel/holder-service.rs

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -60,6 +60,13 @@ impl Guest for Component {
6060
result == StreamResult::Dropped && values.is_empty()
6161
}
6262

63+
async fn read_one_and_keep() -> bool {
64+
let mut stream = HELD_STREAM.with(|stream| stream.borrow_mut().take().unwrap());
65+
let (result, values) = stream.read(Vec::with_capacity(1)).await;
66+
HELD_STREAM.with(|held| assert!(held.borrow_mut().replace(stream).is_none()));
67+
result == StreamResult::Complete(1) && values.len() == 1
68+
}
69+
6370
async fn read_one_and_drop() -> bool {
6471
let mut stream = HELD_STREAM.with(|stream| stream.borrow_mut().take().unwrap());
6572
let (result, values) = stream.read(Vec::with_capacity(1)).await;

tests/runtime/moonbit/stream-write-cancel/runner.rs

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -24,20 +24,21 @@ impl Guest for Component {
2424
.is_pending());
2525
holder::wait_write_started().await;
2626
assert!(holder::write_started());
27+
assert!(holder::read_one_and_keep().await);
2728

2829
drop(task);
2930
wait_cancellation_observed().await;
3031
assert!(cancellation_observed());
3132
assert!(holder::writable_dropped().await);
3233
assert_eq!(holder::leaf_live_count(), 0);
33-
assert_eq!(holder::leaf_drop_count(), 1);
34+
assert_eq!(holder::leaf_drop_count(), 2);
3435

3536
start_peer_drop().await;
3637
assert!(holder::read_one_and_drop().await);
3738
wait_producer_finished().await;
3839
assert!(second_write_started());
3940
assert!(producer_finished());
4041
assert_eq!(holder::leaf_live_count(), 0);
41-
assert_eq!(holder::leaf_drop_count(), 5);
42+
assert_eq!(holder::leaf_drop_count(), 6);
4243
}
4344
}

tests/runtime/moonbit/stream-write-cancel/test.mbt

Lines changed: 20 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -25,17 +25,28 @@ let producer_done : (@async-core.Future[Unit], @async-core.Promise[Unit]) =
2525

2626
///|
2727
pub async fn run_until_cancelled(
28-
_background_group : @async-core.TaskGroup[Unit]
28+
background_group : @async-core.TaskGroup[Unit]
2929
) -> Unit {
30+
let writer_started : (@async-core.Future[Unit], @async-core.Promise[Unit]) =
31+
@async-core.Future::new()
3032
let stream = @async-core.Stream::produce(async fn(sink) {
31-
defer {
32-
cancellation_seen.val = true
33-
ignore(cancellation_done.1.complete(()))
34-
}
35-
@holder.mark_write_started()
36-
let data : FixedArray[@holder.Leaf] = [@holder.Leaf::leaf()]
37-
let _ = sink.write_all(data[:])
38-
})
33+
// Return while a child owns a partially accepted write, forcing generated
34+
// finalization to wait for the child's canonical buffer cleanup.
35+
background_group.spawn_bg(async fn() {
36+
defer {
37+
cancellation_seen.val = true
38+
ignore(cancellation_done.1.complete(()))
39+
}
40+
@holder.mark_write_started()
41+
ignore(writer_started.1.complete(()))
42+
let data : FixedArray[@holder.Leaf] = [
43+
@holder.Leaf::leaf(),
44+
@holder.Leaf::leaf(),
45+
]
46+
let _ = sink.write_all(data[:])
47+
})
48+
writer_started.0.get()
49+
}, cleanup=fn(value) { value.drop() })
3950
@holder.hold(stream)
4051
never.0.get()
4152
}

tests/runtime/moonbit/stream-write-cancel/test.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -44,7 +44,9 @@ impl Guest for Component {
4444
let (mut writer, reader) = wit_stream::new::<holder::Leaf>();
4545
holder::hold(reader).await;
4646
holder::mark_write_started();
47-
let _ = writer.write_one(holder::Leaf::new()).await;
47+
let _ = writer
48+
.write(vec![holder::Leaf::new(), holder::Leaf::new()])
49+
.await;
4850
std::future::pending::<()>().await;
4951
}
5052

tests/runtime/moonbit/stream-write-cancel/test.wit

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ interface holder {
1313
write-started: func() -> bool;
1414
wait-write-started: async func();
1515
writable-dropped: async func() -> bool;
16+
read-one-and-keep: async func() -> bool;
1617
read-one-and-drop: async func() -> bool;
1718
leaf-live-count: func() -> u32;
1819
leaf-drop-count: func() -> u32;

0 commit comments

Comments
 (0)