blob: dc7aa456fa73f98d0cb08664db340f942cdc5630 [file]
// Copyright 2026 The Pigweed Authors
//
// Licensed under the Apache License, Version 2.0 (the "License"); you may not
// use this file except in compliance with the License. You may obtain a copy of
// the License at
//
// https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
// WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
// License for the specific language governing permissions and limitations under
// the License.
#include "pw_rpc2/internal/method.h"
#include <cstddef>
#include <cstring>
#include <type_traits>
#include <utility>
#include "pw_allocator/testing.h"
#include "pw_assert/check.h"
#include "pw_async2/coro.h"
#include "pw_async2/dispatcher_for_test.h"
#include "pw_async2/poll.h"
#include "pw_bytes/span.h"
#include "pw_rpc2/internal/method_future.h"
#include "pw_rpc2/internal/method_invoker.h"
#include "pw_rpc2/internal/server_call.h"
#include "pw_rpc2/internal/server_connection_task.h"
#include "pw_rpc2/internal/test_utils.h"
#include "pw_rpc2/method_type.h"
#include "pw_rpc2/server.h"
#include "pw_rpc2/service.h"
#include "pw_transport/socket.h"
#include "pw_unit_test/framework.h"
namespace pw::rpc2::internal {
namespace {
class MockFuture {
public:
using value_type = void;
explicit MockFuture(int pends_before_ready = 0)
: pends_remaining_(pends_before_ready) {}
bool is_pendable() const { return pends_remaining_ >= 0; }
bool is_complete() const { return pends_remaining_ == 0; }
async2::Poll<void> Pend(async2::Context& cx) {
if (pends_remaining_ <= 0) {
return async2::Ready();
}
--pends_remaining_;
if (pends_remaining_ == 0) {
return async2::Ready();
}
cx.ReEnqueue();
return async2::Pending();
}
private:
int pends_remaining_;
};
struct StubMsg {
int value = 0;
struct Serializer {
static size_t MaxEncodedSize(const StubMsg&) { return sizeof(int); }
static pw::StatusWithSize Serialize(const StubMsg& msg,
pw::span<std::byte> dest) {
if (dest.size() < sizeof(int)) {
return pw::StatusWithSize::ResourceExhausted();
}
std::memcpy(dest.data(), &msg.value, sizeof(int));
return pw::StatusWithSize(pw::OkStatus(), sizeof(int));
}
template <typename T>
static pw::Result<StubMsg> Deserialize(pw::span<const std::byte> source) {
if (source.size() < sizeof(int)) {
return pw::Status::DataLoss();
}
StubMsg msg;
std::memcpy(&msg.value, source.data(), sizeof(int));
return msg;
}
};
};
class TestService : public Service {
public:
TestService() : Service(1, {}) {}
MockFuture RawUnaryFuture(pw::ConstBuf request, RawUnaryWriter responder) {
last_request_size_ = request.size();
auto res_fut = responder.ReserveFinish(request.size());
return MockFuture(1);
}
MockFuture RawServerStreamingFuture(pw::ConstBuf request, RawWriter writer) {
last_request_size_ = request.size();
(void)writer;
return MockFuture(1);
}
MockFuture RawClientStreamingFuture(RawReader reader,
RawUnaryWriter responder) {
(void)reader;
(void)responder;
return MockFuture(1);
}
MockFuture RawBidiStreamingFuture(RawReader reader, RawWriter writer) {
(void)reader;
(void)writer;
return MockFuture(1);
}
async2::Coro<void> RawUnaryCoro(async2::CoroContext,
pw::ConstBuf request,
RawUnaryWriter responder) {
last_request_size_ = request.size();
auto res_fut = responder.ReserveFinish(request.size());
co_return;
}
async2::Coro<void> RawServerStreamingCoro(async2::CoroContext,
pw::ConstBuf request,
RawWriter writer) {
last_request_size_ = request.size();
(void)writer;
co_return;
}
async2::Coro<void> RawClientStreamingCoro(async2::CoroContext,
RawReader reader,
RawUnaryWriter responder) {
(void)reader;
(void)responder;
co_return;
}
async2::Coro<void> RawBidiStreamingCoro(async2::CoroContext,
RawReader reader,
RawWriter writer) {
(void)reader;
(void)writer;
co_return;
}
MockFuture TypedUnaryFuture(const StubMsg& req,
UnaryWriter<StubMsg> responder) {
last_typed_value_ = req.value;
(void)responder;
return MockFuture(1);
}
async2::Coro<void> TypedUnaryCoro(async2::CoroContext,
StubMsg req,
UnaryWriter<StubMsg> responder) {
last_typed_value_ = req.value;
(void)responder;
co_return;
}
void OnRawUnary(size_t size) { last_request_size_ = size; }
// A const method taking its request by value. Exercises the const and
// by-value axes of `MethodTraits` end to end.
MockFuture TypedUnaryFutureConst(StubMsg req,
UnaryWriter<StubMsg> responder) const {
last_const_value_ = req.value;
(void)responder;
return MockFuture(1);
}
size_t last_request_size() const { return last_request_size_; }
int last_typed_value() const { return last_typed_value_; }
int last_const_value() const { return last_const_value_; }
private:
size_t last_request_size_ = 0;
int last_typed_value_ = 0;
mutable int last_const_value_ = 0;
};
// Signature classification (`MethodTraits`) is checked in
// method_traits_test.cc, which needs only `method_traits.h`.
class StatelessUnaryFuture {
public:
using value_type = void;
StatelessUnaryFuture() = default;
StatelessUnaryFuture(pw::ConstBuf request, RawUnaryWriter responder)
: request_len_(request.size()) {
(void)responder;
}
bool is_pendable() const { return !done_; }
bool is_complete() const { return done_; }
async2::Poll<void> Pend(async2::Context&) {
done_ = true;
return async2::Ready();
}
size_t request_len() const { return request_len_; }
private:
size_t request_len_ = 0;
bool done_ = false;
};
class StatefulUnaryFuture {
public:
using value_type = void;
StatefulUnaryFuture() = default;
StatefulUnaryFuture(TestService& svc,
pw::ConstBuf request,
RawUnaryWriter responder) {
(void)responder;
svc.OnRawUnary(request.size());
}
bool is_pendable() const { return !done_; }
bool is_complete() const { return done_; }
async2::Poll<void> Pend(async2::Context&) {
done_ = true;
return async2::Ready();
}
private:
bool done_ = false;
};
class MethodInvokerTest : public ::testing::Test {
protected:
MethodInvokerTest() {
auto [conn, raw_conn] = test::MakeMockConnection(conn_alloc_);
connection_ = conn;
raw_conn_ = raw_conn;
connection_task_ = conn_alloc_.MakeShared<ServerConnectionTask>(
EstablishedConnection{connection_}, conn_alloc_, server_);
}
~MethodInvokerTest() override {
// The dispatcher is declared before the objects the connection task
// refers to, so it can no longer unpost the task on the way out. Do it
// here instead, while the server and the allocators are still alive.
if (connection_task_ != nullptr) {
connection_task_->Deregister();
}
}
/// Allocates a server call and registers it with the connection, exactly as
/// `ServerConnectionTask::HandleIncomingRequest()` does.
ServerCall& AdoptCall(uint32_t call_id, const Method& method) {
auto call_res = ServerCall::Allocate(
*connection_task_, call_id, service_, method, alloc_);
PW_CHECK_OK(call_res.status());
return **call_res;
}
/// Invokes `method` on `call` and starts the task that runs it, exactly
/// as `ServerConnectionTask::HandleIncomingRequest()` does.
ServerError InvokeIntoCall(const Method& method,
ServerCall& call,
pw::ConstBuf&& payload) {
const ServerError error = method.Invoke(service_, call, std::move(payload));
if (error == ServerError::kOk) {
dispatcher_.Post(call);
}
return error;
}
/// Runs the dispatcher, which polls the connection and every call started on
/// it.
///
/// A connection never completes on its own --- it is always waiting for the
/// next packet --- so this runs until the dispatcher stalls rather than to
/// completion.
void RunConnection() {
if (!posted_) {
dispatcher_.PostShared(connection_task_);
posted_ = true;
}
dispatcher_.RunUntilStalled();
}
pw::allocator::test::AllocatorForTest<4096> conn_alloc_;
pw::allocator::test::AllocatorForTest<4096> alloc_;
// Declared before the server, which is bound to it for the server's life.
async2::DispatcherForTest dispatcher_;
Server server_{conn_alloc_, dispatcher_};
transport::ReliableDatagramSocket connection_;
test::MockConnection* raw_conn_ = nullptr;
SharedPtr<ServerConnectionTask> connection_task_;
TestService service_;
bool posted_ = false;
};
TEST_F(MethodInvokerTest, CreateMethodForFutures) {
using RawUnary =
RawMethodInvoker<&TestService::RawUnaryFuture, MethodType::kUnary>;
constexpr Method unary = RawUnary::CreateMethod<TestService>(1);
EXPECT_EQ(unary.id(), 1u);
EXPECT_EQ(unary.future_storage_size(),
sizeof(async2::internal::BoxedFutureImpl<void, MockFuture>));
using RawServerStreaming =
RawMethodInvoker<&TestService::RawServerStreamingFuture,
MethodType::kServerStreaming>;
constexpr Method s_stream = RawServerStreaming::CreateMethod<TestService>(2);
EXPECT_EQ(s_stream.id(), 2u);
EXPECT_EQ(s_stream.future_storage_size(),
sizeof(async2::internal::BoxedFutureImpl<void, MockFuture>));
using RawClientStreaming =
RawMethodInvoker<&TestService::RawClientStreamingFuture,
MethodType::kClientStreaming>;
constexpr Method c_stream = RawClientStreaming::CreateMethod<TestService>(3);
EXPECT_EQ(c_stream.id(), 3u);
EXPECT_EQ(c_stream.future_storage_size(),
sizeof(async2::internal::BoxedFutureImpl<void, MockFuture>));
using RawBidi = RawMethodInvoker<&TestService::RawBidiStreamingFuture,
MethodType::kBidirectionalStreaming>;
constexpr Method bidi = RawBidi::CreateMethod<TestService>(4);
EXPECT_EQ(bidi.id(), 4u);
EXPECT_EQ(bidi.future_storage_size(),
sizeof(async2::internal::BoxedFutureImpl<void, MockFuture>));
}
TEST_F(MethodInvokerTest, CreateMethodForCoros) {
using RawUnaryCoro =
RawMethodInvoker<&TestService::RawUnaryCoro, MethodType::kUnary>;
constexpr Method unary = RawUnaryCoro::CreateMethod<TestService>(1);
EXPECT_EQ(unary.id(), 1u);
EXPECT_EQ(unary.future_storage_size(),
sizeof(MethodFutureImpl<async2::Coro<void>>));
using RawBidiCoro = RawMethodInvoker<&TestService::RawBidiStreamingCoro,
MethodType::kBidirectionalStreaming>;
constexpr Method bidi = RawBidiCoro::CreateMethod<TestService>(2);
EXPECT_EQ(bidi.id(), 2u);
EXPECT_EQ(bidi.future_storage_size(),
sizeof(MethodFutureImpl<async2::Coro<void>>));
}
TEST_F(MethodInvokerTest, InvokeRawUnaryFutureIntoCall) {
using Invoker =
RawMethodInvoker<&TestService::RawUnaryFuture, MethodType::kUnary>;
constexpr Method method = Invoker::CreateMethod<TestService>(1);
ServerCall& call = AdoptCall(/*call_id=*/1, method);
auto buf = pw::Buf::Allocate(conn_alloc_, 8);
pw::ConstBuf payload(std::move(buf));
ServerError error = InvokeIntoCall(method, call, std::move(payload));
EXPECT_EQ(error, ServerError::kOk);
EXPECT_EQ(service_.last_request_size(), 8u);
RunConnection();
EXPECT_EQ(alloc_.GetAllocated(), 0u);
}
TEST_F(MethodInvokerTest, InvokeRawUnaryCoroIntoCall) {
using Invoker =
RawMethodInvoker<&TestService::RawUnaryCoro, MethodType::kUnary>;
constexpr Method method = Invoker::CreateMethod<TestService>(2);
ServerCall& call = AdoptCall(/*call_id=*/2, method);
auto buf = pw::Buf::Allocate(conn_alloc_, 16);
pw::ConstBuf payload(std::move(buf));
ServerError error = InvokeIntoCall(method, call, std::move(payload));
EXPECT_EQ(error, ServerError::kOk);
RunConnection();
EXPECT_EQ(service_.last_request_size(), 16u);
EXPECT_EQ(alloc_.GetAllocated(), 0u);
}
TEST_F(MethodInvokerTest, InvokeTypedUnaryFutureIntoCall) {
using Invoker = MethodInvoker<&TestService::TypedUnaryFuture,
MethodType::kUnary,
StubMsg,
StubMsg>;
constexpr Method method = Invoker::CreateMethod<TestService>(3);
ServerCall& call = AdoptCall(/*call_id=*/3, method);
int val = 12345;
auto buf = pw::Buf::Allocate(conn_alloc_, sizeof(int));
std::memcpy(buf.data(), &val, sizeof(int));
pw::ConstBuf payload(std::move(buf));
ServerError error = InvokeIntoCall(method, call, std::move(payload));
EXPECT_EQ(error, ServerError::kOk);
EXPECT_EQ(service_.last_typed_value(), 12345);
RunConnection();
EXPECT_EQ(alloc_.GetAllocated(), 0u);
}
TEST_F(MethodInvokerTest, InvokeConstMethodWithRequestByValue) {
using Invoker = MethodInvoker<&TestService::TypedUnaryFutureConst,
MethodType::kUnary,
StubMsg,
StubMsg>;
constexpr Method method = Invoker::CreateMethod<TestService>(14);
ServerCall& call = AdoptCall(/*call_id=*/14, method);
int val = 4242;
auto buf = pw::Buf::Allocate(conn_alloc_, sizeof(int));
std::memcpy(buf.data(), &val, sizeof(int));
pw::ConstBuf payload(std::move(buf));
ServerError error = InvokeIntoCall(method, call, std::move(payload));
EXPECT_EQ(error, ServerError::kOk);
EXPECT_EQ(service_.last_const_value(), 4242);
RunConnection();
EXPECT_EQ(alloc_.GetAllocated(), 0u);
}
TEST_F(MethodInvokerTest, InvokeTypedUnaryFutureDeserializationFailure) {
using Invoker = MethodInvoker<&TestService::TypedUnaryFuture,
MethodType::kUnary,
StubMsg,
StubMsg>;
constexpr Method method = Invoker::CreateMethod<TestService>(4);
ServerCall& call = AdoptCall(/*call_id=*/4, method);
// Provide only 2 bytes when sizeof(int) = 4 is required.
auto buf = pw::Buf::Allocate(conn_alloc_, 2);
pw::ConstBuf payload(std::move(buf));
ServerError error = InvokeIntoCall(method, call, std::move(payload));
EXPECT_EQ(error, ServerError::kInvalidRequestPayload);
// The server retires a call whose invocation failed, which frees it: no
// future was ever set, so there is nothing else holding it.
connection_task_->RetireServerCall(call);
EXPECT_EQ(alloc_.GetAllocated(), 0u);
}
TEST_F(MethodInvokerTest, InvokeTypedUnaryCoroIntoCall) {
using Invoker = MethodInvoker<&TestService::TypedUnaryCoro,
MethodType::kUnary,
StubMsg,
StubMsg>;
constexpr Method method = Invoker::CreateMethod<TestService>(5);
ServerCall& call = AdoptCall(/*call_id=*/5, method);
int val = 9999;
auto buf = pw::Buf::Allocate(conn_alloc_, sizeof(int));
std::memcpy(buf.data(), &val, sizeof(int));
pw::ConstBuf payload(std::move(buf));
ServerError error = InvokeIntoCall(method, call, std::move(payload));
EXPECT_EQ(error, ServerError::kOk);
RunConnection();
EXPECT_EQ(service_.last_typed_value(), 9999);
EXPECT_EQ(alloc_.GetAllocated(), 0u);
}
TEST_F(MethodInvokerTest, InvokeRawServerStreamingFutureIntoCall) {
using Invoker = RawMethodInvoker<&TestService::RawServerStreamingFuture,
MethodType::kServerStreaming>;
constexpr Method method = Invoker::CreateMethod<TestService>(6);
ServerCall& call = AdoptCall(/*call_id=*/6, method);
auto buf = pw::Buf::Allocate(conn_alloc_, 4);
pw::ConstBuf payload(std::move(buf));
ServerError error = InvokeIntoCall(method, call, std::move(payload));
EXPECT_EQ(error, ServerError::kOk);
RunConnection();
EXPECT_EQ(alloc_.GetAllocated(), 0u);
}
TEST_F(MethodInvokerTest, InvokeRawClientStreamingFutureIntoCall) {
using Invoker = RawMethodInvoker<&TestService::RawClientStreamingFuture,
MethodType::kClientStreaming>;
constexpr Method method = Invoker::CreateMethod<TestService>(7);
ServerCall& call = AdoptCall(/*call_id=*/7, method);
pw::ConstBuf payload;
ServerError error = InvokeIntoCall(method, call, std::move(payload));
EXPECT_EQ(error, ServerError::kOk);
RunConnection();
EXPECT_EQ(alloc_.GetAllocated(), 0u);
}
TEST_F(MethodInvokerTest, InvokeRawBidiStreamingFutureIntoCall) {
using Invoker = RawMethodInvoker<&TestService::RawBidiStreamingFuture,
MethodType::kBidirectionalStreaming>;
constexpr Method method = Invoker::CreateMethod<TestService>(8);
ServerCall& call = AdoptCall(/*call_id=*/8, method);
pw::ConstBuf payload;
ServerError error = InvokeIntoCall(method, call, std::move(payload));
EXPECT_EQ(error, ServerError::kOk);
RunConnection();
EXPECT_EQ(alloc_.GetAllocated(), 0u);
}
TEST_F(MethodInvokerTest, InvokeRawServerStreamingCoroIntoCall) {
using Invoker = RawMethodInvoker<&TestService::RawServerStreamingCoro,
MethodType::kServerStreaming>;
constexpr Method method = Invoker::CreateMethod<TestService>(9);
ServerCall& call = AdoptCall(/*call_id=*/9, method);
auto buf = pw::Buf::Allocate(conn_alloc_, 4);
pw::ConstBuf payload(std::move(buf));
ServerError error = InvokeIntoCall(method, call, std::move(payload));
EXPECT_EQ(error, ServerError::kOk);
RunConnection();
EXPECT_EQ(alloc_.GetAllocated(), 0u);
}
TEST_F(MethodInvokerTest, InvokeRawClientStreamingCoroIntoCall) {
using Invoker = RawMethodInvoker<&TestService::RawClientStreamingCoro,
MethodType::kClientStreaming>;
constexpr Method method = Invoker::CreateMethod<TestService>(10);
ServerCall& call = AdoptCall(/*call_id=*/10, method);
pw::ConstBuf payload;
ServerError error = InvokeIntoCall(method, call, std::move(payload));
EXPECT_EQ(error, ServerError::kOk);
RunConnection();
EXPECT_EQ(alloc_.GetAllocated(), 0u);
}
TEST_F(MethodInvokerTest, InvokeRawBidiStreamingCoroIntoCall) {
using Invoker = RawMethodInvoker<&TestService::RawBidiStreamingCoro,
MethodType::kBidirectionalStreaming>;
constexpr Method method = Invoker::CreateMethod<TestService>(11);
ServerCall& call = AdoptCall(/*call_id=*/11, method);
pw::ConstBuf payload;
ServerError error = InvokeIntoCall(method, call, std::move(payload));
EXPECT_EQ(error, ServerError::kOk);
RunConnection();
EXPECT_EQ(alloc_.GetAllocated(), 0u);
}
TEST_F(MethodInvokerTest, FutureMethodInvokerCreateMethod) {
using Invoker =
RawFutureMethodInvoker<StatelessUnaryFuture, MethodType::kUnary>;
constexpr Method method = Invoker::CreateMethod<TestService>(1);
EXPECT_EQ(method.id(), 1u);
EXPECT_EQ(
method.future_storage_size(),
sizeof(async2::internal::BoxedFutureImpl<void, StatelessUnaryFuture>));
}
TEST_F(MethodInvokerTest, InvokeStatelessTypeMethodIntoCall) {
using Invoker =
RawFutureMethodInvoker<StatelessUnaryFuture, MethodType::kUnary>;
constexpr Method method = Invoker::CreateMethod<TestService>(12);
ServerCall& call = AdoptCall(/*call_id=*/12, method);
auto buf = pw::Buf::Allocate(conn_alloc_, 4);
pw::ConstBuf payload(std::move(buf));
ServerError error = InvokeIntoCall(method, call, std::move(payload));
EXPECT_EQ(error, ServerError::kOk);
RunConnection();
EXPECT_EQ(alloc_.GetAllocated(), 0u);
}
TEST_F(MethodInvokerTest, InvokeStatefulTypeMethodIntoCall) {
using Invoker =
RawFutureMethodInvoker<StatefulUnaryFuture, MethodType::kUnary>;
constexpr Method method = Invoker::CreateMethod<TestService>(13);
ServerCall& call = AdoptCall(/*call_id=*/13, method);
auto buf = pw::Buf::Allocate(conn_alloc_, 6);
pw::ConstBuf payload(std::move(buf));
ServerError error = InvokeIntoCall(method, call, std::move(payload));
EXPECT_EQ(error, ServerError::kOk);
EXPECT_EQ(service_.last_request_size(), 6u);
RunConnection();
EXPECT_EQ(alloc_.GetAllocated(), 0u);
}
TEST(MethodTest, StorageSizeBounds) {
constexpr auto stub_invoke =
[](Service&, ServerCall&, ConstBuf&&) -> ServerError {
return ServerError::kOk;
};
constexpr Method m1(101, MethodType::kUnary, 0, stub_invoke);
static_assert(m1.id() == 101);
static_assert(m1.future_storage_size() == 0);
EXPECT_EQ(m1.id(), 101u);
EXPECT_EQ(m1.future_storage_size(), 0u);
constexpr Method m2(102, MethodType::kServerStreaming, 128, stub_invoke);
static_assert(m2.id() == 102);
static_assert(m2.future_storage_size() == 128);
EXPECT_EQ(m2.id(), 102u);
EXPECT_EQ(m2.future_storage_size(), 128u);
// The maximum supported size.
constexpr Method m3(103,
MethodType::kBidirectionalStreaming,
Method::kMaxSizeBytes,
stub_invoke);
static_assert(m3.id() == 103);
static_assert(m3.future_storage_size() == Method::kMaxSizeBytes);
EXPECT_EQ(m3.id(), 103u);
EXPECT_EQ(m3.future_storage_size(), Method::kMaxSizeBytes);
}
TEST(MethodTest, StoresMethodType) {
constexpr auto stub_invoke =
[](Service&, ServerCall&, ConstBuf&&) -> ServerError {
return ServerError::kOk;
};
constexpr Method unary(1, MethodType::kUnary, 0, stub_invoke);
constexpr Method server_stream(
2, MethodType::kServerStreaming, 0, stub_invoke);
constexpr Method client_stream(
3, MethodType::kClientStreaming, 0, stub_invoke);
constexpr Method bidi(4, MethodType::kBidirectionalStreaming, 0, stub_invoke);
static_assert(unary.type() == MethodType::kUnary);
static_assert(server_stream.type() == MethodType::kServerStreaming);
static_assert(client_stream.type() == MethodType::kClientStreaming);
static_assert(bidi.type() == MethodType::kBidirectionalStreaming);
static_assert(!HasClientStream(unary.type()));
static_assert(!HasServerStream(unary.type()));
static_assert(!HasClientStream(server_stream.type()));
static_assert(HasServerStream(server_stream.type()));
static_assert(HasClientStream(client_stream.type()));
static_assert(!HasServerStream(client_stream.type()));
static_assert(HasClientStream(bidi.type()));
static_assert(HasServerStream(bidi.type()));
}
TEST_F(MethodInvokerTest,
CoroutineFrameAllocationFailureReturnsResourceExhaustedWithoutCancel) {
using UnaryCoroInvoker =
RawMethodInvoker<&TestService::RawUnaryCoro, MethodType::kUnary>;
constexpr Method unary_method =
UnaryCoroInvoker::CreateMethod<TestService>(20);
ServerCall& unary_call = AdoptCall(/*call_id=*/20, unary_method);
auto buf = pw::Buf::Allocate(conn_alloc_, 4);
pw::ConstBuf payload(std::move(buf));
// Exhaust the connection allocator so coroutine frame allocation fails.
conn_alloc_.Exhaust();
ServerError error =
InvokeIntoCall(unary_method, unary_call, std::move(payload));
EXPECT_EQ(error, ServerError::kFailedToAllocateCallResourcesWhileRunning);
// The responder destructor must NOT have closed the write side or queued a
// CANCELLED packet, leaving the call open for ServerConnectionTask to send
// RESOURCE_EXHAUSTED.
EXPECT_FALSE(unary_call.is_write_closed());
EXPECT_EQ(raw_conn_->written_packet_count(), 0u);
// ServerConnectionTask queues the error and closes the write side before
// retiring a call whose invocation failed. Closing it here keeps retirement
// from ending the call itself.
unary_call.CloseWrite();
connection_task_->RetireServerCall(unary_call);
using StreamCoroInvoker =
RawMethodInvoker<&TestService::RawServerStreamingCoro,
MethodType::kServerStreaming>;
constexpr Method stream_method =
StreamCoroInvoker::CreateMethod<TestService>(21);
ServerCall& stream_call = AdoptCall(/*call_id=*/21, stream_method);
pw::ConstBuf stream_payload;
ServerError stream_error =
InvokeIntoCall(stream_method, stream_call, std::move(stream_payload));
EXPECT_EQ(stream_error,
ServerError::kFailedToAllocateCallResourcesWhileRunning);
// The writer destructor must NOT have closed the write side or queued a
// StreamEnd (EOF) packet.
EXPECT_FALSE(stream_call.is_write_closed());
EXPECT_EQ(raw_conn_->written_packet_count(), 0u);
// ServerConnectionTask queues the error and closes the write side before
// retiring a call whose invocation failed. Closing it here keeps retirement
// from ending the call itself.
stream_call.CloseWrite();
connection_task_->RetireServerCall(stream_call);
}
} // namespace
} // namespace pw::rpc2::internal