blob: 92b44b20ce76d112f7b41f4c72bf5a7c57c9d48e [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/packet_reservation.h"
#include <cstddef>
#include "pw_assert/check.h"
#include "pw_async2/await.h"
#include "pw_bytes/span.h"
#include "pw_rpc2/writer.h"
#include "pw_status/status.h"
// These members are defined out of line so that their preconditions can be
// checked with `PW_CHECK`/`PW_DCHECK`, which carry a message but may not be
// used in headers.
namespace pw::rpc2 {
namespace internal {
async2::Poll<Status> WriteFutureBase::PendWrite(async2::Context& cx,
const void* payload,
SerializeFn serialize) {
PW_CHECK(is_pendable());
PW_AWAIT(auto res_result, res_fut_, cx);
if (!res_result.ok()) {
return async2::Ready(res_result.status());
}
PacketReservation res = std::move(*res_result);
if (serialize != nullptr) {
StatusWithSize serialize_status = serialize(payload, res.payload());
if (!serialize_status.ok()) {
res.Drop();
return async2::Ready(serialize_status.status());
}
res.TruncatePayload(serialize_status.size());
} else {
// `serialize` is null only for header-only packets (`WriteFuture<void>`),
// which are reserved with a payload size of 0.
PW_DASSERT(res.size() == 0);
}
return async2::Ready(res.Commit());
}
} // namespace internal
Status PacketReservation::Commit() {
if (state_ != State::kActive) {
return Status::FailedPrecondition();
}
state_ = State::kFinished;
// The call ended while this reservation was held: its terminal packet was
// sent first, or it was cancelled, completed by the peer, or retired. The
// peer has forgotten the call, so this packet would only arrive as a stray.
if (call_ != nullptr && call_->is_write_closed()) {
reservation_.Cancel();
const Status status =
call_->is_completed() && !call_->completion_status().ok()
? call_->completion_status()
: Status::FailedPrecondition();
ReleaseCall(/*committed=*/false);
return status;
}
auto encode_res = packet_.EncodeHeader(reservation_, payload_size_);
// The header always fits: the reservation was sized to hold it, and
// `payload_size_` is bounded by `TruncatePayload()`.
PW_CHECK_OK(encode_res.status(),
"Failed to encode the header of an RPC packet that was reserved "
"with room for it");
if (!reservation_.Commit(*encode_res)) {
ReleaseCall(/*committed=*/false);
return Status::Unavailable();
}
ReleaseCall(/*committed=*/true);
return OkStatus();
}
void PacketReservation::ReleaseCall(bool committed) {
if (call_ == nullptr) {
return;
}
if (committed && packet_.type().is_start()) {
call_->MarkStarted();
}
if (packet_.closes_stream()) {
if (committed) {
call_->CommitTerminalWrite();
} else {
call_->AbandonTerminalWrite();
}
}
call_ = nullptr;
}
void PacketReservation::CheckPayloadAccess() const {
PW_DCHECK(state_ == State::kActive,
"Payload accessed on a packet reservation that was already "
"committed, dropped, or moved from");
PW_DCHECK(packet_.type().has_payload(),
"Control packets (StreamEnd, Error) do not carry payload");
}
pw::ConstByteSpan PacketReservation::PayloadSpan(size_t size) const {
CheckPayloadAccess();
if (size == 0) {
return {};
}
return pw::ConstByteSpan(reservation_)
.subspan(packet_.payload_offset(), size);
}
} // namespace pw::rpc2