Select

So far, when we wanted to add concurrency to the system, we spawned a new task. We will now cover some additional ways to concurrently execute asynchronous code with Tokio.

tokio::select!

The tokio::select! macro allows waiting on multiple async computations and returns when a single computation completes.

For example:

  1. use tokio::sync::oneshot;
  2. #[tokio::main]
  3. async fn main() {
  4. let (tx1, rx1) = oneshot::channel();
  5. let (tx2, rx2) = oneshot::channel();
  6. tokio::spawn(async {
  7. let _ = tx1.send("one");
  8. });
  9. tokio::spawn(async {
  10. let _ = tx2.send("two");
  11. });
  12. tokio::select! {
  13. val = rx1 => {
  14. println!("rx1 completed first with {:?}", val);
  15. }
  16. val = rx2 => {
  17. println!("rx2 completed first with {:?}", val);
  18. }
  19. }
  20. }

Two oneshot channels are used. Either channel could complete first. The select! statement awaits on both channels and binds val to the value returned by the task. When either tx1 or tx2 complete, the associated block is executed.

The branch that does not complete is dropped. In the example, the computation is awaiting the oneshot::Receiver for each channel. The oneshot::Receiver for the channel that did not complete yet is dropped.

Cancellation

With asynchronous Rust, cancellation is performed by dropping a future. Recall from “Async in depth”, async Rust operation are implemented using futures and futures are lazy. The operation only proceeds when the future is polled. If the future is dropped, the operation cannot proceed because all associated state has been dropped.

That said, sometimes an asynchronous operation will spawn background tasks or start other operation that run in the background. For example, in the above example, a task is spawned to send a message back. Usually, the task will perform some computation to generate the value.

Futures or other types can implement Drop to cleanup background resources. Tokio’s oneshot::Receiver implements Drop by sending a closed notification to the Sender half. The sender half can receive this notification and abort the in-progress operation by dropping it.

  1. use tokio::sync::oneshot;
  2. async fn some_operation() -> String {
  3. // Compute value here
  4. }
  5. #[tokio::main]
  6. async fn main() {
  7. let (mut tx1, rx1) = oneshot::channel();
  8. let (tx2, rx2) = oneshot::channel();
  9. tokio::spawn(async {
  10. // Select on the operation and the oneshot's
  11. // `closed()` notification.
  12. tokio::select! {
  13. val = some_operation() => {
  14. let _ = tx1.send(val);
  15. }
  16. _ = tx1.closed() => {
  17. // `some_operation()` is canceled, the
  18. // task completes and `tx1` is dropped.
  19. }
  20. }
  21. });
  22. tokio::spawn(async {
  23. let _ = tx2.send("two");
  24. });
  25. tokio::select! {
  26. val = rx1 => {
  27. println!("rx1 completed first with {:?}", val);
  28. }
  29. val = rx2 => {
  30. println!("rx2 completed first with {:?}", val);
  31. }
  32. }
  33. }

The Future implementation

To help better understand how select! works, let’s look at what a hypothetical Future implementation would look like. This is a simplified version. In practice, select! includes additional functionality like randomly selecting the branch to poll first.

  1. use tokio::sync::oneshot;
  2. use std::future::Future;
  3. use std::pin::Pin;
  4. use std::task::{Context, Poll};
  5. struct MySelect {
  6. rx1: oneshot::Receiver<&'static str>,
  7. rx2: oneshot::Receiver<&'static str>,
  8. }
  9. impl Future for MySelect {
  10. type Output = ();
  11. fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<()> {
  12. if let Poll::Ready(val) = Pin::new(&mut self.rx1).poll(cx) {
  13. println!("rx1 completed first with {:?}", val);
  14. return Poll::Ready(());
  15. }
  16. if let Poll::Ready(val) = Pin::new(&mut self.rx2).poll(cx) {
  17. println!("rx2 completed first with {:?}", val);
  18. return Poll::Ready(());
  19. }
  20. Poll::Pending
  21. }
  22. }
  23. #[tokio::main]
  24. async fn main() {
  25. let (tx1, rx1) = oneshot::channel();
  26. let (tx2, rx2) = oneshot::channel();
  27. // use tx1 and tx2
  28. MySelect {
  29. rx1,
  30. rx2,
  31. }.await;
  32. }

