Split PayloadTransferUpdate for incoming and outgoing sessions.

PiperOrigin-RevId: 657257085
This commit is contained in:
Francis Tsui
2024-07-29 10:57:51 -07:00
committed by Copybara-Service
parent 861ed1b64b
commit cd9ce0d924
5 changed files with 115 additions and 82 deletions
+74 -79
View File
@@ -817,7 +817,8 @@ void NearbySharingServiceImpl::Accept(
bool accept_success = incoming_session->AcceptTransfer(
context_->GetClock(),
absl::bind_front(
&NearbySharingServiceImpl::OnPayloadTransferUpdate, this));
&NearbySharingServiceImpl::IncomingPayloadTransferUpdate,
this));
std::move(status_codes_callback)(
accept_success ? StatusCodes::kOk
: StatusCodes::kOutOfOrderApiCall);
@@ -2947,7 +2948,7 @@ void NearbySharingServiceImpl::OnReceiveConnectionResponse(
std::optional<nearby::sharing::service::proto::V1Frame> frame) {
OnFrameRead(share_target_id, std::move(frame));
},
absl::bind_front(&NearbySharingServiceImpl::OnPayloadTransferUpdate,
absl::bind_front(&NearbySharingServiceImpl::OutgoingPayloadTransferUpdate,
this));
}
@@ -2972,7 +2973,7 @@ void NearbySharingServiceImpl::OnStorageCheckCompleted(
NL_LOG(INFO) << __func__ << ": Auto-accepting self share.";
session.AcceptTransfer(
context_->GetClock(),
absl::bind_front(&NearbySharingServiceImpl::OnPayloadTransferUpdate,
absl::bind_front(&NearbySharingServiceImpl::IncomingPayloadTransferUpdate,
this));
OnTransferStarted(/*is_incoming=*/true);
}
@@ -3099,9 +3100,9 @@ std::optional<ShareTarget> NearbySharingServiceImpl::CreateShareTarget(
return target;
}
void NearbySharingServiceImpl::OnPayloadTransferUpdate(
void NearbySharingServiceImpl::IncomingPayloadTransferUpdate(
int64_t share_target_id, TransferMetadata metadata) {
ShareSession* session = GetShareSession(share_target_id);
IncomingShareSession* session = GetIncomingShareSession(share_target_id);
if (!session) {
// ShareTarget already disconnected.
NL_LOG(WARNING)
@@ -3109,69 +3110,87 @@ void NearbySharingServiceImpl::OnPayloadTransferUpdate(
<< share_target_id;
return;
}
bool is_in_progress =
metadata.status() == TransferMetadata::Status::kInProgress;
// kInProgress status is logged extensively elsewhere so avoid the spam.
if (!is_in_progress) {
if (metadata.status() != TransferMetadata::Status::kInProgress) {
NL_VLOG(1) << __func__ << ": Nearby Share service: "
<< "Payload transfer update for share target with ID "
<< share_target_id << ": "
<< TransferMetadata::StatusToString(metadata.status());
}
bool payload_incomplete = false;
if (session->IsIncoming()) {
IncomingShareSession* incoming_session =
GetIncomingShareSession(share_target_id);
// Update file paths during progress. It may impact transfer speed.
// TODO: b/289290115 - Revisit UpdateFilePath to enhance transfer speed for
// MacOS.
if (update_file_paths_in_progress_) {
session->UpdateFilePayloadPaths();
}
// Update file paths during progress. It may impact transfer speed.
// TODO: b/289290115 - Revisit UpdateFilePath to enhance transfer speed for
// MacOS.
if (update_file_paths_in_progress_) {
incoming_session->UpdateFilePayloadPaths();
if (metadata.status() == TransferMetadata::Status::kComplete) {
if (!session->FinalizePayloads()) {
metadata = TransferMetadataBuilder()
.set_status(TransferMetadata::Status::kIncompletePayloads)
.build();
}
if (metadata.status() == TransferMetadata::Status::kComplete) {
if (!incoming_session->FinalizePayloads()) {
payload_incomplete = true;
}
fast_initiation_scanner_cooldown_timer_ = std::make_unique<ThreadTimer>(
*service_thread_, "fast_initiation_scanner_cooldown_timer",
kFastInitiationScannerCooldown, [this]() {
fast_initiation_scanner_cooldown_timer_.reset();
InvalidateFastInitiationScanning();
});
} else if (metadata.status() == TransferMetadata::Status::kCancelled) {
NL_VLOG(1) << __func__ << ": Update file paths for cancelled transfer";
if (!update_file_paths_in_progress_) {
incoming_session->UpdateFilePayloadPaths();
}
fast_initiation_scanner_cooldown_timer_ = std::make_unique<ThreadTimer>(
*service_thread_, "fast_initiation_scanner_cooldown_timer",
kFastInitiationScannerCooldown, [this]() {
fast_initiation_scanner_cooldown_timer_.reset();
InvalidateFastInitiationScanning();
});
} else if (metadata.status() == TransferMetadata::Status::kCancelled) {
NL_VLOG(1) << __func__ << ": Update file paths for cancelled transfer";
if (!update_file_paths_in_progress_) {
session->UpdateFilePayloadPaths();
}
}
// Make sure to call this before calling Disconnect, or we risk losing some
// transfer updates in the receive case due to the Disconnect call cleaning up
// share targets.
session->UpdateTransferMetadata(
payload_incomplete
? TransferMetadataBuilder()
.set_status(TransferMetadata::Status::kIncompletePayloads)
.build()
: metadata);
session->UpdateTransferMetadata(metadata);
if (payload_incomplete ||
TransferMetadata::IsFinalStatus(metadata.status())) {
// final status already sent, no need to send again on disconnect.
session->set_disconnect_status(TransferMetadata::Status::kUnknown);
if (TransferMetadata::IsFinalStatus(metadata.status())) {
// Cancellation has its own disconnection strategy, possibly adding a
// delay before disconnection to provide the other party time to process
// the cancellation.
if (metadata.status() != TransferMetadata::Status::kCancelled) {
session->Disconnect();
}
}
// Cancellation has its own disconnection strategy, possibly adding a
// delay before disconnection to provide the other party time to process
// the cancellation.
if (TransferMetadata::IsFinalStatus(metadata.status()) &&
metadata.status() != TransferMetadata::Status::kCancelled) {
Disconnect(share_target_id, metadata);
}
void NearbySharingServiceImpl::OutgoingPayloadTransferUpdate(
int64_t share_target_id, TransferMetadata metadata) {
OutgoingShareSession* session = GetOutgoingShareSession(share_target_id);
if (!session) {
// ShareTarget already disconnected.
NL_LOG(WARNING)
<< "Received payload update after share target disconnected: "
<< share_target_id;
return;
}
// kInProgress status is logged extensively elsewhere so avoid the spam.
if (metadata.status() != TransferMetadata::Status::kInProgress) {
NL_VLOG(1) << __func__ << ": Nearby Share service: "
<< "Payload transfer update for share target with ID "
<< share_target_id << ": "
<< TransferMetadata::StatusToString(metadata.status());
}
// Make sure to call this before calling Disconnect, or we risk losing some
// transfer updates in the receive case due to the Disconnect call cleaning up
// share targets.
session->UpdateTransferMetadata(metadata);
if (TransferMetadata::IsFinalStatus(metadata.status())) {
// Cancellation has its own disconnection strategy, possibly adding a
// delay before disconnection to provide the other party time to process
// the cancellation.
if (metadata.status() != TransferMetadata::Status::kCancelled) {
DisconnectOutgoing(*session, metadata);
}
}
}
@@ -3198,40 +3217,18 @@ void NearbySharingServiceImpl::RemoveIncomingPayloads(
file_handler_.DeleteFilesFromDisk(std::move(files_for_deletion), []() {});
}
void NearbySharingServiceImpl::Disconnect(int64_t share_target_id,
TransferMetadata metadata) {
ShareSession* session = GetShareSession(share_target_id);
if (!session) {
NL_LOG(WARNING)
<< __func__
<< ": Failed to disconnect. No share session found for target - "
<< share_target_id;
return;
}
std::string endpoint_id = session->endpoint_id();
void NearbySharingServiceImpl::DisconnectOutgoing(OutgoingShareSession& session,
TransferMetadata metadata) {
std::string endpoint_id = session.endpoint_id();
// Failed to send or receive. No point in continuing, so disconnect
// immediately.
if (metadata.status() != TransferMetadata::Status::kComplete) {
if (session->IsConnected()) {
session->connection()->Close();
} else {
nearby_connections_manager_->Disconnect(endpoint_id);
}
return;
}
// Files received successfully. Receivers can immediately cancel.
if (session->IsIncoming()) {
if (session->IsConnected()) {
session->connection()->Close();
} else {
nearby_connections_manager_->Disconnect(endpoint_id);
}
session.Disconnect();
return;
}
// Disconnect after a timeout to make sure any pending payloads are sent.
// This can happen after the session has been destroyed.
auto timer = std::make_unique<ThreadTimer>(
*service_thread_, "disconnection_timeout_alarm",
kOutgoingDisconnectionDelay, [this, endpoint_id]() {
@@ -3240,8 +3237,6 @@ void NearbySharingServiceImpl::Disconnect(int64_t share_target_id,
});
disconnection_timeout_alarms_[endpoint_id] = std::move(timer);
session->set_disconnect_status(TransferMetadata::Status::kUnknown);
}
IncomingShareSession& NearbySharingServiceImpl::CreateIncomingShareSession(
+7 -3
View File
@@ -367,10 +367,14 @@ class NearbySharingServiceImpl
const std::optional<NearbyShareDecryptedPublicCertificate>& certificate,
bool is_incoming);
void OnPayloadTransferUpdate(int64_t share_target_id,
TransferMetadata metadata);
void IncomingPayloadTransferUpdate(
int64_t share_target_id, TransferMetadata metadata);
void OutgoingPayloadTransferUpdate(
int64_t share_target_id, TransferMetadata metadata);
void RemoveIncomingPayloads(const IncomingShareSession& session);
void Disconnect(int64_t share_target_id, TransferMetadata metadata);
void DisconnectOutgoing(OutgoingShareSession& session,
TransferMetadata metadata);
IncomingShareSession& CreateIncomingShareSession(
const ShareTarget& share_target, absl::string_view endpoint_id,
+8
View File
@@ -131,6 +131,14 @@ bool ShareSession::OnConnected(const NearbySharingDecoder& decoder,
return true;
}
void ShareSession::Disconnect() {
if (connection_ == nullptr) {
return;
}
// Do not clear connection_ here. It will be cleared in OnDisconnect().
connection_->Close();
}
void ShareSession::Abort(TransferMetadata::Status status) {
NL_DCHECK(TransferMetadata::IsFinalStatus(status))
<< "Abort should only be called with a final status";
+2
View File
@@ -144,6 +144,8 @@ class ShareSession {
void SetTokenForTests(std::string token) { token_ = std::move(token); }
void Disconnect();
protected:
virtual void InvokeTransferUpdateCallback(
const TransferMetadata& metadata) = 0;
+24
View File
@@ -441,5 +441,29 @@ TEST(ShareSessionTest, AbortConnected) {
EXPECT_TRUE(disconnected);
}
TEST(ShareSessionTest, Disconnect) {
FakeNearbyConnectionsManager connections_manager;
NearbySharingDecoderImpl nearby_sharing_decoder;
ShareTarget share_target;
TestShareSession session(std::string(kEndpointId), share_target);
FakeNearbyConnection connection;
bool disconnected = false;
connection.SetDisconnectionListener(
[&disconnected]() { disconnected = true; });
EXPECT_TRUE(session.OnConnected(nearby_sharing_decoder, absl::Now(),
&connections_manager, &connection));
session.Disconnect();
EXPECT_TRUE(disconnected);
}
TEST(ShareSessionTest, DisconnectNotConnected) {
ShareTarget share_target;
TestShareSession session(std::string(kEndpointId), share_target);
session.Disconnect();
}
} // namespace
} // namespace nearby::sharing