Internal change

PiperOrigin-RevId: 372029711
This commit is contained in:
hai007
2021-05-04 17:33:22 -07:00
committed by Copybara-Service
parent 6332010f2a
commit e7cbcabe0f
+59 -141
View File
@@ -57,24 +57,16 @@ bool PayloadManager::SendPayloadLoop(
// Update the still-active recipients of this payload.
if (available_endpoint_ids.empty()) {
NEARBY_LOG(INFO,
"PayloadManager short-circuiting payload_id=%" PRIX64
" after sending %" PRIX64
" bytes because none of the endpoints are available anymore.",
pending_payload.GetInternalPayload()->GetId(),
next_chunk_offset);
NEARBY_LOG(INFO, "No more available endpoints: payload_id=%" PRIX64,
pending_payload.GetInternalPayload()->GetId());
return false;
}
// Check if the payload has been cancelled by the client and, if so,
// notify the remaining recipients.
if (pending_payload.IsLocallyCanceled()) {
NEARBY_LOG(INFO,
"Aborting send of payload_id=%x" PRIx64 " at offset %x" PRIx64
" since it is marked canceled.",
static_cast<std::int64_t>(
pending_payload.GetInternalPayload()->GetId()),
next_chunk_offset);
NEARBY_LOG(INFO, "Payload canceled locally: payload_id=%" PRIX64,
pending_payload.GetInternalPayload()->GetId());
HandleFinishedOutgoingPayload(
client, available_endpoint_ids, payload_header, next_chunk_offset,
proto::connections::PayloadStatus::LOCAL_CANCELLATION);
@@ -101,9 +93,8 @@ bool PayloadManager::SendPayloadLoop(
pending_payload.GetInternalPayload()->GetTotalSize() > 0 &&
pending_payload.GetInternalPayload()->GetTotalSize() <
next_chunk_offset) {
NEARBY_LOG(INFO, "Payload xfer failed: payload_id=%x" PRIx64,
static_cast<std::int64_t>(
pending_payload.GetInternalPayload()->GetId()));
NEARBY_LOG(INFO, "Payload xfer failed: payload_id=%" PRIX64,
pending_payload.GetInternalPayload()->GetId());
HandleFinishedOutgoingPayload(
client, available_endpoint_ids, payload_header, next_chunk_offset,
proto::connections::PayloadStatus::LOCAL_ERROR);
@@ -117,8 +108,8 @@ bool PayloadManager::SendPayloadLoop(
// Check whether at least one endpoint failed.
if (!failed_endpoint_ids.empty()) {
NEARBY_LOG(INFO,
"Payload xfer: endpoints failed: payload_id=%x" PRIx64
"; endpoint_ids={%s}",
"Payload xfer: endpoints failed: payload_id=%" PRIX64
"; ids={%s}",
static_cast<std::int64_t>(payload_header.id()),
ToString(failed_endpoint_ids).c_str());
HandleFinishedOutgoingPayload(
@@ -138,21 +129,14 @@ bool PayloadManager::SendPayloadLoop(
payload_chunk.offset(), payload_chunk.body().size());
}
}
NEARBY_LOG(VERBOSE,
"PayloadManager done sending chunk at offset %x" PRIx64
" of payload_id=%x" PRIx64,
next_chunk_offset,
static_cast<std::int64_t>(
pending_payload.GetInternalPayload()->GetId()));
next_chunk_offset += next_chunk_size;
if (!next_chunk_size) {
// That was the last chunk, we're outta here.
NEARBY_LOG(INFO,
"Payload xfer done: payload_id=%x" PRIx64 "; size=%" PRId64,
static_cast<std::int64_t>(
pending_payload.GetInternalPayload()->GetId()),
next_chunk_offset);
NEARBY_LOG(
INFO, "Payload xfer done: payload_id=%" PRIX64 "; size=%" PRId64,
pending_payload.GetInternalPayload()->GetId(), next_chunk_offset);
return false;
}
}
@@ -221,8 +205,7 @@ Payload::Id PayloadManager::CreateOutgoingPayload(
Payload payload, const EndpointIds& endpoint_ids) {
auto internal_payload{CreateOutgoingInternalPayload(std::move(payload))};
Payload::Id payload_id = internal_payload->GetId();
NEARBY_LOG(INFO, "CreateOutgoingPayload: payload_id=%x" PRIx64,
static_cast<std::int64_t>(payload_id));
NEARBY_LOG(INFO, "CreateOutgoingPayload: payload_id=%" PRIX64, payload_id);
MutexLock lock(&mutex_);
pending_payloads_.StartTrackingPayload(
payload_id, absl::make_unique<PendingPayload>(std::move(internal_payload),
@@ -322,9 +305,9 @@ void PayloadManager::SendPayload(ClientProxy* client,
// already checked whether or not we can work with this Payload type.
if (!executor) {
NEARBY_LOG(INFO,
"PayloadManager failed to determine the right executor for "
"outgoing payload_id=%" PRIx64 ", payload_type=%d",
static_cast<std::int64_t>(payload.GetId()), payload.GetType());
"PayloadManager::SendPayload: unsupported: id=%" PRIX64
", type=%d",
payload.GetId(), payload.GetType());
return;
}
@@ -339,14 +322,7 @@ void PayloadManager::SendPayload(ClientProxy* client,
payload_type]() {
if (shutdown_.Get()) return;
PendingPayload* pending_payload = GetPayload(payload_id);
if (!pending_payload) {
NEARBY_LOG(INFO,
"PayloadManager failed to create InternalPayload for outgoing "
"payload_id=%x" PRIx64
", payload_type=%d, aborting sendPayload().",
static_cast<std::int64_t>(payload_id), payload_type);
return;
}
if (!pending_payload) return;
auto* internal_payload = pending_payload->GetInternalPayload();
if (!internal_payload) return;
PayloadTransferFrame::PayloadHeader payload_header{
@@ -364,9 +340,8 @@ void PayloadManager::SendPayload(ClientProxy* client,
});
});
NEARBY_LOG(INFO,
"PayloadManager: xfer scheduled: self=%p; payload_id=%x" PRIx64
", payload_type=%d",
this, static_cast<std::int64_t>(payload_id), payload_type);
"PayloadManager: xfer scheduled: self=%p; id=%" PRIX64 ", type=%d",
this, payload_id, payload_type);
}
PayloadManager::PendingPayload* PayloadManager::GetPayload(
@@ -379,19 +354,14 @@ Status PayloadManager::CancelPayload(ClientProxy* client,
Payload::Id payload_id) {
PendingPayload* canceled_payload = GetPayload(payload_id);
if (!canceled_payload) {
NEARBY_LOG(INFO,
"Client requested cancel for unknown payload_id=%x" PRIx64
", ignoring.",
static_cast<std::int64_t>(payload_id));
NEARBY_LOG(INFO, "PayloadManager: not found; payload_id=%" PRIX64,
payload_id);
return {Status::kPayloadUnknown};
}
// Mark the payload as canceled.
canceled_payload->MarkLocallyCanceled();
NEARBY_LOG(INFO,
"Cancelling %s payload_id=%x" PRIx64 " at request of client.",
(canceled_payload->IsIncoming() ? "incoming" : "outgoing"),
static_cast<std::int64_t>(payload_id));
NEARBY_LOG(INFO, "PayloadManager: canceled; id=%" PRIX64, payload_id);
// Return SUCCESS immediately. Remaining cleanup and updates will be sent in
// SendPayload() or OnIncomingFrame()
@@ -407,20 +377,19 @@ void PayloadManager::OnIncomingFrame(
switch (frame.packet_type()) {
case PayloadTransferFrame::CONTROL:
NEARBY_LOG(
INFO,
"PayloadManager::OnIncomingFrame [CONTROL]: self=%p; endpoint_id=%s",
this, from_endpoint_id.c_str());
NEARBY_LOG(INFO,
"PayloadManager::OnIncomingFrame [CONTROL]: self=%p; id=%s",
this, from_endpoint_id.c_str());
ProcessControlPacket(to_client, from_endpoint_id, frame);
break;
case PayloadTransferFrame::DATA:
ProcessDataPacket(to_client, from_endpoint_id, frame);
break;
default:
NEARBY_LOG(WARNING,
"PayloadManager: invalid frame; remote endpoint: self=%p; "
"endpoint_id=%s",
this, from_endpoint_id.c_str());
NEARBY_LOG(
WARNING,
"PayloadManager: invalid frame; remote endpoint: self=%p; id=%s",
this, from_endpoint_id.c_str());
break;
}
}
@@ -481,10 +450,7 @@ PayloadManager::EndpointInfoStatusToPayloadStatus(EndpointInfo::Status status) {
case EndpointInfo::Status::kAvailable:
return proto::connections::PayloadStatus::SUCCESS;
default:
NEARBY_LOG(
INFO,
"PayloadManager: Unknown PayloadStatus for EndpointInfo.Status=%d",
status);
NEARBY_LOG(INFO, "PayloadManager: unknown status=%d", status);
return proto::connections::PayloadStatus::UNKNOWN_PAYLOAD_STATUS;
}
}
@@ -574,8 +540,7 @@ PayloadManager::PendingPayload* PayloadManager::CreateIncomingPayload(
}
Payload::Id payload_id = internal_payload->GetId();
NEARBY_LOG(INFO, "CreateIncomingPayload: payload_id=%x" PRIx64,
static_cast<std::int64_t>(payload_id));
NEARBY_LOG(INFO, "CreateIncomingPayload: payload_id=%" PRIX64, payload_id);
MutexLock lock(&mutex_);
pending_payloads_.StartTrackingPayload(
payload_id,
@@ -681,9 +646,9 @@ void PayloadManager::HandleFinishedOutgoingPayload(
PayloadTransferFrame::ControlMessage::PAYLOAD_ERROR);
break;
case proto::connections::PayloadStatus::LOCAL_CANCELLATION:
NEARBY_LOG(
INFO, "Sending PAYLOAD_CANCEL to receiver side; payload_id=%x" PRIx64,
static_cast<std::int64_t>(payload_header.id()));
NEARBY_LOG(INFO,
"Sending PAYLOAD_CANCEL to receiver side; payload_id=%" PRIX64,
static_cast<std::int64_t>(payload_header.id()));
SendControlMessage(
finished_endpoint_ids, payload_header,
num_bytes_successfully_transferred,
@@ -701,10 +666,7 @@ void PayloadManager::HandleFinishedOutgoingPayload(
// No special handling needed for these.
break;
default:
NEARBY_LOG(INFO,
"PayloadManager: Unhandled finished outgoing payload with "
"payload_status=%d",
status);
NEARBY_LOG(INFO, "PayloadManager: unknown status=%d", status);
break;
}
}
@@ -727,10 +689,7 @@ void PayloadManager::HandleFinishedIncomingPayload(
PayloadTransferFrame::ControlMessage::PAYLOAD_CANCELED);
break;
default:
NEARBY_LOG(INFO,
"Unhandled finished incoming payload_id=%x" PRIx64
" with payload_status=%d!",
static_cast<std::int64_t>(payload_header.id()), status);
// TODO(tracyzhou): Add logging.
break;
}
}
@@ -750,8 +709,7 @@ void PayloadManager::HandleSuccessfulOutgoingChunk(
PendingPayload* pending_payload = GetPayload(payload_header.id());
if (!pending_payload || !pending_payload->GetEndpoint(endpoint_id)) {
NEARBY_LOG(INFO,
"HandleSuccessfulOutgoingChunk: endpoint not found: "
"endpoint_id=%s",
"HandleSuccessfulOutgoingChunk: endpoint not found: id=%s",
endpoint_id.c_str());
return;
}
@@ -793,8 +751,8 @@ void PayloadManager::DestroyPendingPayload(Payload::Id payload_id) {
const char* direction = is_incoming ? "incoming" : "outgoing";
NEARBY_LOG(INFO,
"PayloadManager: destroying %s pending payload: "
"self=%p; payload_id=%x" PRIx64,
direction, this, static_cast<std::int64_t>(payload_id));
"self=%p; id=%" PRIX64,
direction, this, payload_id);
pending->Close();
pending.reset();
}
@@ -841,23 +799,12 @@ void PayloadManager::ProcessDataPacket(
*payload_transfer_frame.mutable_payload_header();
PayloadTransferFrame::PayloadChunk& payload_chunk =
*payload_transfer_frame.mutable_payload_chunk();
NEARBY_LOG(VERBOSE,
"PayloadManager got data OfflineFrame for payload_id=%x" PRIx64
" from endpoint_id=%s at offset %x" PRIx64,
static_cast<std::int64_t>(payload_header.id()),
from_endpoint_id.c_str(), payload_chunk.offset());
PendingPayload* pending_payload;
if (payload_chunk.offset() == 0) {
pending_payload =
CreateIncomingPayload(payload_transfer_frame, from_endpoint_id);
if (!pending_payload) {
NEARBY_LOG(WARNING,
"PayloadManager failed to create InternalPayload from "
"PayloadTransferFrame with ID %x" PRIx64
" and type %d, aborting receipt.",
static_cast<std::int64_t>(payload_header.id()),
payload_header.type());
// Send the error to the remote endpoint.
SendControlMessage({from_endpoint_id}, payload_header,
payload_chunk.offset(),
@@ -871,11 +818,8 @@ void PayloadManager::ProcessDataPacket(
[to_client, from_endpoint_id, pending_payload]()
RUN_ON_PAYLOAD_STATUS_UPDATE_THREAD() {
NEARBY_LOG(INFO,
"PayloadManager received new payload_id=%x" PRIx64
" from endpoint_id=%s",
static_cast<std::int64_t>(
pending_payload->GetInternalPayload()->GetId()),
from_endpoint_id.c_str());
"ProcessDataPacket [new]: id=%s; payload_id=%" PRIX64,
from_endpoint_id.c_str(), pending_payload->GetId());
to_client->OnPayload(
from_endpoint_id,
pending_payload->GetInternalPayload()->ReleasePayload());
@@ -883,11 +827,10 @@ void PayloadManager::ProcessDataPacket(
} else {
pending_payload = GetPayload(payload_header.id());
if (!pending_payload) {
NEARBY_LOG(
WARNING,
"ProcessDataPacket: [missing] endpoint_id=%s; payload_id=%x" PRIx64,
from_endpoint_id.c_str(),
static_cast<std::int64_t>(payload_header.id()));
NEARBY_LOG(WARNING,
"ProcessDataPacket: [missing] id=%s; payload_id=%" PRIX64,
from_endpoint_id.c_str(),
static_cast<std::int64_t>(payload_header.id()));
return;
}
}
@@ -895,11 +838,8 @@ void PayloadManager::ProcessDataPacket(
if (pending_payload->IsLocallyCanceled()) {
// This incoming payload was canceled by the client. Drop this frame and do
// all the cleanup. See go/nc-cancel-payload
NEARBY_LOG(
INFO,
"ProcessDataPacket: [cancel] endpoint_id=%s; payload_id=%x" PRIx64,
from_endpoint_id.c_str(),
static_cast<std::int64_t>(pending_payload->GetId()));
NEARBY_LOG(INFO, "ProcessDataPacket: [cancel] id=%s; payload_id=%" PRIX64,
from_endpoint_id.c_str(), pending_payload->GetId());
HandleFinishedIncomingPayload(
to_client, from_endpoint_id, payload_header, payload_chunk.offset(),
proto::connections::PayloadStatus::LOCAL_CANCELLATION);
@@ -919,11 +859,9 @@ void PayloadManager::ProcessDataPacket(
if (pending_payload->GetInternalPayload()
->AttachNextChunk(ByteArray(std::move(*payload_chunk.mutable_body())))
.Raised()) {
NEARBY_LOG(
ERROR,
"ProcessDataPacket: [data: error] endpoint_id=%s; payload_id=%x" PRIx64,
from_endpoint_id.c_str(),
static_cast<std::int64_t>(pending_payload->GetId()));
NEARBY_LOG(WARNING,
"ProcessDataPacket: [data: error] id=%s; payload_id=%" PRIX64,
from_endpoint_id.c_str(), pending_payload->GetId());
HandleFinishedIncomingPayload(
to_client, from_endpoint_id, payload_header, payload_chunk.offset(),
proto::connections::PayloadStatus::LOCAL_ERROR);
@@ -945,19 +883,14 @@ void PayloadManager::ProcessControlPacket(
payload_transfer_frame.control_message();
PendingPayload* pending_payload = GetPayload(payload_header.id());
if (!pending_payload) {
NEARBY_LOG(INFO,
"Got ControlMessage for unknown payload_id=%x" PRIx64
", ignoring: %d",
static_cast<std::int64_t>(payload_header.id()),
control_message.event());
// TODO(tracyzhou): Add logging.
return;
}
switch (control_message.event()) {
case PayloadTransferFrame::ControlMessage::PAYLOAD_CANCELED:
if (pending_payload->IsIncoming()) {
NEARBY_LOG(INFO,
"Incoming PAYLOAD_CANCELED: from endpoint_id=%s; self=%p",
NEARBY_LOG(INFO, "Incoming PAYLOAD_CANCELED: from id=%s; self=%p",
from_endpoint_id.c_str(), this);
// No need to mark the pending payload as cancelled, since this is a
// remote cancellation for an incoming payload -- we handle everything
@@ -967,20 +900,13 @@ void PayloadManager::ProcessControlPacket(
control_message.offset(),
ControlMessageEventToPayloadStatus(control_message.event()));
} else {
NEARBY_LOG(INFO,
"Outgoing PAYLOAD_CANCELED: from endpoint_id=%s; self=%p",
NEARBY_LOG(INFO, "Outgoing PAYLOAD_CANCELED: from id=%s; self=%p",
from_endpoint_id.c_str(), this);
// Mark the payload as canceled *for this endpoint*.
pending_payload->SetEndpointStatusFromControlMessage(from_endpoint_id,
control_message);
}
NEARBY_LOG(VERBOSE,
"Marked %s payload_id=" PRIx64
" as canceled at request of endpoint_id=%s.",
(pending_payload->IsIncoming() ? "incoming" : "outgoing"),
static_cast<std::int64_t>(
pending_payload->GetInternalPayload()->GetId()),
from_endpoint_id.c_str());
// TODO(tracyzhou): Add logging.
break;
case PayloadTransferFrame::ControlMessage::PAYLOAD_ERROR:
if (pending_payload->IsIncoming()) {
@@ -994,10 +920,7 @@ void PayloadManager::ProcessControlPacket(
}
break;
default:
NEARBY_LOG(INFO, "Unhandled control message %d for payload_id= %x" PRIx64,
control_message.event(),
static_cast<std::int64_t>(
pending_payload->GetInternalPayload()->GetId()));
// TODO(tracyzhou): Add logging.
break;
}
}
@@ -1020,9 +943,7 @@ PayloadManager::EndpointInfo::ControlMessageEventToEndpointInfoStatus(
case PayloadTransferFrame::ControlMessage::PAYLOAD_CANCELED:
return Status::kCanceled;
default:
NEARBY_LOG(INFO,
"Unknown EndpointInfo.Status for ControlMessage.EventType %d!",
event);
// TODO(tracyzhou): Add logging.
return Status::kUnknown;
}
}
@@ -1030,9 +951,6 @@ PayloadManager::EndpointInfo::ControlMessageEventToEndpointInfoStatus(
void PayloadManager::EndpointInfo::SetStatusFromControlMessage(
const PayloadTransferFrame::ControlMessage& control_message) {
status.Set(ControlMessageEventToEndpointInfoStatus(control_message.event()));
NEARBY_LOG(VERBOSE,
"Marked endpoint %s with status %d based on OOB ControlMessage",
id.c_str(), status.Get());
}
//////////////////////////////// PendingPayload ////////////////////////////////
@@ -1049,7 +967,7 @@ PayloadManager::PendingPayload::PendingPayload(
for (const auto& id : endpoint_ids) {
endpoints_.emplace(id, EndpointInfo{
.id = id,
.status{EndpointInfo::Status::kAvailable},
.status {EndpointInfo::Status::kAvailable},
});
}
}
@@ -1155,8 +1073,8 @@ void PayloadManager::PendingPayloads::StartTrackingPayload(
pending_payloads_.erase(payload_id);
}
auto pair = pending_payloads_.emplace(payload_id, std::move(pending_payload));
NEARBY_LOG(INFO, "StartTrackingPayload: payload_id=%x" PRIx64 "; inserted=%d",
static_cast<std::int64_t>(payload_id), pair.second);
NEARBY_LOG(INFO, "StartTrackingPayload: payload_id=%" PRIX64 "; inserted=%d",
payload_id, pair.second);
}
std::unique_ptr<PayloadManager::PendingPayload>