The MySelect future contains the futures from each branch. When MySelect is polled, the first branch is polled. If it is ready, the value is used and MySelect completes. After .await receives the output from a future, the future is dropped. This results in the futures for both branches to be dropped. As one branch did not complete, the operation is effectively cancelled.

Remember from the previous section:

When a future returns Poll::Pending, it must ensure the waker is signalled at some point in the future. Forgetting to do this results in the task hanging indefinitely.

There is no explicit usage of the Context argument in the MySelect implementation. Instead, the waker requirement is met by passing cx to the inner futures. As the inner future must also meet the waker requirement, by only returning Poll::Pending when receiving Poll::Pending from an inner future, MySelect also meets the waker requirement.

Syntax

The select! macro can handle more than two branches. The current limit is 64 branches. Each branch is structured as:

  1. <pattern> = <async expression> => <handler>,

When the select macro is evaluated, all the <async expression>s are aggregated and executed concurrently. When an expression completes, the result is matched against <pattern>. If the result matches the pattern, then all remaining async expressions are dropped and <handler> is executed. The <handler> expression has access to any bindings established by <pattern>.

The basic case is <pattern> is a variable name, the result of the async expression is bound to the variable name and <handler> has access to that variable. This is why, in the original example, val was used for <pattern> and <handler> was able to access val.

If <pattern> does not match the result of the async computation, then the remaining async expressions continue to execute concurrently until the next one completes. At this time, the same logic is applied to that result.

Because select! takes any async expression, it is possible to define more complicated computations to select on.

Here, we select on the output of a oneshot channel and a TCP connection.

  1. use tokio::net::TcpStream;
  2. use tokio::sync::oneshot;
  3. #[tokio::main]
  4. async fn main() {
  5. let (tx, rx) = oneshot::channel();
  6. // Spawn a task that sends a message over the oneshot
  7. tokio::spawn(async move {
  8. tx.send("done").unwrap();
  9. });
  10. tokio::select! {
  11. socket = TcpStream::connect("localhost:3465") => {
  12. println!("Socket connected {:?}", socket);
  13. }
  14. msg = rx => {
  15. println!("received message first {:?}", msg);
  16. }
  17. }
  18. }

Here, we select on a oneshot and accepting sockets from a TcpListener.

  1. use tokio::net::TcpListener;
  2. use tokio::sync::oneshot;
  3. use std::io;
  4. #[tokio::main]
  5. async fn main() -> io::Result<()> {
  6. let (tx, rx) = oneshot::channel();
  7. tokio::spawn(async move {
  8. tx.send(()).unwrap();
  9. });
  10. let mut listener = TcpListener::bind("localhost:3465").await?;
  11. tokio::select! {
  12. _ = async {
  13. loop {
  14. let (socket, _) = listener.accept().await?;
  15. tokio::spawn(async move { process(socket) });
  16. }
  17. // Help the rust type inferencer out
  18. Ok::<_, io::Error>(())
  19. } => {}
  20. _ = rx => {
  21. println!("terminating accept loop");
  22. }
  23. }
  24. Ok(())
  25. }

The accept loop runs until an error is encountered or rx receives a value. The _ pattern indicates that we have no interest in the return value of the async computation.

Return value

The tokio::select! macro returns the result of the evaluated <handler> expression.

  1. async fn computation1() -> String {
  2. // .. computation
  3. }
  4. async fn computation2() -> String {
  5. // .. computation
  6. }
  7. #[tokio::main]
  8. async fn main() {
  9. let out = tokio::select! {
  10. res1 = computation1() => res1,
  11. res2 = computation2() => res2,
  12. };
  13. println!("Got = {}", out);
  14. }

Because of this, it is required that the <handler> expression for each branch evaluates to the same type. If the output of a select! expression is not needed, it is good practice to have the expression evaluate to ().

Errors

