How to Use asio::experimental::concurrent_channel for Thread-Safe Async Communication
The asio::experimental::concurrent_channel provides a thread-safe, asynchronous FIFO queue that enables coroutines and threads to exchange data using async_send and async_receive without explicit synchronization.
The asio::experimental::concurrent_channel template in the chriskohlhoff/asio repository implements a multi-producer, multi-consumer channel for asynchronous data transfer. This experimental feature builds upon the generic basic_concurrent_channel implementation to provide a simplified interface for common executor patterns. Whether you are coordinating work between thread pools or implementing pipeline architectures, mastering this component enables robust asynchronous communication in modern C++ applications.
Understanding the Channel Architecture
The concurrent_channel is a type alias defined in include/asio/experimental/concurrent_channel.hpp that instantiates basic_concurrent_channel with automatically deduced template arguments. The underlying implementation in include/asio/experimental/basic_concurrent_channel.hpp provides the complete channel logic, including buffer management and thread-safety mechanisms.
A channel requires two template parameters:
- Executor type: Typically
asio::any_io_executorfor flexibility, or a specific executor likeasio::io_context::executor_type - Signature: A function signature describing the value type, such as
void(std::string)orvoid(int)
Instantiating a concurrent_channel
Create a channel by passing an executor to the constructor. The executor determines where completion handlers execute, while the signature defines the transferable data type.
#include <asio.hpp>
#include <asio/experimental/concurrent_channel.hpp>
asio::io_context ctx;
// Channel that transports strings
asio::experimental::concurrent_channel<
asio::any_io_executor,
void(std::string)
> chan(ctx.get_executor());
The any_io_executor type provides type erasure, allowing the channel to work across different executor types without templating your entire call stack.
Core Asynchronous Operations
The channel exposes four primary operations for data transfer and lifecycle management.
Sending Data with async_send and try_send
Use async_send to enqueue values asynchronously. If the channel buffer is full, the operation waits until space becomes available and then invokes the completion handler.
// Asynchronous send with coroutine support
co_await async_send(chan, "value", asio::use_awaitable);
// Non-blocking attempt
bool sent = chan.try_send("value");
Receiving Data with async_receive and try_receive
Use async_receive to dequeue values. If the channel is empty, the operation pauses until a producer sends data.
// Asynchronous receive with error handling
std::string value;
error_code ec;
co_await async_receive(chan, asio::redirect_error(asio::use_awaitable, ec), value);
Closing the Channel
Call close to signal end-of-stream. This causes pending async_receive operations to complete with asio::error::operation_aborted, allowing consumer loops to terminate cleanly.
chan.close();
Practical Implementation Examples
Coroutine-Based Producer-Consumer
This pattern uses C++20 coroutines with use_awaitable to create readable asynchronous code. The producer sends values and closes the channel when finished, while the consumer handles the closure signal.
#include <asio.hpp>
#include <asio/experimental/concurrent_channel.hpp>
#include <iostream>
using namespace asio;
using namespace asio::experimental;
using string_chan = concurrent_channel<any_io_executor, void(std::string)>;
awaitable<void> producer(string_chan& chan)
{
for (int i = 0; i < 5; ++i)
{
co_await async_send(chan, std::to_string(i), use_awaitable);
std::cout << "sent: " << i << "\n";
}
chan.close();
co_return;
}
awaitable<void> consumer(string_chan& chan)
{
for (;;)
{
std::string value;
error_code ec;
co_await async_receive(chan, redirect_error(use_awaitable, ec), value);
if (ec) break;
std::cout << "received: " << value << "\n";
}
co_return;
}
int main()
{
io_context ctx;
string_chan chan(ctx.get_executor());
co_spawn(ctx, producer(chan), detached);
co_spawn(ctx, consumer(chan), detached);
ctx.run();
}
Callback-Based Implementation
For codebases not using coroutines, the channel supports traditional Asio completion handlers.
#include <asio.hpp>
#include <asio/experimental/concurrent_channel.hpp>
#include <iostream>
using namespace asio;
using namespace asio::experimental;
using int_chan = concurrent_channel<any_io_executor, void(int)>;
void send_handler(const error_code& ec, std::size_t)
{
if (!ec) std::cout << "value sent successfully\n";
}
void receive_handler(const error_code& ec, std::size_t, int value)
{
if (!ec) std::cout << "got value: " << value << "\n";
}
int main()
{
io_context ctx;
int_chan chan(ctx.get_executor());
async_send(chan, 42, send_handler);
async_send(chan, 7, send_handler);
async_receive(chan, receive_handler);
async_receive(chan, receive_handler);
ctx.run();
}
Multi-Threaded Communication
The concurrent_channel is thread-safe, allowing producers and consumers to operate on different threads using the same executor.
#include <asio.hpp>
#include <asio/experimental/concurrent_channel.hpp>
#include <thread>
#include <iostream>
using namespace asio;
using namespace asio::experimental;
using double_chan = concurrent_channel<any_io_executor, void(double)>;
int main()
{
io_context ctx;
double_chan chan(ctx.get_executor());
std::thread consumer_thread([&]{ ctx.run(); });
std::thread producer_thread([&chan]{
for (int i = 0; i < 10; ++i)
{
chan.async_send(static_cast<double>(i) * 0.5,
[](const error_code&, std::size_t){});
}
chan.close();
});
for (;;)
{
double v;
error_code ec;
chan.async_receive(redirect_error(use_future, ec), v).wait();
if (ec) break;
std::cout << "got: " << v << "\n";
}
producer_thread.join();
consumer_thread.join();
}
Source Code Structure
The implementation spans two primary headers in the chriskohlhoff/asio repository:
-
include/asio/experimental/concurrent_channel.hpp: Defines theconcurrent_channeltype alias that automatically selects appropriate template arguments forbasic_concurrent_channel. -
include/asio/experimental/basic_concurrent_channel.hpp: Contains the full class template implementation, including the internal synchronization primitives and buffer management logic.
Unit tests demonstrating edge cases and multi-signature channels are available in src/tests/unit/experimental/concurrent_channel.cpp and src/tests/unit/experimental/basic_concurrent_channel.cpp.
Summary
- The asio::experimental::concurrent_channel provides a thread-safe, asynchronous FIFO queue for inter-coroutine and inter-thread communication.
- Instantiate the channel with an executor and a function signature describing the data type, such as
void(std::string). - Use async_send and async_receive for non-blocking data transfer, or try_send and try_receive for immediate attempts.
- Call close to signal channel termination, which propagates
operation_abortedto pending receivers. - The implementation lives in
include/asio/experimental/concurrent_channel.hppandinclude/asio/experimental/basic_concurrent_channel.hpp.
Frequently Asked Questions
What is the difference between concurrent_channel and basic_concurrent_channel?
concurrent_channel is a convenience type alias that instantiates basic_concurrent_channel with deduced template arguments. While basic_concurrent_channel requires explicit template parameters, concurrent_channel automatically selects the appropriate executor and signature types, simplifying the common single-executor use case.
Is concurrent_channel thread-safe?
Yes. The concurrent_channel is designed for multi-threaded environments. Producers can call async_send from any thread that holds a reference to the channel, and consumers can receive from different threads simultaneously. The internal synchronization in basic_concurrent_channel handles all locking automatically.
How do I handle channel closure in receivers?
When close is called, pending async_receive operations complete with the asio::error::operation_aborted error code. Check for this error in your completion handlers or coroutines to break out of receive loops cleanly. This pattern signals end-of-stream without requiring additional synchronization mechanisms.
Can I use concurrent_channel with executors other than io_context?
Yes. While io_context is common, you can use any Asio executor type, including asio::thread_pool or asio::strand. The any_io_executor type erasure allows the channel to work seamlessly across different executor types, making it suitable for heterogeneous execution environments where producers and consumers use different underlying execution resources.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →