How to alter the duration of a Tokio timer stream that has been combined with another stream?

Viewed 189

I have three streams:

  1. a timer stream (defined below)
  2. a broadcast channel stream (defined below)
  3. a WebSocket stream (not defined here but I don't think it's necessary. It's based on tokio-tungstenite)

I join these 3 streams together with a select.

I want commands I send down the broadcast channel to alter the timer - I want to be able to re-adjust the tokio timer once the streams are fused together.

This is my timer stream:

fn timer_stream(dur: u64) -> impl futures::stream::Stream<Item = Input> {
    let task_time = tokio::time::Instant::now();

    tokio::time::interval_at(task_time, Duration::from_secs(dur))
        .map(move |_| Input::Command_first(task_time))
}

I want to be able to change the task_time and the dur.

This is how I combine my streams together:

let mut websocket_timer = select(exchange_websocket_stream, timer_stream(3));
let mut combined = select(websocket_timer, command_receiver);

command_receiver is defined with use tokio::sync::broadcast; like so:

let (tx_tcp_commands, mut rx_tcp_commands) = broadcast::channel::<Input>(16);
command_receiver = tx_tcp_commands.subscribe()

Input is:

enum Input {
    Command_second(tokio::time::Instant),
    Command_first(tokio::time::Instant),
    Thing,
}

I then put the above in a tokio::main executor and do

loop {
    match combined.next().await {
        None => break,
        Some(Input::Command(t)) => etc,
    }
}

This works as a general stream and I can match on commands sent down to command_receiver, I just need to know how to change the timer parameters.

The loop is processing a very high amount of WebSocket messages so it has to be efficient.

How can I accomplish this?

0 Answers
Related