From f05bf7fc407af404d4de8a52b72cb30d459091e3 Mon Sep 17 00:00:00 2001 From: Guogang Li Date: Wed, 14 Jun 2023 11:36:28 -0700 Subject: [PATCH] Fixed transfer sync issue in payload callback PiperOrigin-RevId: 540333269 --- .../flags/nearby_connections_feature_flags.h | 4 ++ connections/implementation/payload_manager.cc | 62 +++++++++++++++++-- connections/implementation/payload_manager.h | 8 +++ 3 files changed, 68 insertions(+), 6 deletions(-) diff --git a/connections/implementation/flags/nearby_connections_feature_flags.h b/connections/implementation/flags/nearby_connections_feature_flags.h index cd2d771e..45d2bb42 100644 --- a/connections/implementation/flags/nearby_connections_feature_flags.h +++ b/connections/implementation/flags/nearby_connections_feature_flags.h @@ -40,6 +40,10 @@ constexpr auto kBlePeripheralLostTimeoutMillis = constexpr auto kEnableGattQueryInThread = flags::Flag(kConfigPackage, "45415261", false); +// Enable/Disable payload manager to skip chunk update. +constexpr auto kEnablePayloadManagerToSkipChunkUpdate = + flags::Flag(kConfigPackage, "45415729", false); + // LINT.ThenChange( // //depot/google3/location/nearby/cpp/sharing/clients/cpp/nearby_sharing_service_adapter_dart.h, // //depot/google3/location/nearby/cpp/sharing/clients/cpp/nearby_sharing_service_adapter_dart.cc, diff --git a/connections/implementation/payload_manager.cc b/connections/implementation/payload_manager.cc index 760c9240..fe42ce23 100644 --- a/connections/implementation/payload_manager.cc +++ b/connections/implementation/payload_manager.cc @@ -26,7 +26,9 @@ #include "absl/strings/str_cat.h" #include "absl/time/time.h" #include "connections/implementation/analytics/throughput_recorder.h" +#include "connections/implementation/flags/nearby_connections_feature_flags.h" #include "connections/implementation/internal_payload_factory.h" +#include "internal/flags/nearby_flags.h" #include "internal/platform/count_down_latch.h" #include "internal/platform/logging.h" #include "internal/platform/mutex_lock.h" @@ -872,6 +874,12 @@ void PayloadManager::HandleSuccessfulOutgoingChunk( const PayloadTransferFrame::PayloadHeader& payload_header, std::int32_t payload_chunk_flags, std::int64_t payload_chunk_offset, std::int64_t payload_chunk_body_size) { + if (NearbyFlags::GetInstance().GetBoolFlag( + config_package_nearby::nearby_connections_feature:: + kEnablePayloadManagerToSkipChunkUpdate)) { + MutexLock lock(&chunk_update_mutex_); + ++outgoing_chunk_update_count_; + } RunOnStatusUpdateThread( "outgoing-chunk-success", [this, client, endpoint_id, payload_header, payload_chunk_flags, @@ -879,6 +887,27 @@ void PayloadManager::HandleSuccessfulOutgoingChunk( payload_chunk_body_size]() RUN_ON_PAYLOAD_STATUS_UPDATE_THREAD() { // Make sure we're still tracking this payload and its associated // endpoint. + bool is_last_chunk = + (payload_chunk_flags & + PayloadTransferFrame::PayloadChunk::LAST_CHUNK) != 0; + + if (NearbyFlags::GetInstance().GetBoolFlag( + config_package_nearby::nearby_connections_feature:: + kEnablePayloadManagerToSkipChunkUpdate)) { + MutexLock lock(&chunk_update_mutex_); + --outgoing_chunk_update_count_; + if (outgoing_chunk_update_count_ > 0 && payload_header.has_type() && + payload_header.type() == + PayloadTransferFrame::PayloadTransferFrame::PayloadHeader:: + FILE) { + if (!is_last_chunk && payload_chunk_offset != 0) { + NEARBY_LOGS(INFO) << "Skip the outgoing chunk update with offset=" + << payload_chunk_offset; + return; + } + } + } + PendingPayload* pending_payload = GetPayload(payload_header.id()); if (!pending_payload || !pending_payload->GetEndpoint(endpoint_id)) { NEARBY_LOGS(INFO) @@ -888,9 +917,6 @@ void PayloadManager::HandleSuccessfulOutgoingChunk( return; } - bool is_last_chunk = - (payload_chunk_flags & - PayloadTransferFrame::PayloadChunk::LAST_CHUNK) != 0; PayloadProgressInfo update{ payload_header.id(), is_last_chunk ? PayloadProgressInfo::Status::kSuccess @@ -944,20 +970,44 @@ void PayloadManager::HandleSuccessfulIncomingChunk( const PayloadTransferFrame::PayloadHeader& payload_header, std::int32_t payload_chunk_flags, std::int64_t payload_chunk_offset, std::int64_t payload_chunk_body_size) { + if (NearbyFlags::GetInstance().GetBoolFlag( + config_package_nearby::nearby_connections_feature:: + kEnablePayloadManagerToSkipChunkUpdate)) { + MutexLock lock(&chunk_update_mutex_); + ++incoming_chunk_update_count_; + } RunOnStatusUpdateThread( "incoming-chunk-success", [this, client, endpoint_id, payload_header, payload_chunk_flags, payload_chunk_offset, payload_chunk_body_size]() RUN_ON_PAYLOAD_STATUS_UPDATE_THREAD() { // Make sure we're still tracking this payload. + bool is_last_chunk = + (payload_chunk_flags & + PayloadTransferFrame::PayloadChunk::LAST_CHUNK) != 0; + + if (NearbyFlags::GetInstance().GetBoolFlag( + config_package_nearby::nearby_connections_feature:: + kEnablePayloadManagerToSkipChunkUpdate)) { + MutexLock lock(&chunk_update_mutex_); + --incoming_chunk_update_count_; + if (incoming_chunk_update_count_ > 0 && payload_header.has_type() && + payload_header.type() == + PayloadTransferFrame::PayloadTransferFrame::PayloadHeader:: + FILE) { + if (!is_last_chunk && payload_chunk_offset != 0) { + NEARBY_LOGS(INFO) << "Skip the incoming chunk update with offset=" + << payload_chunk_offset; + return; + } + } + } + PendingPayload* pending_payload = GetPayload(payload_header.id()); if (!pending_payload) { return; } - bool is_last_chunk = - (payload_chunk_flags & - PayloadTransferFrame::PayloadChunk::LAST_CHUNK) != 0; PayloadProgressInfo update{ payload_header.id(), is_last_chunk ? PayloadProgressInfo::Status::kSuccess diff --git a/connections/implementation/payload_manager.h b/connections/implementation/payload_manager.h index cef11252..6be21457 100644 --- a/connections/implementation/payload_manager.h +++ b/connections/implementation/payload_manager.h @@ -329,6 +329,14 @@ class PayloadManager : public EndpointManager::FrameProcessor { SingleThreadExecutor payload_status_update_executor_; EndpointManager* endpoint_manager_; + + // When callback processing cannot keep the speed of callback update, the + // callback thread will be lag to the real transfer. In order to keep sync + // between callback and sending/receiving threads, we will skip non-important + // callbacks during file transfer. + mutable Mutex chunk_update_mutex_; + int outgoing_chunk_update_count_ ABSL_GUARDED_BY(chunk_update_mutex_) = 0; + int incoming_chunk_update_count_ ABSL_GUARDED_BY(chunk_update_mutex_) = 0; }; } // namespace connections