Internal change

PiperOrigin-RevId: 369699298
This commit is contained in:
hai007
2021-04-21 11:37:55 -07:00
committed by Copybara-Service
parent e0293911c7
commit 70fe292868
+159 -77
View File
@@ -57,16 +57,24 @@ bool PayloadManager::SendPayloadLoop(
// Update the still-active recipients of this payload.
if (available_endpoint_ids.empty()) {
NEARBY_LOG(INFO, "No more available endpoints: payload_id=%" PRIX64,
pending_payload.GetInternalPayload()->GetId());
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);
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, "Payload canceled locally: payload_id=%" PRIX64,
pending_payload.GetInternalPayload()->GetId());
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);
HandleFinishedOutgoingPayload(
client, available_endpoint_ids, payload_header, next_chunk_offset,
proto::connections::PayloadStatus::LOCAL_CANCELLATION);
@@ -93,8 +101,9 @@ bool PayloadManager::SendPayloadLoop(
pending_payload.GetInternalPayload()->GetTotalSize() > 0 &&
pending_payload.GetInternalPayload()->GetTotalSize() <
next_chunk_offset) {
NEARBY_LOG(INFO, "Payload xfer failed: payload_id=%" PRIX64,
pending_payload.GetInternalPayload()->GetId());
NEARBY_LOG(INFO, "Payload xfer failed: payload_id=%x" PRIx64,
static_cast<std::int64_t>(
pending_payload.GetInternalPayload()->GetId()));
HandleFinishedOutgoingPayload(
client, available_endpoint_ids, payload_header, next_chunk_offset,
proto::connections::PayloadStatus::LOCAL_ERROR);
@@ -108,8 +117,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=%" PRIX64
"; ids={%s}",
"Payload xfer: endpoints failed: payload_id=%x" PRIx64
"; endpoint_ids={%s}",
static_cast<std::int64_t>(payload_header.id()),
ToString(failed_endpoint_ids).c_str());
HandleFinishedOutgoingPayload(
@@ -129,14 +138,21 @@ 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=%" PRIX64 "; size=%" PRId64,
pending_payload.GetInternalPayload()->GetId(), next_chunk_offset);
NEARBY_LOG(INFO,
"Payload xfer done: payload_id=%x" PRIx64 "; size=%" PRId64,
static_cast<std::int64_t>(
pending_payload.GetInternalPayload()->GetId()),
next_chunk_offset);
return false;
}
}
@@ -205,7 +221,8 @@ 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=%" PRIX64, payload_id);
NEARBY_LOG(INFO, "CreateOutgoingPayload: payload_id=%x" PRIx64,
static_cast<std::int64_t>(payload_id));
MutexLock lock(&mutex_);
pending_payloads_.StartTrackingPayload(
payload_id, absl::make_unique<PendingPayload>(std::move(internal_payload),
@@ -304,9 +321,9 @@ void PayloadManager::SendPayload(ClientProxy* client,
// already checked whether or not we can work with this Payload type.
if (!executor) {
NEARBY_LOG(INFO,
"PayloadManager::SendPayload: unsupported: id=%" PRIX64
", type=%d",
payload.GetId(), payload.GetType());
"PayloadManager failed to determine the right executor for "
"outgoing payload_id=%" PRIx64 ", payload_type=%d",
static_cast<std::int64_t>(payload.GetId()), payload.GetType());
return;
}
@@ -317,10 +334,17 @@ void PayloadManager::SendPayload(ClientProxy* client,
Payload::Type payload_type = payload.GetType();
Payload::Id payload_id =
CreateOutgoingPayload(std::move(payload), endpoint_ids);
executor->Execute([this, client, endpoint_ids, payload_id]() {
executor->Execute([this, client, endpoint_ids, payload_id, payload_type]() {
if (shutdown_.Get()) return;
PendingPayload* pending_payload = GetPayload(payload_id);
if (!pending_payload) return;
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;
}
auto* internal_payload = pending_payload->GetInternalPayload();
if (!internal_payload) return;
PayloadTransferFrame::PayloadHeader payload_header{
@@ -337,8 +361,9 @@ void PayloadManager::SendPayload(ClientProxy* client,
});
});
NEARBY_LOG(INFO,
"PayloadManager: xfer scheduled: self=%p; id=%" PRIX64 ", type=%d",
this, payload_id, payload_type);
"PayloadManager: xfer scheduled: self=%p; payload_id=%x" PRIx64
", payload_type=%d",
this, static_cast<std::int64_t>(payload_id), payload_type);
}
PayloadManager::PendingPayload* PayloadManager::GetPayload(
@@ -351,14 +376,19 @@ Status PayloadManager::CancelPayload(ClientProxy* client,
Payload::Id payload_id) {
PendingPayload* canceled_payload = GetPayload(payload_id);
if (!canceled_payload) {
NEARBY_LOG(INFO, "PayloadManager: not found; payload_id=%" PRIX64,
payload_id);
NEARBY_LOG(INFO,
"Client requested cancel for unknown payload_id=%x" PRIx64
", ignoring.",
static_cast<std::int64_t>(payload_id));
return {Status::kPayloadUnknown};
}
// Mark the payload as canceled.
canceled_payload->MarkLocallyCanceled();
NEARBY_LOG(INFO, "PayloadManager: canceled; id=%" PRIX64, payload_id);
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));
// Return SUCCESS immediately. Remaining cleanup and updates will be sent in
// SendPayload() or OnIncomingFrame()
@@ -374,19 +404,20 @@ void PayloadManager::OnIncomingFrame(
switch (frame.packet_type()) {
case PayloadTransferFrame::CONTROL:
NEARBY_LOG(INFO,
"PayloadManager::OnIncomingFrame [CONTROL]: self=%p; id=%s",
this, from_endpoint_id.c_str());
NEARBY_LOG(
INFO,
"PayloadManager::OnIncomingFrame [CONTROL]: self=%p; endpoint_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; id=%s",
this, from_endpoint_id.c_str());
NEARBY_LOG(WARNING,
"PayloadManager: invalid frame; remote endpoint: self=%p; "
"endpoint_id=%s",
this, from_endpoint_id.c_str());
break;
}
}
@@ -446,7 +477,10 @@ PayloadManager::EndpointInfoStatusToPayloadStatus(EndpointInfo::Status status) {
case EndpointInfo::Status::kAvailable:
return proto::connections::PayloadStatus::SUCCESS;
default:
NEARBY_LOG(INFO, "PayloadManager: unknown status=%d", status);
NEARBY_LOG(
INFO,
"PayloadManager: Unknown PayloadStatus for EndpointInfo.Status=%d",
status);
return proto::connections::PayloadStatus::UNKNOWN_PAYLOAD_STATUS;
}
}
@@ -536,7 +570,8 @@ PayloadManager::PendingPayload* PayloadManager::CreateIncomingPayload(
}
Payload::Id payload_id = internal_payload->GetId();
NEARBY_LOG(INFO, "CreateIncomingPayload: payload_id=%" PRIX64, payload_id);
NEARBY_LOG(INFO, "CreateIncomingPayload: payload_id=%x" PRIx64,
static_cast<std::int64_t>(payload_id));
MutexLock lock(&mutex_);
pending_payloads_.StartTrackingPayload(
payload_id,
@@ -589,24 +624,24 @@ void PayloadManager::SendClientCallbacksForFinishedIncomingPayload(
ClientProxy* client, const std::string& endpoint_id,
const PayloadTransferFrame::PayloadHeader& payload_header,
std::int64_t offset_bytes, proto::connections::PayloadStatus status) {
RunOnStatusUpdateThread(
[this, client, endpoint_id, payload_header, offset_bytes,
status]() RUN_ON_PAYLOAD_STATUS_UPDATE_THREAD() {
// Make sure we're still tracking this payload.
PendingPayload* pending_payload = GetPayload(payload_header.id());
if (!pending_payload) {
return;
}
RunOnStatusUpdateThread([this, client, endpoint_id, payload_header,
offset_bytes,
status]() RUN_ON_PAYLOAD_STATUS_UPDATE_THREAD() {
// Make sure we're still tracking this payload.
PendingPayload* pending_payload = GetPayload(payload_header.id());
if (!pending_payload) {
return;
}
// Unless we never started tracking this payload (meaning we failed to
// even create the InternalPayload), notify the client (and close it).
PayloadProgressInfo update{
payload_header.id(),
PayloadManager::PayloadStatusToTransferUpdateStatus(status),
payload_header.total_size(), offset_bytes};
NotifyClientOfIncomingPayloadProgressInfo(client, endpoint_id, update);
DestroyPendingPayload(payload_header.id());
});
// Unless we never started tracking this payload (meaning we failed to
// even create the InternalPayload), notify the client (and close it).
PayloadProgressInfo update{
payload_header.id(),
PayloadManager::PayloadStatusToTransferUpdateStatus(status),
payload_header.total_size(), offset_bytes};
NotifyClientOfIncomingPayloadProgressInfo(client, endpoint_id, update);
DestroyPendingPayload(payload_header.id());
});
}
void PayloadManager::SendControlMessage(
@@ -639,9 +674,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=%" PRIX64,
static_cast<std::int64_t>(payload_header.id()));
NEARBY_LOG(
INFO, "Sending PAYLOAD_CANCEL to receiver side; payload_id=%x" PRIx64,
static_cast<std::int64_t>(payload_header.id()));
SendControlMessage(
finished_endpoint_ids, payload_header,
num_bytes_successfully_transferred,
@@ -659,7 +694,10 @@ void PayloadManager::HandleFinishedOutgoingPayload(
// No special handling needed for these.
break;
default:
NEARBY_LOG(INFO, "PayloadManager: unknown status=%d", status);
NEARBY_LOG(INFO,
"PayloadManager: Unhandled finished outgoing payload with "
"payload_status=%d",
status);
break;
}
}
@@ -682,7 +720,10 @@ void PayloadManager::HandleFinishedIncomingPayload(
PayloadTransferFrame::ControlMessage::PAYLOAD_CANCELED);
break;
default:
// TODO(tracyzhou): Add logging.
NEARBY_LOG(INFO,
"Unhandled finished incoming payload_id=%x" PRIx64
" with payload_status=%d!",
static_cast<std::int64_t>(payload_header.id()), status);
break;
}
}
@@ -701,7 +742,8 @@ 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: id=%s",
"HandleSuccessfulOutgoingChunk: endpoint not found: "
"endpoint_id=%s",
endpoint_id.c_str());
return;
}
@@ -743,8 +785,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; id=%" PRIX64,
direction, this, payload_id);
"self=%p; payload_id=%x" PRIx64,
direction, this, static_cast<std::int64_t>(payload_id));
pending->Close();
pending.reset();
}
@@ -790,12 +832,23 @@ 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(),
@@ -808,8 +861,11 @@ void PayloadManager::ProcessDataPacket(
[to_client, from_endpoint_id, pending_payload]()
RUN_ON_PAYLOAD_STATUS_UPDATE_THREAD() {
NEARBY_LOG(INFO,
"ProcessDataPacket [new]: id=%s; payload_id=%" PRIX64,
from_endpoint_id.c_str(), pending_payload->GetId());
"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());
to_client->OnPayload(
from_endpoint_id,
pending_payload->GetInternalPayload()->ReleasePayload());
@@ -817,10 +873,11 @@ void PayloadManager::ProcessDataPacket(
} else {
pending_payload = GetPayload(payload_header.id());
if (!pending_payload) {
NEARBY_LOG(WARNING,
"ProcessDataPacket: [missing] id=%s; payload_id=%" PRIX64,
from_endpoint_id.c_str(),
static_cast<std::int64_t>(payload_header.id()));
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()));
return;
}
}
@@ -828,8 +885,11 @@ 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] id=%s; payload_id=%" PRIX64,
from_endpoint_id.c_str(), pending_payload->GetId());
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()));
HandleFinishedIncomingPayload(
to_client, from_endpoint_id, payload_header, payload_chunk.offset(),
proto::connections::PayloadStatus::LOCAL_CANCELLATION);
@@ -849,9 +909,11 @@ void PayloadManager::ProcessDataPacket(
if (pending_payload->GetInternalPayload()
->AttachNextChunk(ByteArray(std::move(*payload_chunk.mutable_body())))
.Raised()) {
NEARBY_LOG(WARNING,
"ProcessDataPacket: [data: error] id=%s; payload_id=%" PRIX64,
from_endpoint_id.c_str(), pending_payload->GetId());
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()));
HandleFinishedIncomingPayload(
to_client, from_endpoint_id, payload_header, payload_chunk.offset(),
proto::connections::PayloadStatus::LOCAL_ERROR);
@@ -873,14 +935,19 @@ void PayloadManager::ProcessControlPacket(
payload_transfer_frame.control_message();
PendingPayload* pending_payload = GetPayload(payload_header.id());
if (!pending_payload) {
// TODO(tracyzhou): Add logging.
NEARBY_LOG(INFO,
"Got ControlMessage for unknown payload_id=%x" PRIx64
", ignoring: %d",
static_cast<std::int64_t>(payload_header.id()),
control_message.event());
return;
}
switch (control_message.event()) {
case PayloadTransferFrame::ControlMessage::PAYLOAD_CANCELED:
if (pending_payload->IsIncoming()) {
NEARBY_LOG(INFO, "Incoming PAYLOAD_CANCELED: from id=%s; self=%p",
NEARBY_LOG(INFO,
"Incoming PAYLOAD_CANCELED: from endpoint_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
@@ -890,13 +957,20 @@ void PayloadManager::ProcessControlPacket(
control_message.offset(),
ControlMessageEventToPayloadStatus(control_message.event()));
} else {
NEARBY_LOG(INFO, "Outgoing PAYLOAD_CANCELED: from id=%s; self=%p",
NEARBY_LOG(INFO,
"Outgoing PAYLOAD_CANCELED: from endpoint_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);
}
// TODO(tracyzhou): Add logging.
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());
break;
case PayloadTransferFrame::ControlMessage::PAYLOAD_ERROR:
if (pending_payload->IsIncoming()) {
@@ -910,7 +984,10 @@ void PayloadManager::ProcessControlPacket(
}
break;
default:
// TODO(tracyzhou): Add logging.
NEARBY_LOG(INFO, "Unhandled control message %d for payload_id= %x" PRIx64,
control_message.event(),
static_cast<std::int64_t>(
pending_payload->GetInternalPayload()->GetId()));
break;
}
}
@@ -933,7 +1010,9 @@ PayloadManager::EndpointInfo::ControlMessageEventToEndpointInfoStatus(
case PayloadTransferFrame::ControlMessage::PAYLOAD_CANCELED:
return Status::kCanceled;
default:
// TODO(tracyzhou): Add logging.
NEARBY_LOG(INFO,
"Unknown EndpointInfo.Status for ControlMessage.EventType %d!",
event);
return Status::kUnknown;
}
}
@@ -941,6 +1020,9 @@ 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 ////////////////////////////////
@@ -957,7 +1039,7 @@ PayloadManager::PendingPayload::PendingPayload(
for (const auto& id : endpoint_ids) {
endpoints_.emplace(id, EndpointInfo{
.id = id,
.status {EndpointInfo::Status::kAvailable},
.status{EndpointInfo::Status::kAvailable},
});
}
}
@@ -1062,8 +1144,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=%" PRIX64 "; inserted=%d",
payload_id, pair.second);
NEARBY_LOG(INFO, "StartTrackingPayload: payload_id=%x" PRIx64 "; inserted=%d",
static_cast<std::int64_t>(payload_id), pair.second);
}
std::unique_ptr<PayloadManager::PendingPayload>