Using the ? operator propagates the error from the expression. How this works depends on whether ? is used from an async expression or from a handler. Using ? in an async expression propagates the error out of the async expression. This makes the output of the async expression a Result. Using ? from a handler immediately propagates the error out of the select! expression. Let’s look at the accept loop example again:

  1. use tokio::net::TcpListener;
  2. use tokio::sync::oneshot;
  3. use std::io;
  4. #[tokio::main]
  5. async fn main() -> io::Result<()> {
  6. // [setup `rx` oneshot channel]
  7. let listener = TcpListener::bind("localhost:3465").await?;
  8. tokio::select! {
  9. res = async {
  10. loop {
  11. let (socket, _) = listener.accept().await?;
  12. tokio::spawn(async move { process(socket) });
  13. }
  14. // Help the rust type inferencer out
  15. Ok::<_, io::Error>(())
  16. } => {
  17. res?;
  18. }
  19. _ = rx => {
  20. println!("terminating accept loop");
  21. }
  22. }
  23. Ok(())
  24. }

Notice listener.accept().await?. The ? operator propagates the error out of that expression and to the res binding. On an error, res will be set to Err(_). Then, in the handler, the ? operator is used again. The res? statement will propagate an error out of the main function.

Pattern matching

Recall that the select! macro branch syntax was defined as:

  1. <pattern> = <async expression> => <handler>,

So far, we have only used variable bindings for <pattern>. However, any Rust pattern can be used. For example, say we are receiving from multiple MPSC channels, we might do something like this:

  1. use tokio::sync::mpsc;
  2. #[tokio::main]
  3. async fn main() {
  4. let (mut tx1, mut rx1) = mpsc::channel(128);
  5. let (mut tx2, mut rx2) = mpsc::channel(128);
  6. tokio::spawn(async move {
  7. // Do something w/ `tx1` and `tx2`
  8. });
  9. tokio::select! {
  10. Some(v) = rx1.recv() => {
  11. println!("Got {:?} from rx1", v);
  12. }
  13. Some(v) = rx2.recv() => {
  14. println!("Got {:?} from rx2", v);
  15. }
  16. else => {
  17. println!("Both channels closed");
  18. }
  19. }
  20. }

In this example, the select! expression waits on receiving a value from rx1 and rx2. If a channel closes, recv() returns None. This does not match the pattern and the branch is disabled. The select! expression will continue waiting on the remaining branches.

Notice that this select! expression includes an else branch. The select! expression must evaluate to a value. When using pattern matching, it is possible that none of the branches match their associated patterns. If this happens, the else branch is evaluated.

Borrowing

When spawning tasks, the spawned async expression must own all of its data. The select! macro does not have this limitation. Each branch’s async expression may borrow data and operate concurrently. Following Rust’s borrow rules, multiple async expressions may immutably borrow a single piece of data or a single async expression may mutably borrow a piece of data.

Let’s look at some examples. Here, we simultaneously send the same data to two different TCP destinations.

  1. use tokio::io::AsyncWriteExt;
  2. use tokio::net::TcpStream;
  3. use std::io;
  4. use std::net::SocketAddr;
  5. async fn race(
  6. data: &[u8],
  7. addr1: SocketAddr,
  8. addr2: SocketAddr
  9. ) -> io::Result<()> {
  10. tokio::select! {
  11. Ok(_) = async {
  12. let mut socket = TcpStream::connect(addr1).await?;
  13. socket.write_all(data).await?;
  14. Ok::<_, io::Error>(())
  15. } => {}
  16. Ok(_) = async {
  17. let mut socket = TcpStream::connect(addr2).await?;
  18. socket.write_all(data).await?;
  19. Ok::<_, io::Error>(())
  20. } => {}
  21. else => {}
  22. };
  23. Ok(())
  24. }

The data variable is being borrowed immutably from both async expressions. When one of the operations completes successfully, the other one is dropped. Because we pattern match on Ok(_), if an expression fails, the other one continues to execute.

When it comes to each branch’s <handler>, select! guarantees that only a single <handler> runs. Because of this, each <handler> may mutably borrow the same data.

