slightly change example
Some checks reported errors
continuous-integration/drone/push Build was killed
continuous-integration/drone/pr Build was killed

This commit is contained in:
Alex 2022-09-12 17:19:26 +02:00
parent 0f799a7768
commit f0326607ee
Signed by: lx
GPG key ID: 0E496D15096376BE

View file

@ -142,6 +142,7 @@ impl Example {
example_field, example_field,
hex::encode(id) hex::encode(id)
); );
// Fake data stream with some delays in item production
let stream = let stream =
Box::pin(stream::iter([100, 200, 300, 400]).then(|x| async move { Box::pin(stream::iter([100, 200, 300, 400]).then(|x| async move {
tokio::time::sleep(Duration::from_millis(500)).await; tokio::time::sleep(Duration::from_millis(500)).await;
@ -191,16 +192,10 @@ impl StreamingEndpointHandler<ExampleMessage> for Example {
msg.msg() msg.msg()
); );
let source_stream = msg.take_stream().unwrap(); let source_stream = msg.take_stream().unwrap();
// Return same stream with 300ms delay
let new_stream = Box::pin(source_stream.then(|x| async move { let new_stream = Box::pin(source_stream.then(|x| async move {
info!(
"Handler: stream got bytes {:?}",
x.as_ref().map(|b| b.len())
);
tokio::time::sleep(Duration::from_millis(300)).await; tokio::time::sleep(Duration::from_millis(300)).await;
Ok(Bytes::from(vec![ x
10u8;
x.map(|b| b.len()).unwrap_or(1422) * 2
]))
})); }));
Resp::new(ExampleResponse { Resp::new(ExampleResponse {
example_field: false, example_field: false,