[Safe-to-disconnect] Fix a bug in ReplaceChannelForEndpoint()

PiperOrigin-RevId: 679298374
This commit is contained in:
hai007
2024-09-26 14:37:57 -07:00
committed by Copybara-Service
parent 1a87241ffd
commit 53b09d452b
5 changed files with 125 additions and 86 deletions
@@ -19,6 +19,8 @@
#include <utility>
#include "absl/time/time.h"
#include "connections/implementation/client_proxy.h"
#include "connections/implementation/endpoint_channel.h"
#include "connections/implementation/offline_frames.h"
#include "internal/platform/condition_variable.h"
#include "internal/platform/feature_flags.h"
@@ -37,10 +39,10 @@ const absl::Duration kDataTransferDelay = absl::Milliseconds(500);
}
EndpointChannelManager::~EndpointChannelManager() {
NEARBY_LOGS(INFO) << "Initiating shutdown of EndpointChannelManager.";
LOG(INFO) << "Initiating shutdown of EndpointChannelManager.";
MutexLock lock(&mutex_);
channel_state_.DestroyAll();
NEARBY_LOGS(INFO) << "EndpointChannelManager has shut down.";
LOG(INFO) << "EndpointChannelManager has shut down.";
}
void EndpointChannelManager::RegisterChannelForEndpoint(
@@ -48,22 +50,30 @@ void EndpointChannelManager::RegisterChannelForEndpoint(
std::unique_ptr<EndpointChannel> channel) {
MutexLock lock(&mutex_);
NEARBY_LOGS(INFO) << "EndpointChannelManager registered channel of type "
LOG(INFO) << "EndpointChannelManager registered channel of type "
<< channel->GetType() << " to endpoint " << endpoint_id;
SetActiveEndpointChannel(client, endpoint_id, std::move(channel),
true /* enable_encryption */);
NEARBY_LOGS(INFO) << "Registered channel: id=" << endpoint_id;
LOG(INFO) << "Registered channel: id=" << endpoint_id;
}
void EndpointChannelManager::ReplaceChannelForEndpoint(
ClientProxy* client, const std::string& endpoint_id,
std::unique_ptr<EndpointChannel> channel, bool enable_encryption) {
MutexLock lock(&mutex_);
if (client->IsSafeToDisconnectEnabled(endpoint_id) &&
channel_state_.IsWaitingForSafeToDisconnectTimeout(endpoint_id)) {
LOG(WARNING)
<< "EndpointChannelManager failed to replace endpoint " << endpoint_id
<< "'s channel with type " << channel->GetType()
<< " because the endpoint is waiting for active channel closure.";
return;
}
auto* endpoint = channel_state_.LookupEndpointData(endpoint_id);
if (endpoint != nullptr && endpoint->channel == nullptr) {
NEARBY_LOGS(INFO) << "EndpointChannelManager is missing channel while "
LOG(INFO) << "EndpointChannelManager is missing channel while "
"trying to update: endpoint "
<< endpoint_id;
}
@@ -88,7 +98,7 @@ std::shared_ptr<EndpointChannel> EndpointChannelManager::GetChannelForEndpoint(
auto* endpoint = channel_state_.LookupEndpointData(endpoint_id);
if (endpoint == nullptr) {
NEARBY_LOGS(INFO) << "No channel info for endpoint " << endpoint_id;
LOG(INFO) << "No channel info for endpoint " << endpoint_id;
return {};
}
@@ -196,7 +206,7 @@ void EndpointChannelManager::ChannelState::UpdateEncryptionContextForEndpoint(
void EndpointChannelManager::ChannelState::UpdateSafeToDisconnectForEndpoint(
const std::string& endpoint_id,
bool safe_to_disconnect_enabled) {
NEARBY_LOGS(INFO) << "[safe-to-disconnect] "
LOG(INFO) << "[safe-to-disconnect] "
"UpdateSafeToDisconnectForEndpoint for: "
<< endpoint_id << " " << safe_to_disconnect_enabled;
@@ -208,7 +218,7 @@ bool EndpointChannelManager::ChannelState::GetSafeToDisconnectForEndpoint(
const std::string& endpoint_id) {
auto item = endpoints_.find(endpoint_id);
if (item == endpoints_.end()) return false;
NEARBY_LOGS(INFO) << "[safe-to-disconnect] GetSafeToDisconnectForEndpoint: "
LOG(INFO) << "[safe-to-disconnect] GetSafeToDisconnectForEndpoint: "
<< item->second.safe_to_disconnect_enabled;
return item->second.safe_to_disconnect_enabled;
}
@@ -230,18 +240,18 @@ bool EndpointChannelManager::ChannelState::RemoveEndpoint(
// we resume to ensure the thread won't hang when trying to write to it.
channel->Resume();
NEARBY_LOGS(INFO) << "[safe-to-disconnect] Sending DISCONNECTION frame"
LOG(INFO) << "[safe-to-disconnect] Sending DISCONNECTION frame"
" with request 0, ack 0";
channel->Write(
parser::ForDisconnection(/* request_safe_to_disconnect */ false,
/* ack_safe_to_disconnect */ false));
NEARBY_LOGS(INFO)
LOG(INFO)
<< "EndpointChannelManager reported the disconnection to endpoint "
<< endpoint_id;
SystemClock::Sleep(kDataTransferDelay);
}
NEARBY_LOGS(INFO) << "Remove Endpoint: " << endpoint_id;
LOG(INFO) << "Remove Endpoint: " << endpoint_id;
endpoints_.erase(item);
return true;
}
@@ -251,7 +261,7 @@ bool EndpointChannelManager::ChannelState::isWifiLanConnected() const {
auto channel = endpoint.second.channel;
if (channel) {
if (channel->GetMedium() == Medium::WIFI_LAN) {
NEARBY_LOGS(INFO) << "Found WIFI_LAN Medium for endpoint:"
LOG(INFO) << "Found WIFI_LAN Medium for endpoint:"
<< endpoint.first;
return true;
}
@@ -266,7 +276,7 @@ void EndpointChannelManager::ChannelState::MarkEndpointStopWaitToDisconnect(
bool notify_stop_waiting) {
auto item = endpoints_.find(endpoint_id);
if (item == endpoints_.end()) return;
NEARBY_LOGS(INFO) << "[safe-to-disconnect] is_safe_to_disconnect= "
LOG(INFO) << "[safe-to-disconnect] is_safe_to_disconnect= "
<< is_safe_to_disconnect
<< ", notify_stop_waiting= " << notify_stop_waiting
<< " for endpoint: " << endpoint_id;
@@ -275,7 +285,7 @@ void EndpointChannelManager::ChannelState::MarkEndpointStopWaitToDisconnect(
item->second.is_safe_to_disconnect = is_safe_to_disconnect;
if (!item->second.timeout_to_disconnected_enabled) return;
if (notify_stop_waiting) {
NEARBY_LOGS(INFO) << "[safe-to-disconnect] Notify stop "
LOG(INFO) << "[safe-to-disconnect] Notify stop "
"waiting before timeout.";
item->second.timeout_to_disconnected.Notify();
item->second.timeout_to_disconnected_notified = true;
@@ -287,7 +297,7 @@ bool EndpointChannelManager::ChannelState::CreateNewTimeoutDisconnectedState(
const std::string& endpoint_id, absl::Duration timeout_millis) {
auto item = endpoints_.find(endpoint_id);
if (item == endpoints_.end()) return false;
NEARBY_LOGS(INFO) << "[safe-to-disconnect] "
LOG(INFO) << "[safe-to-disconnect] "
"Create TimeoutDisconnectedState for endpoint: "
<< endpoint_id;
{
@@ -295,7 +305,7 @@ bool EndpointChannelManager::ChannelState::CreateNewTimeoutDisconnectedState(
item->second.timeout_to_disconnected_enabled = true;
item->second.timeout_to_disconnected_notified = false;
item->second.timeout_to_disconnected.Wait(timeout_millis);
NEARBY_LOGS(INFO) << "[safe-to-disconnect] Wait is done with "
LOG(INFO) << "[safe-to-disconnect] Wait is done with "
<< (item->second.timeout_to_disconnected_notified
? "notification"
: "timeout");
@@ -306,6 +316,19 @@ bool EndpointChannelManager::ChannelState::CreateNewTimeoutDisconnectedState(
}
return true;
}
bool EndpointChannelManager::ChannelState::IsWaitingForSafeToDisconnectTimeout(
const std::string& endpoint_id) {
auto item = endpoints_.find(endpoint_id);
if (item == endpoints_.end()) return false;
{
MutexLock lock(&item->second.timeout_to_disconnected_mutex);
LOG(INFO) << "[safe-to-disconnect] "
"IsWaitingForSafeToDisconnectTimeout for endpoint: "
<< endpoint_id << ": "
<< item->second.timeout_to_disconnected_enabled;
return (item->second.timeout_to_disconnected_enabled);
}
}
bool EndpointChannelManager::ChannelState::IsSafeToDisconnect(
const std::string& endpoint_id) {
@@ -314,7 +337,7 @@ bool EndpointChannelManager::ChannelState::IsSafeToDisconnect(
if (item == endpoints_.end()) return true;
{
MutexLock lock(&item->second.timeout_to_disconnected_mutex);
NEARBY_LOGS(INFO)
LOG(INFO)
<< "[safe-to-disconnect] Get SafeToDisconnect status for endpoint: "
<< endpoint_id << ": " << item->second.is_safe_to_disconnect;
return (item->second.is_safe_to_disconnect);
@@ -343,7 +366,7 @@ bool EndpointChannelManager::UnregisterChannelForEndpoint(
safe_to_disconnect_enabled, result)) {
return false;
}
NEARBY_LOGS(INFO)
LOG(INFO)
<< "EndpointChannelManager unregistered channel for endpoint "
<< endpoint_id;
return true;
@@ -195,6 +195,7 @@ class EndpointChannelManager final {
bool notify_stop_waiting);
bool CreateNewTimeoutDisconnectedState(const std::string& endpoint_id,
absl::Duration timeout_millis);
bool IsWaitingForSafeToDisconnectTimeout(const std::string& endpoint_id);
bool IsSafeToDisconnect(const std::string& endpoint_id);
void RemoveTimeoutDisconnectedState(const std::string& endpoint_id);
@@ -62,7 +62,9 @@ constexpr auto kEnableSafeToDisconnect =
// by default, enable Wi-Fi Hotspot client.
constexpr auto kEnableWifiHotspotClient =
flags::Flag<bool>(kConfigPackage, "45648734", true);
// Enable/Disable payload-received-ack feature.
// Set the safe-to-disconnect version.
// Enable 1. safe-to-disconnect check 2. reserved 3. auto-reconnect 4.
// auto-resume 5. non-distance-constraint-recovery 6. payload_ack
constexpr auto kSafeToDisconnectVersion =
flags::Flag<int64_t>(kConfigPackage, "45425841", 0);
// When true, use stable endpoint ID.
+77 -66
View File
@@ -82,7 +82,7 @@ bool PayloadManager::SendPayloadLoop(
// Update the still-active recipients of this payload.
if (available_endpoint_ids.empty()) {
NEARBY_LOGS(INFO)
LOG(INFO)
<< "PayloadManager short-circuiting payload_id="
<< pending_payload.GetInternalPayload()->GetId() << " after sending "
<< next_chunk_offset
@@ -93,7 +93,7 @@ bool PayloadManager::SendPayloadLoop(
// Check if the payload has been cancelled by the client and, if so,
// notify the remaining recipients.
if (pending_payload.IsLocallyCanceled()) {
NEARBY_LOGS(INFO) << "Aborting send of payload_id="
LOG(INFO) << "Aborting send of payload_id="
<< pending_payload.GetInternalPayload()->GetId()
<< " at offset " << next_chunk_offset
<< " since it is marked canceled.";
@@ -113,7 +113,7 @@ bool PayloadManager::SendPayloadLoop(
pending_payload.GetInternalPayload()->SkipToOffset(resume_offset);
if (!real_offset.ok()) {
// Stop sending since it may cause remote file merging failed.
NEARBY_LOGS(WARNING) << "PayloadManager failed to skip offset "
LOG(WARNING) << "PayloadManager failed to skip offset "
<< resume_offset << " on payload_id "
<< pending_payload.GetInternalPayload()->GetId();
HandleFinishedOutgoingPayload(
@@ -144,7 +144,7 @@ bool PayloadManager::SendPayloadLoop(
pending_payload.GetInternalPayload()->GetTotalSize() > 0 &&
pending_payload.GetInternalPayload()->GetTotalSize() <
next_chunk_offset) {
NEARBY_LOGS(INFO) << "Payload xfer failed: payload_id="
LOG(INFO) << "Payload xfer failed: payload_id="
<< pending_payload.GetInternalPayload()->GetId();
HandleFinishedOutgoingPayload(
client, available_endpoint_ids, payload_header, next_chunk_offset,
@@ -162,7 +162,7 @@ bool PayloadManager::SendPayloadLoop(
payload_header, payload_chunk, available_endpoint_ids, packet_meta_data);
// Check whether at least one endpoint failed.
if (!failed_endpoint_ids.empty()) {
NEARBY_LOGS(INFO) << "Payload xfer: endpoints failed: payload_id="
LOG(INFO) << "Payload xfer: endpoints failed: payload_id="
<< payload_header.id() << "; endpoint_ids={"
<< ToString(failed_endpoint_ids) << "}",
HandleFinishedOutgoingPayload(client, failed_endpoint_ids,
@@ -196,7 +196,7 @@ bool PayloadManager::SendPayloadLoop(
if (!next_chunk_size) {
// That was the last chunk, we're outta here.
NEARBY_LOGS(INFO) << "Payload xfer done: payload_id="
LOG(INFO) << "Payload xfer done: payload_id="
<< pending_payload.GetInternalPayload()->GetId()
<< "; size=" << next_chunk_offset;
ThroughputRecorderContainer::GetInstance()
@@ -297,7 +297,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_LOGS(INFO) << "CreateOutgoingPayload: payload_id=" << payload_id;
LOG(INFO) << "CreateOutgoingPayload: payload_id=" << payload_id;
MutexLock lock(&mutex_);
pending_payloads_.StartTrackingPayload(
payload_id,
@@ -316,7 +316,7 @@ PayloadManager::PayloadManager(EndpointManager& endpoint_manager)
}
void PayloadManager::CancelAllPayloads() {
NEARBY_LOGS(INFO) << "PayloadManager: canceling payloads; self=" << this;
LOG(INFO) << "PayloadManager: canceling payloads; self=" << this;
{
MutexLock lock(&mutex_);
int pending_outgoing_payloads = 0;
@@ -332,7 +332,7 @@ void PayloadManager::CancelAllPayloads() {
}
}
if (shutdown_barrier_) {
NEARBY_LOGS(INFO) << "PayloadManager: waiting for pending outgoing "
LOG(INFO) << "PayloadManager: waiting for pending outgoing "
"payloads; self="
<< this;
shutdown_barrier_->Await();
@@ -346,11 +346,11 @@ void PayloadManager::DisconnectFromEndpointManager() {
}
PayloadManager::~PayloadManager() {
NEARBY_LOGS(INFO) << "PayloadManager: going down; self=" << this;
LOG(INFO) << "PayloadManager: going down; self=" << this;
ThroughputRecorderContainer::GetInstance().Shutdown();
DisconnectFromEndpointManager();
CancelAllPayloads();
NEARBY_LOGS(INFO) << "PayloadManager: turn down payload executors; self="
LOG(INFO) << "PayloadManager: turn down payload executors; self="
<< this;
bytes_payload_executor_.Shutdown();
stream_payload_executor_.Shutdown();
@@ -362,7 +362,7 @@ PayloadManager::~PayloadManager() {
RunOnStatusUpdateThread(
"~payload-manager",
[this, &stop_latch]() RUN_ON_PAYLOAD_STATUS_UPDATE_THREAD() {
NEARBY_LOGS(INFO) << "PayloadManager: stop tracking payloads; self="
LOG(INFO) << "PayloadManager: stop tracking payloads; self="
<< this;
MutexLock lock(&mutex_);
pending_payloads_.StopTrackingAllPayloads();
@@ -370,19 +370,19 @@ PayloadManager::~PayloadManager() {
});
stop_latch.Await();
NEARBY_LOGS(INFO) << "PayloadManager: turn down notification executor; self="
LOG(INFO) << "PayloadManager: turn down notification executor; self="
<< this;
// Stop all the ongoing Runnables (as gracefully as possible).
payload_status_update_executor_.Shutdown();
NEARBY_LOGS(INFO) << "PayloadManager: down; self=" << this;
LOG(INFO) << "PayloadManager: down; self=" << this;
}
bool PayloadManager::NotifyShutdown() {
MutexLock lock(&mutex_);
if (!shutdown_.Get()) return false;
if (!shutdown_barrier_) return false;
NEARBY_LOGS(INFO) << "PayloadManager [shutdown mode]";
LOG(INFO) << "PayloadManager [shutdown mode]";
shutdown_barrier_->CountDown();
return true;
}
@@ -391,7 +391,7 @@ void PayloadManager::SendPayload(ClientProxy* client,
const EndpointIds& endpoint_ids,
Payload payload) {
if (shutdown_.Get()) return;
NEARBY_LOGS(INFO) << "SendPayload: endpoint_ids={" << ToString(endpoint_ids)
LOG(INFO) << "SendPayload: endpoint_ids={" << ToString(endpoint_ids)
<< "}";
// Before transfer to internal payload, retrieves the Payload size for
// analytics.
@@ -417,7 +417,7 @@ void PayloadManager::SendPayload(ClientProxy* client,
RecordInvalidPayloadAnalytics(client, endpoint_ids, payload.GetId(),
payload.GetType(), payload.GetOffset(),
payload_total_size);
NEARBY_LOGS(INFO)
LOG(INFO)
<< "PayloadManager failed to determine the right executor for "
"outgoing payload_id="
<< payload.GetId() << ", payload_type=" << ToString(payload.GetType());
@@ -445,7 +445,7 @@ void PayloadManager::SendPayload(ClientProxy* client,
RecordInvalidPayloadAnalytics(client, endpoint_ids, payload_id,
payload_type, resume_offset,
payload_total_size);
NEARBY_LOGS(INFO)
LOG(INFO)
<< "PayloadManager failed to create InternalPayload for outgoing "
"payload_id="
<< payload_id << ", payload_type=" << ToString(payload_type)
@@ -484,7 +484,7 @@ void PayloadManager::SendPayload(ClientProxy* client,
DestroyPendingPayload(payload_id);
});
});
NEARBY_LOGS(INFO) << "PayloadManager: xfer scheduled: self=" << this
LOG(INFO) << "PayloadManager: xfer scheduled: self=" << this
<< "; payload_id=" << payload_id
<< ", payload_type=" << ToString(payload_type);
}
@@ -498,14 +498,14 @@ Status PayloadManager::CancelPayload(ClientProxy* client,
Payload::Id payload_id) {
PendingPayloadHandle canceled_payload = GetPayload(payload_id);
if (!canceled_payload) {
NEARBY_LOGS(INFO) << "Client requested cancel for unknown payload_id="
LOG(INFO) << "Client requested cancel for unknown payload_id="
<< payload_id << ", ignoring.";
return {Status::kPayloadUnknown};
}
// Mark the payload as canceled.
canceled_payload->MarkLocallyCanceled();
NEARBY_LOGS(INFO) << "Cancelling "
LOG(INFO) << "Cancelling "
<< (canceled_payload->IsIncoming() ? "incoming"
: "outgoing")
<< " payload_id=" << payload_id << " at request of client.";
@@ -539,7 +539,7 @@ void PayloadManager::OnIncomingFrame(
is_last);
}
}
NEARBY_LOGS(INFO)
LOG(INFO)
<< "PayloadManager skipped process payloads before PCP connected, "
<< frame.payload_header().id();
return;
@@ -547,7 +547,7 @@ void PayloadManager::OnIncomingFrame(
switch (frame.packet_type()) {
case PayloadTransferFrame::CONTROL:
NEARBY_LOGS(INFO) << "PayloadManager::OnIncomingFrame [CONTROL]: self="
LOG(INFO) << "PayloadManager::OnIncomingFrame [CONTROL]: self="
<< this << "; endpoint_id=" << from_endpoint_id;
ProcessControlPacket(to_client, from_endpoint_id, frame);
break;
@@ -556,13 +556,13 @@ void PayloadManager::OnIncomingFrame(
packet_meta_data);
break;
case PayloadTransferFrame::PAYLOAD_ACK:
NEARBY_LOGS(INFO) << "[safe-to-disconnect][PAYLOAD_RECEIVED_ACK] sender "
LOG(INFO) << "[safe-to-disconnect][PAYLOAD_RECEIVED_ACK] sender "
"received payload ack from "
<< from_endpoint_id;
ProcessPayloadAckPacket(from_endpoint_id, frame);
break;
default:
NEARBY_LOGS(WARNING)
LOG(WARNING)
<< "PayloadManager: invalid frame; remote endpoint: self=" << this
<< "; endpoint_id=" << from_endpoint_id;
break;
@@ -648,7 +648,7 @@ PayloadManager::EndpointInfoStatusToPayloadStatus(EndpointInfo::Status status) {
case EndpointInfo::Status::kAvailable:
return location::nearby::proto::connections::PayloadStatus::SUCCESS;
default:
NEARBY_LOGS(INFO) << "PayloadManager: Unknown PayloadStatus";
LOG(INFO) << "PayloadManager: Unknown PayloadStatus";
return location::nearby::proto::connections::PayloadStatus::
UNKNOWN_PAYLOAD_STATUS;
}
@@ -664,7 +664,7 @@ PayloadManager::ControlMessageEventToPayloadStatus(
return location::nearby::proto::connections::PayloadStatus::
REMOTE_CANCELLATION;
default:
NEARBY_LOGS(INFO) << "PayloadManager: unknown event=" << event;
LOG(INFO) << "PayloadManager: unknown event=" << event;
return location::nearby::proto::connections::PayloadStatus::
UNKNOWN_PAYLOAD_STATUS;
}
@@ -755,7 +755,7 @@ PayloadManager::PendingPayloadHandle PayloadManager::CreateIncomingPayload(
}
Payload::Id payload_id = internal_payload->GetId();
NEARBY_LOGS(INFO) << "CreateIncomingPayload: payload_id=" << payload_id;
LOG(INFO) << "CreateIncomingPayload: payload_id=" << payload_id;
pending_payloads_.StartTrackingPayload(
payload_id,
std::make_unique<PendingPayload>(
@@ -765,7 +765,7 @@ PayloadManager::PendingPayloadHandle PayloadManager::CreateIncomingPayload(
}
void PayloadManager::OnPendingPayloadDestroy(const PendingPayload* payload) {
NEARBY_LOGS(INFO) << "PayloadManager: destroying " << payload->ToString()
LOG(INFO) << "PayloadManager: destroying " << payload->ToString()
<< " self=" << this;
ThroughputRecorderContainer::GetInstance().StopTPRecorder(
payload->GetId(), payload->IsIncoming()
@@ -879,14 +879,13 @@ void PayloadManager::SendPayloadReceivedAck(ClientProxy* client,
"send_payload_ack", [this, &pending_payload, endpoint_id]() {
endpoint_manager_->SendPayloadAck(pending_payload.GetId(),
{endpoint_id});
NEARBY_LOGS(INFO) << "[safe-to-disconnect] Send "
LOG(INFO) << "[safe-to-disconnect] Send "
"PAYLOAD_RECEIVED_ACK frame to: "
<< endpoint_id << " done";
});
// Send the PAYLOAD_RECEIVED_ACK to the remote endpoint for the sender asap.
NEARBY_LOGS(INFO) << "[safe-to-disconnect] [PAYLOAD_RECEIVED_ACK] "
"isLastChunk, receiver send payload ack to "
<< endpoint_id;
LOG(INFO) << "[safe-to-disconnect] " << pending_payload.GetId()
<< " isLastChunk, receiver send ack to " << endpoint_id;
}
bool PayloadManager::WaitForReceivedAck(
@@ -899,7 +898,7 @@ bool PayloadManager::WaitForReceivedAck(
return true;
}
NEARBY_LOGS(INFO) << "[safe-to-disconnect] Last Chunk, sender wait for "
LOG(INFO) << "[safe-to-disconnect] Last Chunk, sender wait for "
"PAYLOAD_RECEIVED_ACK frame from: "
<< endpoint_id;
while (true) {
@@ -907,7 +906,7 @@ bool PayloadManager::WaitForReceivedAck(
GetPayload(payload_header.id());
// Make sure we're still tracking this payload and its associated endpoint.
if (!latest_pending_payload) {
NEARBY_LOGS(INFO) << "[safe-to-disconnect] short-circuiting "
LOG(INFO) << "[safe-to-disconnect] short-circuiting "
"latest_pending_payload is null for "
<< payload_header.id() << ", stop wait ack.";
return false;
@@ -915,7 +914,7 @@ bool PayloadManager::WaitForReceivedAck(
auto* endpoint_info = latest_pending_payload->GetEndpoint(endpoint_id);
if (endpoint_info == nullptr) {
NEARBY_LOGS(INFO) << "[safe-to-disconnect] short-circuiting "
LOG(INFO) << "[safe-to-disconnect] short-circuiting "
"endpointInfo is null for "
<< payload_header.id() << ", stop wait ack.";
return false;
@@ -927,7 +926,7 @@ bool PayloadManager::WaitForReceivedAck(
payload_chunk_offset,
location::nearby::proto::connections::
PayloadStatus::LOCAL_CANCELLATION);
NEARBY_LOGS(INFO) << "[safe-to-disconnect] short-circuiting local "
LOG(INFO) << "[safe-to-disconnect] short-circuiting local "
"payload cancellation for "
<< payload_header.id() << ", stop wait ack.";
return false;
@@ -938,7 +937,7 @@ bool PayloadManager::WaitForReceivedAck(
HandleFinishedOutgoingPayload(
client, {endpoint_id}, payload_header, payload_chunk_offset,
EndpointInfoStatusToPayloadStatus(endpoint_info->status.Get()));
NEARBY_LOGS(INFO) << "[safe-to-disconnect] short-circuiting remote "
LOG(INFO) << "[safe-to-disconnect] short-circuiting remote "
"payload cancellation for "
<< payload_header.id() << ", stop wait ack.";
return false;
@@ -946,6 +945,9 @@ bool PayloadManager::WaitForReceivedAck(
{
MutexLock lock(&endpoint_info->payload_received_ack_mutex);
if (endpoint_info->is_payload_received_ack) {
LOG(INFO) << "[safe-to-disconnect][PAYLOAD_RECEIVED_ACK] sender already"
" received payload ack from "
<< endpoint_id << ", stop wait PAYLOAD_RECEIVED_ACK.";
endpoint_info->is_payload_received_ack = false;
return true;
}
@@ -953,18 +955,27 @@ bool PayloadManager::WaitForReceivedAck(
FeatureFlags::GetInstance()
.GetFlags()
.wait_payload_received_ack_millis);
endpoint_info->is_payload_received_ack = false;
if (!wait_exception.Ok()) {
NEARBY_LOGS(INFO)
endpoint_info->is_payload_received_ack = false;
LOG(INFO)
<< "[safe-to-disconnect][PAYLOAD_RECEIVED_ACK] sender wait for "
"received payload ack from "
<< endpoint_id << " end with exception: " << wait_exception.value;
return false;
}
NEARBY_LOGS(INFO) << "[safe-to-disconnect][PAYLOAD_RECEIVED_ACK] sender "
"received payload ack from "
<< endpoint_id;
return true;
if (endpoint_info->is_payload_received_ack) {
LOG(INFO) << "[safe-to-disconnect][PAYLOAD_RECEIVED_ACK] Received "
"notification that sender "
"received payload ack from "
<< endpoint_id;
endpoint_info->is_payload_received_ack = false;
return true;
} else {
LOG(INFO) << "[safe-to-disconnect][PAYLOAD_RECEIVED_ACK] sender doesn't"
" received payload ack from "
<< endpoint_id << ", end with timeout.";
return false;
}
}
}
return true;
@@ -1000,7 +1011,7 @@ void PayloadManager::HandleFinishedOutgoingPayload(
break;
case location::nearby::proto::connections::PayloadStatus::
LOCAL_CANCELLATION:
NEARBY_LOGS(INFO)
LOG(INFO)
<< "Sending PAYLOAD_CANCEL to receiver side; payload_id="
<< payload_header.id();
SendControlMessage(
@@ -1022,7 +1033,7 @@ void PayloadManager::HandleFinishedOutgoingPayload(
// No special handling needed for these.
break;
default:
NEARBY_LOGS(INFO)
LOG(INFO)
<< "PayloadManager: Unhandled finished outgoing payload with "
"payload_status="
<< status;
@@ -1050,7 +1061,7 @@ void PayloadManager::HandleFinishedIncomingPayload(
PayloadTransferFrame::ControlMessage::PAYLOAD_CANCELED);
break;
default:
NEARBY_LOGS(INFO) << "Unhandled finished incoming payload_id="
LOG(INFO) << "Unhandled finished incoming payload_id="
<< payload_header.id()
<< " with payload_status=" << status;
break;
@@ -1089,7 +1100,7 @@ void PayloadManager::HandleSuccessfulOutgoingChunk(
PayloadTransferFrame::PayloadTransferFrame::PayloadHeader::
FILE) {
if (!is_last_chunk && payload_chunk_offset != 0) {
NEARBY_LOGS(INFO) << "Skip the outgoing chunk update with offset="
LOG(INFO) << "Skip the outgoing chunk update with offset="
<< payload_chunk_offset;
client->GetAnalyticsRecorder().OnPayloadChunkSent(
endpoint_id, payload_header.id(), payload_chunk_body_size);
@@ -1100,7 +1111,7 @@ void PayloadManager::HandleSuccessfulOutgoingChunk(
PendingPayloadHandle pending_payload = GetPayload(payload_header.id());
if (!pending_payload || !pending_payload->GetEndpoint(endpoint_id)) {
NEARBY_LOGS(INFO)
LOG(INFO)
<< "HandleSuccessfulOutgoingChunk: endpoint not found: "
"endpoint_id="
<< endpoint_id;
@@ -1173,7 +1184,7 @@ void PayloadManager::HandleSuccessfulIncomingChunk(
PayloadTransferFrame::PayloadTransferFrame::PayloadHeader::
FILE) {
if (!is_last_chunk && payload_chunk_offset != 0) {
NEARBY_LOGS(INFO) << "Skip the incoming chunk update with offset="
LOG(INFO) << "Skip the incoming chunk update with offset="
<< payload_chunk_offset;
client->GetAnalyticsRecorder().OnPayloadChunkReceived(
endpoint_id, payload_header.id(), payload_chunk_body_size);
@@ -1225,7 +1236,7 @@ void PayloadManager::ProcessDataPacket(
<< payload_chunk.offset();
// We explicitly deny payloads with ID 0.
if (payload_header.id() == 0) {
NEARBY_LOGS(WARNING) << "Denying payload with ID 0 for endpoint_id="
LOG(WARNING) << "Denying payload with ID 0 for endpoint_id="
<< from_endpoint_id << ", aborting receipt.";
// Send the error to the remote endpoint.
SendControlMessage({from_endpoint_id}, payload_header,
@@ -1255,7 +1266,7 @@ void PayloadManager::ProcessDataPacket(
pending_payload =
CreateIncomingPayload(payload_transfer_frame, from_endpoint_id);
if (!pending_payload) {
NEARBY_LOGS(WARNING)
LOG(WARNING)
<< "PayloadManager failed to create InternalPayload from "
"PayloadTransferFrame with payload_id="
<< payload_header.id() << " and type " << payload_header.type()
@@ -1273,7 +1284,7 @@ void PayloadManager::ProcessDataPacket(
pending_payload = GetPayload(payload_id)]()
RUN_ON_PAYLOAD_STATUS_UPDATE_THREAD() {
if (!pending_payload) return;
NEARBY_LOGS(INFO)
LOG(INFO)
<< "PayloadManager received new payload_id="
<< pending_payload->GetInternalPayload()->GetId()
<< " from endpoint_id=" << from_endpoint_id;
@@ -1286,7 +1297,7 @@ void PayloadManager::ProcessDataPacket(
}
if (!pending_payload) {
NEARBY_LOGS(WARNING) << "ProcessDataPacket: [missing] endpoint_id="
LOG(WARNING) << "ProcessDataPacket: [missing] endpoint_id="
<< from_endpoint_id
<< "; payload_id=" << payload_header.id();
return;
@@ -1294,7 +1305,7 @@ 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_LOGS(INFO) << "ProcessDataPacket: [cancel] endpoint_id="
LOG(INFO) << "ProcessDataPacket: [cancel] endpoint_id="
<< from_endpoint_id
<< "; payload_id=" << pending_payload->GetId();
HandleFinishedIncomingPayload(to_client, from_endpoint_id, payload_header,
@@ -1319,7 +1330,7 @@ void PayloadManager::ProcessDataPacket(
if (pending_payload->GetInternalPayload()
->AttachNextChunk(ByteArray(std::move(*payload_chunk.mutable_body())))
.Raised()) {
NEARBY_LOGS(ERROR) << "ProcessDataPacket: [data: error] endpoint_id="
LOG(ERROR) << "ProcessDataPacket: [data: error] endpoint_id="
<< from_endpoint_id
<< "; payload_id=" << pending_payload->GetId();
HandleFinishedIncomingPayload(
@@ -1357,7 +1368,7 @@ void PayloadManager::ProcessControlPacket(
payload_transfer_frame.control_message();
PendingPayloadHandle pending_payload = GetPayload(payload_header.id());
if (!pending_payload) {
NEARBY_LOGS(INFO) << "Got ControlMessage for unknown payload_id="
LOG(INFO) << "Got ControlMessage for unknown payload_id="
<< payload_header.id()
<< ", ignoring: " << control_message.event();
return;
@@ -1366,7 +1377,7 @@ void PayloadManager::ProcessControlPacket(
switch (control_message.event()) {
case PayloadTransferFrame::ControlMessage::PAYLOAD_CANCELED:
if (pending_payload->IsIncoming()) {
NEARBY_LOGS(INFO) << "Incoming PAYLOAD_CANCELED: from endpoint_id="
LOG(INFO) << "Incoming PAYLOAD_CANCELED: from endpoint_id="
<< from_endpoint_id << "; self=" << this;
// No need to mark the pending payload as cancelled, since this is a
// remote cancellation for an incoming payload -- we handle everything
@@ -1376,7 +1387,7 @@ void PayloadManager::ProcessControlPacket(
control_message.offset(),
ControlMessageEventToPayloadStatus(control_message.event()));
} else {
NEARBY_LOGS(INFO) << "Outgoing PAYLOAD_CANCELED: from endpoint_id="
LOG(INFO) << "Outgoing PAYLOAD_CANCELED: from endpoint_id="
<< from_endpoint_id << "; self=" << this;
// Mark the payload as canceled *for this endpoint*.
pending_payload->SetEndpointStatusFromControlMessage(from_endpoint_id,
@@ -1400,7 +1411,7 @@ void PayloadManager::ProcessControlPacket(
}
break;
default:
NEARBY_LOGS(INFO) << "Unhandled control message "
LOG(INFO) << "Unhandled control message "
<< control_message.event() << " for payload_id="
<< pending_payload->GetInternalPayload()->GetId();
break;
@@ -1413,19 +1424,19 @@ void PayloadManager::ProcessPayloadAckPacket(
auto payload_header = payload_transfer_frame.payload_header();
PendingPayloadHandle pending_payload = GetPayload(payload_header.id());
if (!pending_payload) {
NEARBY_LOGS(INFO) << "[safe-to-disconnect][PAYLOAD_RECEIVED_ACK] "
LOG(INFO) << "[safe-to-disconnect][PAYLOAD_RECEIVED_ACK] "
"short-circuiting got payload "
"ack for unknown payload "
<< payload_header.id() << ", ignoring";
return;
}
if (pending_payload->IsIncoming()) {
NEARBY_LOGS(INFO) << "[safe-to-disconnect][PAYLOAD_RECEIVED_ACK] "
LOG(INFO) << "[safe-to-disconnect][PAYLOAD_RECEIVED_ACK] "
"short-circuiting got Payload "
"ack for incoming payload "
<< payload_header.id() << ", ignoring";
}
NEARBY_LOGS(INFO)
LOG(INFO)
<< "[safe-to-disconnect][PAYLOAD_RECEIVED_ACK] sender received payload "
<< payload_header.id() << " ack from " << from_endpoint_id;
pending_payload->MarkReceivedAckFromEndpoint(from_endpoint_id);
@@ -1492,7 +1503,7 @@ PayloadManager::EndpointInfo::ControlMessageEventToEndpointInfoStatus(
case PayloadTransferFrame::ControlMessage::PAYLOAD_CANCELED:
return Status::kCanceled;
default:
NEARBY_LOGS(INFO)
LOG(INFO)
<< "Unknown EndpointInfo.Status for ControlMessage.EventType "
<< event;
return Status::kUnknown;
@@ -1645,7 +1656,7 @@ void PayloadManager::PendingPayloads::StartTrackingPayload(
// If the |payload_id| is being re-used, always prefer the newer payload.
Remove(pending_payloads_.find(payload_id));
NEARBY_LOGS(INFO) << "StartTrackingPayload: " << pending_payload->ToString();
LOG(INFO) << "StartTrackingPayload: " << pending_payload->ToString();
pending_payload->IncRefCount();
pending_payloads_[payload_id] = std::move(pending_payload);
}
@@ -1653,7 +1664,7 @@ void PayloadManager::PendingPayloads::StartTrackingPayload(
void PayloadManager::PendingPayloads::StopTrackingPayload(
Payload::Id payload_id) {
MutexLock lock(&mutex_);
NEARBY_LOGS(INFO) << "StopTrackingPayload " << payload_id;
LOG(INFO) << "StopTrackingPayload " << payload_id;
Remove(pending_payloads_.find(payload_id));
}
+3 -1
View File
@@ -71,6 +71,8 @@ class FeatureFlags {
// DiscoverPeripheralTracker flow.
bool enable_invoking_legacy_device_discovered_cb = false;
// Enable 1. safe-to-disconnect check 2. reserved 3. auto-reconnect 4.
// auto-resume 5. non-distance-constraint-recovery 6. payload_ack
std::int32_t min_nc_version_supports_safe_to_disconnect = 1;
std::int32_t min_nc_version_supports_auto_reconnect = 3;
absl::Duration auto_reconnect_retry_delay_millis = absl::Milliseconds(5000);
@@ -81,7 +83,7 @@ class FeatureFlags {
// Android code won't be able to launch "payload_received_ack" feature for
// in near future, so change "payload_received_ack" version from "2" to "5"
// after auto-reconnect and auto-resume.
std::int32_t min_nc_version_supports_payload_received_ack = 5;
std::int32_t min_nc_version_supports_payload_received_ack = 6;
// If the other part doesn't ack the safe_to_disconnect request, the
// initiator will end the connection in 30s.
absl::Duration safe_to_disconnect_ack_delay_millis =