diff --git a/benches/rpc-tonic/src/main.rs b/benches/rpc-tonic/src/main.rs index 7372e2b0..4427abd1 100644 --- a/benches/rpc-tonic/src/main.rs +++ b/benches/rpc-tonic/src/main.rs @@ -41,14 +41,16 @@ impl BenchService for BenchServiceImpl { let req = request.into_inner(); let count = req.response_count; let payload = req.payload; - let stream = futures::stream::unfold(0, move |i| { - let payload = payload.clone(); - async move { - if i >= count { - return None; - } - Some((Ok(BenchResponse { payload }), i + 1)) + let stream = futures::stream::unfold((payload, 0), move |(payload, i)| async move { + if i >= count { + return None; } + // Each item owns its payload, as a handler that builds a response + // per message would; the copy is part of the measured work. + let item = BenchResponse { + payload: payload.clone(), + }; + Some((Ok(item), (payload, i + 1))) }); Ok(Response::new(Box::pin(stream))) } diff --git a/benches/rpc/src/lib.rs b/benches/rpc/src/lib.rs index 0f80fbd1..0e3da512 100644 --- a/benches/rpc/src/lib.rs +++ b/benches/rpc/src/lib.rs @@ -50,20 +50,17 @@ impl BenchService for BenchServiceImpl { let req = req.to_owned_message(); let count = req.response_count; let payload = req.payload; - let stream = futures::stream::unfold(0, move |i| { - let payload = payload.clone(); - async move { - if i >= count { - return None; - } - Some(( - Ok(BenchResponse { - payload, - ..Default::default() - }), - i + 1, - )) + let stream = futures::stream::unfold((payload, 0), move |(payload, i)| async move { + if i >= count { + return None; } + // Each item owns its payload, as a handler that builds a response + // per message would; the copy is part of the measured work. + let item = BenchResponse { + payload: payload.clone(), + ..Default::default() + }; + Some((Ok(item), (payload, i + 1))) }); Response::stream_ok(stream) }