For example this modifies out in both handlers:

  1. use tokio::sync::oneshot;
  2. #[tokio::main]
  3. async fn main() {
  4. let (tx1, rx1) = oneshot::channel();
  5. let (tx2, rx2) = oneshot::channel();
  6. let mut out = String::new();
  7. tokio::spawn(async move {
  8. // Send values on `tx1` and `tx2`.
  9. });
  10. tokio::select! {
  11. _ = rx1 => {
  12. out.push_str("rx1 completed");
  13. }
  14. _ = rx2 => {
  15. out.push_str("rx2 completed");
  16. }
  17. }
  18. println!("{}", out);
  19. }

Loops

The select! macro is often used in loops. This section will go over some examples to show common ways of using the select! macro in a loop. We start by selecting over multiple channels:

  1. use tokio::sync::mpsc;
  2. #[tokio::main]
  3. async fn main() {
  4. let (tx1, mut rx1) = mpsc::channel(128);
  5. let (tx2, mut rx2) = mpsc::channel(128);
  6. let (tx3, mut rx3) = mpsc::channel(128);
  7. loop {
  8. let msg = tokio::select! {
  9. Some(msg) = rx1.recv() => msg,
  10. Some(msg) = rx2.recv() => msg,
  11. Some(msg) = rx3.recv() => msg,
  12. else => { break }
  13. };
  14. println!("Got {}", msg);
  15. }
  16. println!("All channels have been closed.");
  17. }

This example selects over the three channel receivers. When a message is received on any channel, it is written to STDOUT. When a channel is closed, recv() returns with None. By using pattern matching, the select! macro continues waiting on the remaining channels. When all channels are closed, the else branch is evaluated and the loop is terminated.

The select! macro randomly picks branches to check first for readiness. When multiple channels have pending values, a random channel will be picked to receive from. This is to handle the case where the receive loop processes messages slower than they are pushed into the channels, meaning that the channels start to fill up. If select! did not randomly pick a branch to check first, on each iteration of the loop, rx1 would be checked first. If rx1 always contained a new message, the remaining channels would never be checked.

If when select! is evaluated, multiple channels have pending messages, only one channel has a value popped. All other channels remain untouched, and their messages stay in those channels until the next loop iteration. No messages are lost.

Resuming an async operation

Now we will show how to run an asynchronous operation across multiple calls to select!. In this example, we have an MPSC channel with item type i32, and an asynchronous function. We want to run the asynchronous function until it completes or an even integer is received on the channel.

  1. async fn action() {
  2. // Some asynchronous logic
  3. }
  4. #[tokio::main]
  5. async fn main() {
  6. let (mut tx, mut rx) = tokio::sync::mpsc::channel(128);
  7. let operation = action();
  8. tokio::pin!(operation);
  9. loop {
  10. tokio::select! {
  11. _ = &mut operation => break,
  12. Some(v) = rx.recv() => {
  13. if v % 2 == 0 {
  14. break;
  15. }
  16. }
  17. }
  18. }
  19. }

Note how, instead of calling action() in the select! macro, it is called outside the loop. The return of action() is assigned to operation without calling .await. Then we call tokio::pin! on operation.

Inside the select! loop, instead of passing in operation, we pass in &mut operation. The operation variable is tracking the in-flight asynchronous operation. Each iteration of the loop uses the same operation instead of issuing a new call to action().

The other select! branch receives a message from the channel. If the message is even, we are done looping. Otherwise, start the select! again.

This is the first time we use tokio::pin!. We aren’t going to get into the details of pinning yet. The thing to note is that, to .await a reference, the value being referenced must be pinned or implement Unpin.

