Examples#
pw_rpc2: Next gen, zero-copy async RPC with end-to-end backpressure
The Echo example implements a service with one RPC of each type and serves it
over TCP. It consists of four standalone programs, which show the two ways to
write pw_rpc2 code: with C++20 coroutines, or with hand-written futures.
Program |
Description |
|---|---|
|
Implements the service with coroutines. |
|
Implements the service with futures. |
|
Calls each RPC from coroutines. |
|
Calls each RPC from futures and a task. |
Any client works with any server. The source code is in pw_rpc2/cpp/examples/echo.
Run the example#
Start a server:
bazelisk run //pw_rpc2/cpp/examples/echo:server_coro
In another terminal, run a client:
bazelisk run //pw_rpc2/cpp/examples/echo:client_coro
The client makes each RPC once and logs the responses:
INF Echo: hello
INF Repeat: ping
INF Repeat: ping
INF Repeat: ping
INF Collect: one two three
INF EchoStream: alpha
INF EchoStream: beta
The server logs each request:
INF Echo: hello
INF Repeat: ping x3
INF Collect: one
INF Collect: two
INF Collect: three
INF EchoStream: alpha
INF EchoStream: beta
Define the service#
echo.proto defines the Echo service:
syntax = "proto3";
package examples;
// A text message. `echo.pwpb_options` gives `msg` a fixed capacity.
message EchoMessage {
string msg = 1;
}
// Asks the server to send `msg` back `count` times.
message RepeatRequest {
string msg = 1;
uint32 count = 2;
}
// Demonstrates each of the four types of RPC.
service Echo {
// Unary: responds with the request.
rpc Echo(EchoMessage) returns (EchoMessage);
// Server streaming: sends `msg` back `count` times.
rpc Repeat(RepeatRequest) returns (stream EchoMessage);
// Client streaming: after the client finishes its stream, responds with the
// messages joined by spaces.
rpc Collect(stream EchoMessage) returns (EchoMessage);
// Bidirectional streaming: sends each message back as it arrives.
rpc EchoStream(stream EchoMessage) returns (stream EchoMessage);
}
The generated code uses pw_protobuf message structs. By default,
pw_protobuf encodes and decodes string fields with callbacks.
echo.pwpb_options gives the string fields a fixed capacity instead, so the
structs store them inline as pw::InlineString<64>:
examples.EchoMessage.msg max_size:64
examples.RepeatRequest.msg max_size:64
pwpb_proto_library generates the message structs, and
pwpb_rpc2_proto_library generates the service and client code:
pw_proto_filegroup(
name = "echo_proto_and_options",
srcs = ["echo.proto"],
options_files = ["echo.pwpb_options"],
)
proto_library(
name = "echo_proto",
srcs = [":echo_proto_and_options"],
import_prefix = "echo_pb",
strip_import_prefix = "/pw_rpc2/cpp/examples/echo",
)
pwpb_proto_library(
name = "echo_pwpb",
deps = [":echo_proto"],
)
pwpb_rpc2_proto_library(
name = "echo_pwpb_rpc2",
pwpb_proto_library_deps = [":echo_pwpb"],
deps = [":echo_proto"],
)
The generated header, echo_pb/echo.pwpb.rpc2.h, declares a base class for
implementing the service, examples::pw_rpc2::pwpb::Echo::Service, and a
client for calling it, examples::pw_rpc2::pwpb::Echo::Client.
Implement the service#
A service implementation derives from Echo::Service, passing itself as the
template argument. It implements each RPC either as a coroutine member function
or as a nested <Method>Future class. An RPC’s arguments depend on its method
type:
Type |
Arguments |
|---|---|
Unary |
|
Server streaming |
|
Client streaming |
|
Bidirectional streaming |
|
UnaryWriter::Finish()sends the response.Writer::Write()sends a message, andWriter::Finish()ends the stream.Reader::Read()receives the next message, orOUT_OF_RANGEonce the client has finished its stream.
RPC lifetime#
An RPC ends when the coroutine or future that implements it completes. If the RPC does work elsewhere, such as in another task, the coroutine or future must wait for that work to finish before completing. Otherwise, the RPC ends early. Readers and writers that the work holds remain safe to use, but their operations fail.
Conversely, once an RPC is finished or cancelled, pw_rpc2 destroys its
coroutine or future without waiting for it to complete. Do all of the RPC’s
work before finishing it.
Coroutines#
A coroutine RPC takes a CoroContext as its
first argument and returns pw::async2::Coro<void>. See
Coroutines for more about coroutines.
using EchoMessage = examples::pwpb::EchoMessage::Message;
using RepeatRequest = examples::pwpb::RepeatRequest::Message;
// Implements each RPC as a coroutine member function. Requests are taken by
// value so that they stay valid across `co_await`.
class EchoService : public examples::pw_rpc2::pwpb::Echo::Service<EchoService> {
public:
Coro<void> Echo(CoroContext,
EchoMessage request,
pw::rpc2::UnaryWriter<EchoMessage> responder) {
PW_LOG_INFO("Echo: %s", request.msg.c_str());
co_await responder.Finish(std::move(request));
}
Coro<void> Repeat(CoroContext,
RepeatRequest request,
pw::rpc2::Writer<EchoMessage> writer) {
PW_LOG_INFO("Repeat: %s x%u", request.msg.c_str(), request.count);
const EchoMessage response{.msg = request.msg};
for (uint32_t i = 0; i < request.count; ++i) {
pw::Status status = co_await writer.Write(response);
if (!status.ok()) {
co_return; // The call ended early, e.g. the client cancelled it.
}
}
co_await writer.Finish();
}
Coro<void> Collect(CoroContext,
pw::rpc2::Reader<EchoMessage> reader,
pw::rpc2::UnaryWriter<EchoMessage> responder) {
EchoMessage response;
while (true) {
pw::Result<EchoMessage> request = co_await reader.Read();
if (request.status().IsOutOfRange()) {
break; // The client finished its stream.
}
if (!request.ok()) {
co_return;
}
PW_LOG_INFO("Collect: %s", request->msg.c_str());
// Messages that don't fit in the response are truncated.
if (!response.msg.empty()) {
pw::string::Append(response.msg, " ").IgnoreError();
}
pw::string::Append(response.msg, request->msg).IgnoreError();
}
co_await responder.Finish(std::move(response));
}
Coro<void> EchoStream(CoroContext,
pw::rpc2::Reader<EchoMessage> reader,
pw::rpc2::Writer<EchoMessage> writer) {
while (true) {
pw::Result<EchoMessage> request = co_await reader.Read();
if (request.status().IsOutOfRange()) {
break; // The client finished its stream.
}
if (!request.ok()) {
co_return;
}
PW_LOG_INFO("EchoStream: %s", request->msg.c_str());
pw::Status status = co_await writer.Write(std::move(*request));
if (!status.ok()) {
co_return;
}
}
co_await writer.Finish();
}
};
Because the RPCs return Coro<void>, they can’t use PW_CO_TRY, which
returns a status with co_return. Instead, they check each status and
co_return early if the call has ended.
Futures#
Without coroutines, each RPC is a nested <Method>Future class. For each
call, pw_rpc2 constructs the future from the RPC’s arguments and polls it
until it completes. The class must be a future whose value_type is void.
using EchoMessage = examples::pwpb::EchoMessage::Message;
using RepeatRequest = examples::pwpb::RepeatRequest::Message;
// Implements each RPC with a nested `<Method>Future` type. For each call,
// pw_rpc2 constructs the future from the RPC's arguments and polls it until it
// completes.
//
// `<Method>Future` constructors can be declared in two ways, depending on
// whether the future needs a reference to the service instance:
// - `<Method>Future(EchoService& service, <RPC arguments>...)`
// - `<Method>Future(<RPC arguments>...)`
class EchoService : public examples::pw_rpc2::pwpb::Echo::Service<EchoService> {
public:
class EchoFuture {
public:
using value_type = void;
EchoFuture() = default;
EchoFuture(EchoMessage request,
pw::rpc2::UnaryWriter<EchoMessage> responder)
: state_(State::kStarting),
request_(std::move(request)),
responder_(std::move(responder)) {}
bool is_pendable() const {
return state_ != State::kUninitialized && state_ != State::kDone;
}
bool is_complete() const { return state_ == State::kDone; }
Poll<> Pend(Context& cx) {
if (state_ == State::kStarting) {
PW_LOG_INFO("Echo: %s", request_.msg.c_str());
response_future_ = responder_.Finish(std::move(request_));
state_ = State::kResponding;
}
PW_AWAIT(response_future_, cx);
state_ = State::kDone;
return Ready();
}
private:
enum class State { kUninitialized, kStarting, kResponding, kDone };
State state_ = State::kUninitialized;
EchoMessage request_;
pw::rpc2::UnaryWriter<EchoMessage> responder_;
pw::rpc2::WriteFuture<EchoMessage> response_future_;
};
The future for a streaming RPC is a state machine that awaits one operation at
a time. EchoStreamFuture reads a message, writes it back, and repeats until
the client finishes its stream. Then it finishes its own stream:
class EchoStreamFuture {
public:
using value_type = void;
EchoStreamFuture() = default;
EchoStreamFuture(pw::rpc2::Reader<EchoMessage> reader,
pw::rpc2::Writer<EchoMessage> writer)
: state_(State::kReading),
reader_(std::move(reader)),
writer_(std::move(writer)) {}
bool is_pendable() const {
return state_ != State::kUninitialized && state_ != State::kDone;
}
bool is_complete() const { return state_ == State::kDone; }
Poll<> Pend(Context& cx) {
while (true) {
switch (state_) {
case State::kReading: {
if (!read_future_.is_pendable()) {
read_future_ = reader_.Read();
}
PW_AWAIT(pw::Result<EchoMessage> request, read_future_, cx);
if (request.status().IsOutOfRange()) {
state_ = State::kFinishing; // The client finished its stream.
break;
}
if (!request.ok()) {
state_ = State::kDone;
return Ready();
}
PW_LOG_INFO("EchoStream: %s", request->msg.c_str());
write_future_ = writer_.Write(std::move(*request));
state_ = State::kWriting;
break;
}
case State::kWriting: {
PW_AWAIT(pw::Status status, write_future_, cx);
if (!status.ok()) {
state_ = State::kDone;
return Ready();
}
state_ = State::kReading;
break;
}
case State::kFinishing: {
if (!finish_future_.is_pendable()) {
finish_future_ = writer_.Finish();
}
PW_AWAIT(finish_future_, cx);
state_ = State::kDone;
return Ready();
}
case State::kUninitialized:
case State::kDone:
PW_CRASH("Polled a future that is not pendable");
}
}
}
private:
enum class State { kUninitialized, kReading, kWriting, kFinishing, kDone };
State state_ = State::kUninitialized;
pw::rpc2::Reader<EchoMessage> reader_;
pw::rpc2::Writer<EchoMessage> writer_;
pw::rpc2::ReadFuture<EchoMessage> read_future_;
pw::rpc2::WriteFuture<EchoMessage> write_future_;
pw::rpc2::WriteFuture<> finish_future_;
};
See Implementing a composite future for more about writing futures like these.
Serve the service#
A pw::rpc2::Server serves its registered services on the connections that
its listeners accept. Both servers set it up the same way:
pw::Allocator& allocator = pw::allocator::GetLibCAllocator();
BasicDispatcher dispatcher;
pw::transport::FramedTcpListener listener(allocator);
if (!listener.Listen(port).ok()) {
return 1;
}
PW_LOG_INFO("Listening on port %u", static_cast<unsigned>(listener.port()));
EchoService service;
pw::rpc2::Server server(allocator, dispatcher);
PW_CHECK_OK(server.RegisterService(service));
PW_CHECK_OK(server.RegisterListenerBlocking(listener));
server.Start();
// Serves RPCs until the process is terminated.
dispatcher.RunToCompletion();
Start() runs the server on the dispatcher. The service’s RPCs run on the
dispatcher’s thread, so state that only the RPCs access does not need locking.
Call the RPCs#
pw::rpc2::Client::Connect() connects to a server and resolves to a
pw::rpc2::Client, which is a handle to the connection. The generated
Echo::Client binds a pw::rpc2::Client to the Echo service. Both are
cheap to copy, so pass them by value.
The generated client creates functions for each RPC. Each function returns a future that starts the RPC and resolves to:
Type |
Result |
|---|---|
Unary |
|
Server streaming |
|
Client streaming |
|
Bidirectional streaming |
|
Writer::Finish() ends the client’s stream. Reader::Read() resolves to
OUT_OF_RANGE once the server has finished its stream.
Coroutines#
Each call is a coroutine that makes one RPC and logs the responses:
using EchoMessage = examples::pwpb::EchoMessage::Message;
// `pw::rpc2::Client` is a connection to a server, and the generated `Echo`
// client binds a `pw::rpc2::Client` to the `Echo` service. Both are cheap to
// copy, so pass them by value.
using EchoClient = examples::pw_rpc2::pwpb::Echo::Client;
// Unary RPC: sends one request and receives one response.
Coro<pw::Status> CallEcho(CoroContext, EchoClient client) {
PW_CO_TRY_ASSIGN(EchoMessage response,
co_await client.Echo({.msg = "hello"}));
PW_LOG_INFO("Echo: %s", response.msg.c_str());
co_return pw::OkStatus();
}
// Server streaming RPC: sends one request, then reads responses until the
// server finishes its stream.
Coro<pw::Status> CallRepeat(CoroContext, EchoClient client) {
PW_CO_TRY_ASSIGN(pw::rpc2::Reader<EchoMessage> reader,
co_await client.Repeat({.msg = "ping", .count = 3}));
while (true) {
pw::Result<EchoMessage> response = co_await reader.Read();
if (response.status().IsOutOfRange()) {
co_return pw::OkStatus(); // The server finished its stream.
}
PW_CO_TRY(response.status());
PW_LOG_INFO("Repeat: %s", response->msg.c_str());
}
}
// Client streaming RPC: streams requests, then receives one response.
Coro<pw::Status> CallCollect(CoroContext, EchoClient client) {
PW_CO_TRY_ASSIGN(auto call, co_await client.Collect());
pw::rpc2::Writer<EchoMessage>& writer = call.writer();
for (std::string_view word : {"one", "two", "three"}) {
PW_CO_TRY(co_await writer.Write({.msg = word}));
}
PW_CO_TRY(co_await writer.Finish());
PW_CO_TRY_ASSIGN(EchoMessage response, co_await call.response());
PW_LOG_INFO("Collect: %s", response.msg.c_str());
co_return pw::OkStatus();
}
// Bidirectional streaming RPC: streams requests and reads the server's stream
// of responses. The writer and reader are independent; they could also be used
// from separate tasks.
Coro<pw::Status> CallEchoStream(CoroContext, EchoClient client) {
PW_CO_TRY_ASSIGN(auto call, co_await client.EchoStream());
pw::rpc2::Writer<EchoMessage>& writer = call.writer();
for (std::string_view word : {"alpha", "beta"}) {
PW_CO_TRY(co_await writer.Write({.msg = word}));
}
PW_CO_TRY(co_await writer.Finish());
pw::rpc2::Reader<EchoMessage>& reader = call.reader();
while (true) {
pw::Result<EchoMessage> response = co_await reader.Read();
if (response.status().IsOutOfRange()) {
co_return pw::OkStatus(); // The server finished its stream.
}
PW_CO_TRY(response.status());
PW_LOG_INFO("EchoStream: %s", response->msg.c_str());
}
}
A top-level coroutine connects to the server, makes the calls, and closes the connection:
// Connects to the server, makes each type of RPC, and closes the connection.
Coro<pw::Status> RunClient(
CoroContext cx,
Dispatcher& dispatcher,
pw::Allocator& allocator,
pw::transport::ReliableDatagramConnector& connector) {
PW_CO_TRY_ASSIGN(
pw::rpc2::Client client,
co_await pw::rpc2::Client::Connect(dispatcher, allocator, connector));
EchoClient echo_client(client);
PW_CO_TRY(co_await CallEcho(cx, echo_client));
PW_CO_TRY(co_await CallRepeat(cx, echo_client));
PW_CO_TRY(co_await CallCollect(cx, echo_client));
PW_CO_TRY(co_await CallEchoStream(cx, echo_client));
co_return co_await client.Close();
}
main() runs the top-level coroutine as a FutureTask. RunToCompletion() returns once the task has
finished and the connection has closed:
pw::Allocator& allocator = pw::allocator::GetLibCAllocator();
BasicDispatcher dispatcher;
pw::transport::FramedTcpConnector connector(allocator, "127.0.0.1", port);
CoroContext coro_cx(allocator);
FutureTask task(RunClient(coro_cx, dispatcher, allocator, connector));
dispatcher.Post(task);
dispatcher.RunToCompletion();
return task.value().ok() ? 0 : 1;
Futures#
Without coroutines, each streaming call is a composite future. RepeatCall
starts the RPC, then reads responses until the server finishes its stream:
using EchoMessage = examples::pwpb::EchoMessage::Message;
using RepeatRequest = examples::pwpb::RepeatRequest::Message;
// `pw::rpc2::Client` is a connection to a server, and the generated `Echo`
// client binds a `pw::rpc2::Client` to the `Echo` service. Both are cheap to
// copy, so pass them by value.
using EchoClient = examples::pw_rpc2::pwpb::Echo::Client;
// Server streaming RPC: sends one request, then reads responses until the
// server finishes its stream.
class RepeatCall {
public:
using value_type = pw::Status;
RepeatCall() = default;
explicit RepeatCall(EchoClient client)
: state_(State::kCalling), client_(std::move(client)) {}
bool is_pendable() const {
return state_ != State::kUninitialized && state_ != State::kDone;
}
bool is_complete() const { return state_ == State::kDone; }
Poll<pw::Status> Pend(Context& cx) {
while (true) {
switch (state_) {
case State::kCalling: {
if (!call_future_.is_pendable()) {
call_future_ = client_.Repeat({.msg = "ping", .count = 3});
}
PW_AWAIT(pw::Result<pw::rpc2::Reader<EchoMessage>> reader,
call_future_,
cx);
if (!reader.ok()) {
return Complete(reader.status());
}
reader_ = std::move(*reader);
state_ = State::kReading;
break;
}
case State::kReading: {
if (!read_future_.is_pendable()) {
read_future_ = reader_.Read();
}
PW_AWAIT(pw::Result<EchoMessage> response, read_future_, cx);
if (response.status().IsOutOfRange()) {
return Complete(pw::OkStatus()); // The server finished its stream.
}
if (!response.ok()) {
return Complete(response.status());
}
PW_LOG_INFO("Repeat: %s", response->msg.c_str());
break;
}
case State::kUninitialized:
case State::kDone:
PW_CRASH("Polled a future that is not pendable");
}
}
}
private:
enum class State { kUninitialized, kCalling, kReading, kDone };
Poll<pw::Status> Complete(pw::Status status) {
state_ = State::kDone;
return Ready(status);
}
State state_ = State::kUninitialized;
EchoClient client_;
pw::rpc2::ServerStreamFuture<RepeatRequest, EchoMessage> call_future_;
pw::rpc2::Reader<EchoMessage> reader_;
pw::rpc2::ReadFuture<EchoMessage> read_future_;
};
CollectCall and EchoStreamCall follow the same pattern. A
pw::async2::Task connects to the server, runs the calls in sequence, and
closes the connection. The unary Echo call is simpler; the task awaits the
RPC’s response future directly:
// Connects to the server, makes each type of RPC, and closes the connection.
class ClientTask : public Task {
public:
ClientTask(Dispatcher& dispatcher,
pw::Allocator& allocator,
pw::transport::ReliableDatagramConnector& connector)
: dispatcher_(dispatcher), allocator_(allocator), connector_(connector) {}
// Returns the result of the RPCs once the task has completed.
pw::Status status() const { return status_; }
private:
enum class State {
kConnecting,
kEcho,
kRepeat,
kCollect,
kEchoStream,
kClosing,
};
Poll<> DoPend(Context& cx) override {
while (true) {
switch (state_) {
case State::kConnecting: {
if (!connect_future_.is_pendable()) {
connect_future_ =
pw::rpc2::Client::Connect(dispatcher_, allocator_, connector_);
}
PW_AWAIT(pw::Result<pw::rpc2::Client> client, connect_future_, cx);
if (!client.ok()) {
status_ = client.status();
return Ready();
}
client_ = std::move(*client);
echo_client_ = EchoClient(client_);
state_ = State::kEcho;
break;
}
// Unary RPC: sends one request and receives one response.
case State::kEcho: {
if (!echo_future_.is_pendable()) {
echo_future_ = echo_client_.Echo({.msg = "hello"});
}
PW_AWAIT(pw::Result<EchoMessage> response, echo_future_, cx);
if (response.ok()) {
PW_LOG_INFO("Echo: %s", response->msg.c_str());
}
status_ = response.status();
state_ = status_.ok() ? State::kRepeat : State::kClosing;
break;
}
case State::kRepeat: {
if (!repeat_.is_pendable()) {
repeat_ = RepeatCall(echo_client_);
}
PW_AWAIT(status_, repeat_, cx);
state_ = status_.ok() ? State::kCollect : State::kClosing;
break;
}
case State::kCollect: {
if (!collect_.is_pendable()) {
collect_ = CollectCall(echo_client_);
}
PW_AWAIT(status_, collect_, cx);
state_ = status_.ok() ? State::kEchoStream : State::kClosing;
break;
}
case State::kEchoStream: {
if (!echo_stream_.is_pendable()) {
echo_stream_ = EchoStreamCall(echo_client_);
}
PW_AWAIT(status_, echo_stream_, cx);
state_ = State::kClosing;
break;
}
case State::kClosing: {
if (!close_future_.is_pendable()) {
close_future_ = client_.Close();
}
PW_AWAIT(pw::Status close_status, close_future_, cx);
status_.Update(close_status);
return Ready();
}
}
}
}
Dispatcher& dispatcher_;
pw::Allocator& allocator_;
pw::transport::ReliableDatagramConnector& connector_;
State state_ = State::kConnecting;
pw::Status status_;
pw::rpc2::Client client_;
EchoClient echo_client_;
pw::rpc2::ClientFuture connect_future_;
pw::rpc2::UnaryFuture<EchoMessage, EchoMessage> echo_future_;
RepeatCall repeat_;
CollectCall collect_;
EchoStreamCall echo_stream_;
pw::rpc2::ControlFuture close_future_;
};
main() posts the task to a pw::async2::BasicDispatcher and runs the
dispatcher until the task has finished and the connection has closed.