Fixed transfer sync issue in payload callback

PiperOrigin-RevId: 540333269
This commit is contained in:
Guogang Li
2023-06-14 11:37:45 -07:00
committed by Copybara-Service
parent cf8afb73aa
commit f05bf7fc40
3 changed files with 68 additions and 6 deletions
@@ -40,6 +40,10 @@ constexpr auto kBlePeripheralLostTimeoutMillis =
constexpr auto kEnableGattQueryInThread =
flags::Flag<bool>(kConfigPackage, "45415261", false);
// Enable/Disable payload manager to skip chunk update.
constexpr auto kEnablePayloadManagerToSkipChunkUpdate =
flags::Flag<bool>(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,
+56 -6
View File
@@ -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
@@ -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