If we remove the tokio::pin! line and try to compile, we get the following error:

  1. error[E0599]: no method named `poll` found for struct
  2. `std::pin::Pin<&mut &mut impl std::future::Future>`
  3. in the current scope
  4. --> src/main.rs:16:9
  5. |
  6. 16 | / tokio::select! {
  7. 17 | | _ = &mut operation => break,
  8. 18 | | Some(v) = rx.recv() => {
  9. 19 | | if v % 2 == 0 {
  10. ... |
  11. 22 | | }
  12. 23 | | }
  13. | |_________^ method not found in
  14. | `std::pin::Pin<&mut &mut impl std::future::Future>`
  15. |
  16. = note: the method `poll` exists but the following trait bounds
  17. were not satisfied:
  18. `impl std::future::Future: std::marker::Unpin`
  19. which is required by
  20. `&mut impl std::future::Future: std::future::Future`

Although we covered Future in the previous chapter, this error still isn’t very clear. If you hit such an error about Future not being implemented when attempting to call .await on a reference, then the future probably needs to be pinned.

Read more about Pin on the standard library.

Modifying a branch

Let’s look at a slightly more complicated loop. We have:

  1. A channel of i32 values.
  2. An async operation to perform on i32 values.

The logic we want to implement is:

  1. Wait for an even number on the channel.
  2. Start the asynchronous operation using the even number as input.
  3. Wait for the operation, but at the same time listen for more even numbers on the channel.
  4. If a new even number is received before the existing operation completes, abort the existing operation and start it over with the new even number.
  1. async fn action(input: Option<i32>) -> Option<String> {
  2. // If the input is `None`, return `None`.
  3. // This could also be written as `let i = input?;`
  4. let i = match input {
  5. Some(input) => input,
  6. None => return None,
  7. };
  8. // async logic here
  9. }
  10. #[tokio::main]
  11. async fn main() {
  12. let (mut tx, mut rx) = tokio::sync::mpsc::channel(128);
  13. let mut done = false;
  14. let operation = action(None);
  15. tokio::pin!(operation);
  16. tokio::spawn(async move {
  17. let _ = tx.send(1).await;
  18. let _ = tx.send(3).await;
  19. let _ = tx.send(2).await;
  20. });
  21. loop {
  22. tokio::select! {
  23. res = &mut operation, if !done => {
  24. done = true;
  25. if let Some(v) = res {
  26. println!("GOT = {}", v);
  27. return;
  28. }
  29. }
  30. Some(v) = rx.recv() => {
  31. if v % 2 == 0 {
  32. // `.set` is a method on `Pin`.
  33. operation.set(action(Some(v)));
  34. done = false;
  35. }
  36. }
  37. }
  38. }
  39. }

We use a similar strategy as the previous example. The async fn is called outside of the loop and assigned to operation. The operation variable is pinned. The loop selects on both operation and the channel receiver.

Notice how action takes Option<i32> as an argument. Before we receive the first even number, we need to instantiate operation to something. We make action take Option and return Option. If None is passed in, None is returned. The first loop iteration, operation completes immediately with None.

This example uses some new syntax. The first branch includes , if !done. This is a branch precondition. Before explaining how it works, let’s look at what happens if the precondition is omitted. Leaving out , if !done and running the example results in the following output:

  1. thread 'main' panicked at '`async fn` resumed after completion', src/main.rs:1:55
  2. note: run with `RUST_BACKTRACE=1` environment variable to display a backtrace

This error happens when attempting to use operation after it has already completed. Usually, when using .await, the value being awaited is consumed. In this example, we await on a reference. This means operation is still around after it has completed.

To avoid this panic, we must take care to disable the first branch if operation has completed. The done variable is used to track whether or not operation completed. A select! branch may include a precondition. This precondition is checked before select! awaits on the branch. If the condition evaluates to false then the branch is disabled. The done variable is initialized to false. When operation completes, done is set to true. The next loop iteration will disable the operation branch. When an even message is received from the channel, operation is reset and done is set to false.

Per-task concurrency

Both tokio::spawn and select! enable running concurrent asynchronous operations. However, the strategy used to run concurrent operations differs. The tokio::spawn function takes an asynchronous operation and spawns a new task to run it. A task is the object that the Tokio runtime schedules. Two different tasks are scheduled independently by Tokio. They may run simultaneously on different operating system threads. Because of this, a spawned task has the same restriction as a a spawned thread: no borrowing.

The select! macro runs all branches concurrently on the same task. Because all branches of the select! macro are executed on the same task, they will never run simultaneously. The select! macro multiplexes asynchronous operations on a single task.