Skip to content

Commit 7c48901

Browse files
perf(moonbit): tune primitive stream buffering (#1664)
1 parent ce2815c commit 7c48901

3 files changed

Lines changed: 131 additions & 16 deletions

File tree

crates/moonbit/src/async/trait.mbt

Lines changed: 27 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -5,10 +5,11 @@ pub struct Sink[X] {
55
priv is_open : () -> Bool
66
priv has_cleanup : () -> Bool
77
priv cleanup : ((X) -> Unit)?
8+
priv write_window_size : Int
89
}
910

1011
///|
11-
let sink_write_window_size : Int = 64
12+
let default_sink_write_window_size : Int = 64
1213

1314
///|
1415
/// The write callback must either raise before consuming any value or return
@@ -21,8 +22,17 @@ pub fn[X] Sink::from_callbacks(
2122
close : async () -> Unit,
2223
is_open : () -> Bool,
2324
cleanup : ((X) -> Unit)?,
25+
write_window_size : Int,
2426
) -> Sink[X] {
25-
{ write, close, is_open, has_cleanup: () => cleanup is Some(_), cleanup }
27+
guard write_window_size > 0
28+
{
29+
write,
30+
close,
31+
is_open,
32+
has_cleanup: () => cleanup is Some(_),
33+
cleanup,
34+
write_window_size,
35+
}
2636
}
2737

2838
///|
@@ -34,13 +44,13 @@ pub async fn[X] Sink::write(self : Sink[X], data : ArrayView[X]) -> Int {
3444
if data.length() == 0 {
3545
return 0
3646
}
37-
let length = if data.length() < sink_write_window_size {
47+
let length = if data.length() < self.write_window_size {
3848
data.length()
3949
} else {
40-
sink_write_window_size
50+
self.write_window_size
4151
}
4252
// The callback may suspend, so it must not retain the caller's borrowed view.
43-
let buffer = FixedArray::makei(length, i => data[i])
53+
let buffer = FixedArray::from_array(data[:length])
4454
let written = (self.write)(buffer[:])
4555
guard written >= 0 && written <= length
4656
if !(self.is_open)() && written < length && (self.has_cleanup)() {
@@ -60,12 +70,12 @@ pub async fn Sink::write_bytes(self : Sink[Byte], data : BytesView) -> Int {
6070
if data.length() == 0 {
6171
return 0
6272
}
63-
let length = if data.length() < sink_write_window_size {
73+
let length = if data.length() < self.write_window_size {
6474
data.length()
6575
} else {
66-
sink_write_window_size
76+
self.write_window_size
6777
}
68-
let buffer = FixedArray::makei(length, i => data[i])
78+
let buffer = data[:length].to_fixedarray()
6979
let written = (self.write)(buffer[:])
7080
guard written >= 0 && written <= length
7181
if !(self.is_open)() && written < length && (self.has_cleanup)() {
@@ -314,7 +324,7 @@ fn[X] fill_stream_pipe_from_writers(pipe : Ref[StreamPipe[X]]) -> Unit {
314324
guard take_stream_pipe_writer(pipe) is Some(writer) else { return }
315325
guard writer.value is Some(data) else { continue }
316326
let take = if available < data.length() { available } else { data.length() }
317-
let chunk = FixedArray::makei(take, i => data[i])
327+
let chunk = FixedArray::from_array(data[:take])
318328
pipe.val.chunks.push_back(chunk)
319329
pipe.val.buffered = pipe.val.buffered + take
320330
wake_stream_pipe_writer(writer, take)
@@ -415,7 +425,7 @@ async fn[X] stream_pipe_write(
415425
} else {
416426
data.length()
417427
}
418-
let chunk = FixedArray::makei(take, i => data[i])
428+
let chunk = FixedArray::from_array(data[:take])
419429
wake_stream_pipe_reader(reader, Some(chunk), false)
420430
return take
421431
}
@@ -424,13 +434,13 @@ async fn[X] stream_pipe_write(
424434
if pipe.val.capacity > 0 && pipe.val.buffered < pipe.val.capacity {
425435
let available = pipe.val.capacity - pipe.val.buffered
426436
let take = if available < data.length() { available } else { data.length() }
427-
let chunk = FixedArray::makei(take, i => data[i])
437+
let chunk = FixedArray::from_array(data[:take])
428438
pipe.val.chunks.push_back(chunk)
429439
pipe.val.buffered = pipe.val.buffered + take
430440
return take
431441
}
432442
let writer = StreamPipeWriter::{
433-
value: Some(FixedArray::makei(data.length(), i => data[i])),
443+
value: Some(FixedArray::from_array(data)),
434444
accepted: 0,
435445
coro: Some(current_coroutine()),
436446
}
@@ -463,7 +473,9 @@ async fn[X] stream_pipe_read(
463473
Some(head) => {
464474
let available = head.length() - pipe.val.head_pos
465475
let take = if count < available { count } else { available }
466-
let result = FixedArray::makei(take, i => head[pipe.val.head_pos + i])
476+
let result = FixedArray::from_array(
477+
head[pipe.val.head_pos:pipe.val.head_pos + take],
478+
)
467479
pipe.val.head_pos = pipe.val.head_pos + take
468480
pipe.val.buffered = pipe.val.buffered - take
469481
if pipe.val.head_pos >= head.length() {
@@ -484,7 +496,7 @@ async fn[X] stream_pipe_read(
484496
Some(writer) => {
485497
guard writer.value is Some(data) else { continue }
486498
let take = if count < data.length() { count } else { data.length() }
487-
let result = FixedArray::makei(take, i => data[i])
499+
let result = FixedArray::from_array(data[:take])
488500
wake_stream_pipe_writer(writer, take)
489501
return Some(result)
490502
}
@@ -740,6 +752,7 @@ fn[X] stream_pipe_sink(pipe : Ref[StreamPipe[X]]) -> Sink[X] {
740752
None => ()
741753
}
742754
}),
755+
write_window_size: default_sink_write_window_size,
743756
}
744757
}
745758

crates/moonbit/src/async_support.rs

Lines changed: 34 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,30 @@ use super::FunctionBindgen;
2121
use super::InterfaceGenerator;
2222
use super::wasm_type;
2323

24+
const DEFAULT_STREAM_WINDOW_ELEMENTS: usize = 64;
25+
const PRIMITIVE_STREAM_BUFFER_BYTES: usize = 4 * 1024;
26+
27+
fn is_fixed_primitive(resolve: &Resolve, ty: &Type) -> bool {
28+
match ty {
29+
Type::U8
30+
| Type::S8
31+
| Type::U16
32+
| Type::S16
33+
| Type::U32
34+
| Type::S32
35+
| Type::U64
36+
| Type::S64
37+
| Type::F32
38+
| Type::F64
39+
| Type::Char => true,
40+
Type::Id(id) => match &resolve.types[*id].kind {
41+
TypeDefKind::Type(ty) => is_fixed_primitive(resolve, ty),
42+
_ => false,
43+
},
44+
_ => false,
45+
}
46+
}
47+
2448
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
2549
enum PayloadFor {
2650
Future,
@@ -1233,8 +1257,15 @@ impl<'a> InterfaceGenerator<'a> {
12331257
.map(|ty| self.world_gen.sizes.size(ty).size_wasm32())
12341258
.unwrap_or(0);
12351259
let read_chunk_owns_buffer = result_type.is_some_and(|ty| self.is_list_canonical(ty));
1236-
let staging_window = if payload_sites.is_empty() { 64 } else { 1 };
1237-
let max_read_count = 64;
1260+
let primitive_window = result_type
1261+
.filter(|ty| is_fixed_primitive(self.resolve, ty))
1262+
.map(|_| (PRIMITIVE_STREAM_BUFFER_BYTES / elem_size).max(1));
1263+
let staging_window = if payload_sites.is_empty() {
1264+
primitive_window.unwrap_or(DEFAULT_STREAM_WINDOW_ELEMENTS)
1265+
} else {
1266+
1
1267+
};
1268+
let max_read_count = primitive_window.unwrap_or(DEFAULT_STREAM_WINDOW_ELEMENTS);
12381269

12391270
let EndpointPayloadFragments {
12401271
lift,
@@ -2104,6 +2135,7 @@ fn wasm{symbol_name}StreamCommit(handle : Int) -> Unit {{
21042135
() => close_writer_serialized(),
21052136
() => !writer_closed.val,
21062137
Some(cleanup_value),
2138+
{staging_window},
21072139
)
21082140
let relay_source = producer is None
21092141
let run_producer = async fn() -> Unit {{

crates/moonbit/src/lib.rs

Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3445,6 +3445,76 @@ mod tests {
34453445
);
34463446
}
34473447

3448+
#[test]
3449+
fn fixed_primitive_streams_use_four_kibibyte_buffer_budget() {
3450+
for (wit_type, expected_window) in [
3451+
("u8", 4096),
3452+
("s8", 4096),
3453+
("u16", 2048),
3454+
("s16", 2048),
3455+
("u32", 1024),
3456+
("s32", 1024),
3457+
("f32", 1024),
3458+
("char", 1024),
3459+
("u64", 512),
3460+
("s64", 512),
3461+
("f64", 512),
3462+
("bool", 64),
3463+
("string", 64),
3464+
] {
3465+
let wit = format!(
3466+
r#"
3467+
package a:b;
3468+
3469+
interface api {{
3470+
type scalar = {wit_type};
3471+
exchange: func(input: stream<scalar>) -> stream<scalar>;
3472+
}}
3473+
3474+
world runner {{ import api; }}
3475+
"#
3476+
);
3477+
let files = generate(&wit, "runner");
3478+
let generated = files
3479+
.iter()
3480+
.map(|(_, contents)| String::from_utf8_lossy(contents))
3481+
.collect::<Vec<_>>()
3482+
.join("\n");
3483+
assert!(
3484+
generated.contains(&format!("let read_count = if count < {expected_window}")),
3485+
"unexpected read window for {wit_type}: {generated}"
3486+
);
3487+
assert!(
3488+
generated.contains(&format!(
3489+
"let data_len = if data.length() < {expected_window}"
3490+
)),
3491+
"unexpected write window for {wit_type}: {generated}"
3492+
);
3493+
let compact = generated.split_whitespace().collect::<Vec<_>>().join(" ");
3494+
assert!(
3495+
compact.contains(&format!("Some(cleanup_value), {expected_window}, )")),
3496+
"sink did not receive the {expected_window}-element window for {wit_type}: \
3497+
{generated}"
3498+
);
3499+
}
3500+
3501+
let runtime = generate(
3502+
r#"
3503+
package a:b;
3504+
world runner {
3505+
import exchange: func(input: stream<u8>) -> stream<u8>;
3506+
}
3507+
"#,
3508+
"runner",
3509+
);
3510+
let sink = file(&runtime, "async-core/async_trait.mbt");
3511+
assert!(sink.contains("priv write_window_size : Int"));
3512+
assert!(sink.contains("data.length() < self.write_window_size"));
3513+
assert!(sink.contains("FixedArray::from_array(data[:length])"));
3514+
assert!(sink.contains("data[:length].to_fixedarray()"));
3515+
assert!(!sink.contains("FixedArray::makei"));
3516+
}
3517+
34483518
#[test]
34493519
fn async_export_background_group_name_is_deconflicted() {
34503520
let files = generate(

0 commit comments

Comments
 (0)