analytics: Add operation result code for Payload analytics, II

PiperOrigin-RevId: 704252835
This commit is contained in:
Edwin Wu
2024-12-09 05:53:23 -08:00
committed by Copybara-Service
parent c5fed0c0ed
commit f0fb9662e2
4 changed files with 261 additions and 175 deletions
@@ -55,6 +55,7 @@ namespace connections {
namespace {
using ::location::nearby::analytics::proto::ConnectionsLog;
using ::location::nearby::connections::OfflineFrame;
using ::location::nearby::connections::PayloadTransferFrame;
using ::location::nearby::connections::V1Frame;
using ::nearby::analytics::PacketMetaData;
using DisconnectionReason =
+166 -119
View File
@@ -23,22 +23,27 @@
#include <utility>
#include <vector>
#include "absl/container/flat_hash_map.h"
#include "absl/functional/any_invocable.h"
#include "absl/functional/bind_front.h"
#include "absl/strings/str_cat.h"
#include "absl/strings/str_format.h"
#include "absl/time/clock.h"
#include "absl/time/time.h"
#include "connections/implementation/analytics/packet_meta_data.h"
#include "connections/implementation/analytics/throughput_recorder.h"
#include "connections/implementation/client_proxy.h"
#include "connections/implementation/endpoint_channel_manager.h"
#include "connections/implementation/endpoint_manager.h"
#include "connections/implementation/flags/nearby_connections_feature_flags.h"
#include "connections/implementation/internal_payload.h"
#include "connections/implementation/internal_payload_factory.h"
#include "connections/implementation/proto/offline_wire_formats.pb.h"
#include "connections/listeners.h"
#include "connections/medium_selector.h"
#include "connections/payload.h"
#include "connections/payload_type.h"
#include "connections/status.h"
#include "internal/flags/nearby_flags.h"
#include "internal/platform/byte_array.h"
#include "internal/platform/count_down_latch.h"
@@ -52,22 +57,20 @@
namespace nearby {
namespace connections {
using ::location::nearby::connections::OfflineFrame;
using ::location::nearby::connections::V1Frame;
using ::location::nearby::proto::connections::PayloadStatus;
using ::nearby::analytics::PacketMetaData;
using ::nearby::analytics::ThroughputRecorderContainer;
using ::nearby::connections::PayloadDirection;
namespace {
using ::location::nearby::connections::OfflineFrame;
using ::location::nearby::connections::PayloadTransferFrame;
using ::location::nearby::connections::V1Frame;
using ::location::nearby::proto::connections::Medium;
using ::location::nearby::proto::connections::OperationResultCode;
using ::location::nearby::proto::connections::PayloadStatus;
using PacketMetaData = ::nearby::analytics::PacketMetaData;
using ::nearby::analytics::ThroughputRecorderContainer;
using PayloadDirection = ::nearby::connections::PayloadDirection;
constexpr absl::Duration kMinTransferUpdateInterval = absl::Milliseconds(50);
}
} // namespace
// C++14 requires to declare this.
// TODO(apolyudov): remove when migration to c++17 is possible.
constexpr absl::Duration PayloadManager::kWaitCloseTimeout;
// TODO(edwinwu): Add OperationResultCode to analytics.
bool PayloadManager::SendPayloadLoop(
ClientProxy* client, PendingPayload& pending_payload,
PayloadTransferFrame::PayloadHeader& payload_header,
@@ -83,6 +86,7 @@ bool PayloadManager::SendPayloadLoop(
for (const auto& endpoint : unavailable_endpoints) {
HandleFinishedOutgoingPayload(
client, {endpoint->id}, payload_header, next_chunk_offset,
EndpointInfoStatusToOperationResultCode(endpoint->status.Get()),
EndpointInfoStatusToPayloadStatus(endpoint->status.Get()));
}
@@ -101,10 +105,10 @@ bool PayloadManager::SendPayloadLoop(
LOG(INFO) << "Aborting send of payload_id="
<< pending_payload.GetInternalPayload()->GetId() << " at offset "
<< next_chunk_offset << " since it is marked canceled.";
HandleFinishedOutgoingPayload(client, available_endpoint_ids,
payload_header, next_chunk_offset,
location::nearby::proto::connections::
PayloadStatus::LOCAL_CANCELLATION);
HandleFinishedOutgoingPayload(
client, available_endpoint_ids, payload_header, next_chunk_offset,
OperationResultCode::CLIENT_CANCELLATION_LOCAL_CANCEL_PAYLOAD,
PayloadStatus::LOCAL_CANCELLATION);
return false;
}
@@ -120,9 +124,10 @@ bool PayloadManager::SendPayloadLoop(
LOG(WARNING) << "PayloadManager failed to skip offset " << resume_offset
<< " on payload_id "
<< pending_payload.GetInternalPayload()->GetId();
HandleFinishedOutgoingPayload(
client, available_endpoint_ids, payload_header, next_chunk_offset,
location::nearby::proto::connections::PayloadStatus::LOCAL_ERROR);
HandleFinishedOutgoingPayload(client, available_endpoint_ids,
payload_header, next_chunk_offset,
OperationResultCode::IO_FILE_READING_ERROR,
PayloadStatus::LOCAL_ERROR);
return false;
}
NEARBY_VLOG(1) << "PayloadManager successfully skipped "
@@ -152,7 +157,7 @@ bool PayloadManager::SendPayloadLoop(
<< pending_payload.GetInternalPayload()->GetId();
HandleFinishedOutgoingPayload(
client, available_endpoint_ids, payload_header, next_chunk_offset,
location::nearby::proto::connections::PayloadStatus::LOCAL_ERROR);
OperationResultCode::IO_FILE_READING_ERROR, PayloadStatus::LOCAL_ERROR);
return false;
}
@@ -169,10 +174,10 @@ bool PayloadManager::SendPayloadLoop(
LOG(INFO) << "Payload xfer: endpoints failed: payload_id="
<< payload_header.id() << "; endpoint_ids={"
<< ToString(failed_endpoint_ids) << "}",
HandleFinishedOutgoingPayload(client, failed_endpoint_ids,
payload_header, next_chunk_offset,
location::nearby::proto::connections::
PayloadStatus::ENDPOINT_IO_ERROR);
HandleFinishedOutgoingPayload(
client, failed_endpoint_ids, payload_header, next_chunk_offset,
OperationResultCode::CONNECTIVITY_GENERIC_WRITING_CHANNEL_IO_ERROR,
PayloadStatus::ENDPOINT_IO_ERROR);
}
bool is_last_chunk = IsLastChunk(payload_chunk);
// Check whether at least one endpoint succeeded -- if they all failed,
@@ -414,9 +419,10 @@ void PayloadManager::SendPayload(ClientProxy* client,
// with. This should never be reached since the ServiceControllerRouter has
// already checked whether or not we can work with this Payload type.
if (!executor) {
RecordInvalidPayloadAnalytics(client, endpoint_ids, payload.GetId(),
payload.GetType(), payload.GetOffset(),
payload_total_size);
RecordInvalidPayloadAnalytics(
client, endpoint_ids, payload.GetId(), payload.GetType(),
payload.GetOffset(), payload_total_size,
OperationResultCode::NEARBY_GENERIC_OUTGOING_PAYLOAD_CREATION_FAILURE);
LOG(INFO) << "PayloadManager failed to determine the right executor for "
"outgoing payload_id="
<< payload.GetId()
@@ -442,9 +448,11 @@ void PayloadManager::SendPayload(ClientProxy* client,
if (shutdown_.Get()) return;
PendingPayloadHandle pending_payload = GetPayload(payload_id);
if (!pending_payload) {
RecordInvalidPayloadAnalytics(client, endpoint_ids, payload_id,
payload_type, resume_offset,
payload_total_size);
RecordInvalidPayloadAnalytics(
client, endpoint_ids, payload_id, payload_type, resume_offset,
payload_total_size,
OperationResultCode::
NEARBY_GENERIC_OUTGOING_PAYLOAD_CREATION_FAILURE);
LOG(INFO)
<< "PayloadManager failed to create InternalPayload for outgoing "
"payload_id="
@@ -514,11 +522,11 @@ Status PayloadManager::CancelPayload(ClientProxy* client,
}
// @EndpointManagerDataPool
void PayloadManager::OnIncomingFrame(
OfflineFrame& offline_frame, const std::string& from_endpoint_id,
ClientProxy* to_client,
location::nearby::proto::connections::Medium current_medium,
PacketMetaData& packet_meta_data) {
void PayloadManager::OnIncomingFrame(OfflineFrame& offline_frame,
const std::string& from_endpoint_id,
ClientProxy* to_client,
Medium current_medium,
PacketMetaData& packet_meta_data) {
PayloadTransferFrame& frame =
*offline_frame.mutable_v1()->mutable_payload_transfer();
@@ -608,29 +616,34 @@ void PayloadManager::OnEndpointDisconnect(ClientProxy* client,
client->OnPayloadProgress(endpoint_id, update);
PayloadStatus payload_status;
OperationResultCode operation_result_code;
switch (reason) {
case DisconnectionReason::LOCAL_DISCONNECTION:
payload_status = PayloadStatus::LOCAL_CLIENT_DISCONNECTION;
operation_result_code =
OperationResultCode::CLIENT_CANCELLATION_LOCAL_DISCONNECT;
break;
case DisconnectionReason::REMOTE_DISCONNECTION:
payload_status = PayloadStatus::REMOTE_CLIENT_DISCONNECTION;
operation_result_code =
OperationResultCode::CLIENT_CANCELLATION_REMOTE_DISCONNECT;
break;
case DisconnectionReason::IO_ERROR:
default:
payload_status = PayloadStatus::ENDPOINT_IO_ERROR;
// TODO(edwinwu): Add for return result code
operation_result_code = OperationResultCode::DETAIL_UNKNOWN;
break;
}
if (pending_payload->IsIncoming()) {
client->GetAnalyticsRecorder().OnIncomingPayloadDone(
endpoint_id, pending_payload->GetId(), payload_status,
location::nearby::proto::connections::OperationResultCode::
DETAIL_UNKNOWN);
operation_result_code);
} else {
client->GetAnalyticsRecorder().OnOutgoingPayloadDone(
endpoint_id, pending_payload->GetId(), payload_status,
location::nearby::proto::connections::OperationResultCode::
DETAIL_UNKNOWN);
operation_result_code);
}
});
@@ -638,46 +651,69 @@ void PayloadManager::OnEndpointDisconnect(ClientProxy* client,
});
}
location::nearby::proto::connections::PayloadStatus
PayloadManager::EndpointInfoStatusToPayloadStatus(EndpointInfo::Status status) {
PayloadStatus PayloadManager::EndpointInfoStatusToPayloadStatus(
EndpointInfo::Status status) {
switch (status) {
case EndpointInfo::Status::kCanceled:
return location::nearby::proto::connections::PayloadStatus::
REMOTE_CANCELLATION;
return PayloadStatus::REMOTE_CANCELLATION;
case EndpointInfo::Status::kError:
return location::nearby::proto::connections::PayloadStatus::REMOTE_ERROR;
return PayloadStatus::REMOTE_ERROR;
case EndpointInfo::Status::kAvailable:
return location::nearby::proto::connections::PayloadStatus::SUCCESS;
return PayloadStatus::SUCCESS;
default:
LOG(INFO) << "PayloadManager: Unknown PayloadStatus";
return location::nearby::proto::connections::PayloadStatus::
UNKNOWN_PAYLOAD_STATUS;
return PayloadStatus::UNKNOWN_PAYLOAD_STATUS;
}
}
location::nearby::proto::connections::PayloadStatus
PayloadManager::ControlMessageEventToPayloadStatus(
OperationResultCode PayloadManager::EndpointInfoStatusToOperationResultCode(
EndpointInfo::Status status) {
switch (status) {
case EndpointInfo::Status::kCanceled:
return OperationResultCode::CLIENT_CANCELLATION_REMOTE_IN_CANCELED_STATE;
case EndpointInfo::Status::kError:
return OperationResultCode::NEARBY_GENERIC_REMOTE_ENDPOINT_STATUS_ERROR;
case EndpointInfo::Status::kAvailable:
return OperationResultCode::DETAIL_SUCCESS;
default:
LOG(INFO) << "PayloadManager: Unknown PayloadStatus";
return OperationResultCode::DETAIL_UNKNOWN;
}
}
PayloadStatus PayloadManager::ControlMessageEventToPayloadStatus(
PayloadTransferFrame::ControlMessage::EventType event) {
switch (event) {
case PayloadTransferFrame::ControlMessage::PAYLOAD_ERROR:
return location::nearby::proto::connections::PayloadStatus::REMOTE_ERROR;
return PayloadStatus::REMOTE_ERROR;
case PayloadTransferFrame::ControlMessage::PAYLOAD_CANCELED:
return location::nearby::proto::connections::PayloadStatus::
REMOTE_CANCELLATION;
return PayloadStatus::REMOTE_CANCELLATION;
default:
LOG(INFO) << "PayloadManager: unknown event=" << event;
return location::nearby::proto::connections::PayloadStatus::
UNKNOWN_PAYLOAD_STATUS;
return PayloadStatus::UNKNOWN_PAYLOAD_STATUS;
}
}
OperationResultCode PayloadManager::ControlMessageEventToOperationResultCode(
PayloadTransferFrame::ControlMessage::EventType event) {
switch (event) {
case PayloadTransferFrame::ControlMessage::PAYLOAD_ERROR:
return OperationResultCode::NEARBY_GENERIC_REMOTE_REPORT_PAYLOADS_ERROR;
case PayloadTransferFrame::ControlMessage::PAYLOAD_CANCELED:
return OperationResultCode::CLIENT_CANCELLATION_REMOTE_CANCEL_PAYLOAD;
default:
LOG(INFO) << "PayloadManager: unknown event=" << event;
return OperationResultCode::DETAIL_UNKNOWN;
}
}
PayloadProgressInfo::Status PayloadManager::PayloadStatusToTransferUpdateStatus(
location::nearby::proto::connections::PayloadStatus status) {
PayloadStatus status) {
switch (status) {
case location::nearby::proto::connections::LOCAL_CANCELLATION:
case location::nearby::proto::connections::REMOTE_CANCELLATION:
case PayloadStatus::LOCAL_CANCELLATION:
case PayloadStatus::REMOTE_CANCELLATION:
return PayloadProgressInfo::Status::kCanceled;
case location::nearby::proto::connections::SUCCESS:
case PayloadStatus::SUCCESS:
return PayloadProgressInfo::Status::kSuccess;
default:
return PayloadProgressInfo::Status::kFailure;
@@ -747,12 +783,14 @@ PayloadTransferFrame::PayloadChunk PayloadManager::CreatePayloadChunk(
return payload_chunk;
}
PayloadManager::PendingPayloadHandle PayloadManager::CreateIncomingPayload(
const PayloadTransferFrame& frame, const std::string& endpoint_id) {
std::pair<PayloadManager::PendingPayloadHandle, OperationResultCode>
PayloadManager::CreateIncomingPayload(const PayloadTransferFrame& frame,
const std::string& endpoint_id) {
// TODO(edwinwu): Add for return result code
auto internal_payload =
CreateIncomingInternalPayload(frame, custom_save_path_);
if (!internal_payload) {
return PendingPayloadHandle();
return {PendingPayloadHandle(), OperationResultCode::DETAIL_UNKNOWN};
}
Payload::Id payload_id = internal_payload->GetId();
@@ -762,7 +800,8 @@ PayloadManager::PendingPayloadHandle PayloadManager::CreateIncomingPayload(
std::make_unique<PendingPayload>(
std::move(internal_payload), EndpointIds{endpoint_id}, true,
absl::bind_front(&PayloadManager::OnPendingPayloadDestroy, this)));
return pending_payloads_.GetPayload(payload_id);
return {pending_payloads_.GetPayload(payload_id),
OperationResultCode::DETAIL_UNKNOWN};
}
void PayloadManager::OnPendingPayloadDestroy(const PendingPayload* payload) {
@@ -781,13 +820,13 @@ void PayloadManager::OnPendingPayloadDestroy(const PendingPayload* payload) {
void PayloadManager::SendClientCallbacksForFinishedOutgoingPayload(
ClientProxy* client, const EndpointIds& finished_endpoint_ids,
const PayloadTransferFrame::PayloadHeader& payload_header,
std::int64_t num_bytes_successfully_transferred,
location::nearby::proto::connections::PayloadStatus status) {
std::int64_t num_bytes_successfully_transferred, PayloadStatus status,
OperationResultCode operation_result_code) {
RunOnStatusUpdateThread(
"outgoing-payload-callbacks",
[this, client, finished_endpoint_ids, payload_header,
num_bytes_successfully_transferred,
status]() RUN_ON_PAYLOAD_STATUS_UPDATE_THREAD() {
num_bytes_successfully_transferred, status,
operation_result_code]() RUN_ON_PAYLOAD_STATUS_UPDATE_THREAD() {
// Make sure we're still tracking this payload.
PendingPayloadHandle pending_payload = GetPayload(payload_header.id());
if (!pending_payload) {
@@ -808,11 +847,14 @@ void PayloadManager::SendClientCallbacksForFinishedOutgoingPayload(
// Notify the client.
client->OnPayloadProgress(endpoint_id, update);
// TODO(edwinwu): Add for return result code
// Mark this payload as done for analytics.
client->GetAnalyticsRecorder().OnOutgoingPayloadDone(
endpoint_id, payload_header.id(), status,
location::nearby::proto::connections::OperationResultCode::
DETAIL_UNKNOWN);
(operation_result_code == OperationResultCode::DETAIL_UNKNOWN &&
status == PayloadStatus::ENDPOINT_IO_ERROR)
? OperationResultCode::DETAIL_UNKNOWN
: operation_result_code);
}
// Remove these endpoints from our tracking list for this payload.
@@ -828,12 +870,12 @@ void PayloadManager::SendClientCallbacksForFinishedOutgoingPayload(
void PayloadManager::SendClientCallbacksForFinishedIncomingPayload(
ClientProxy* client, const std::string& endpoint_id,
const PayloadTransferFrame::PayloadHeader& payload_header,
std::int64_t offset_bytes,
location::nearby::proto::connections::PayloadStatus status) {
std::int64_t offset_bytes, PayloadStatus status,
OperationResultCode operation_result_code) {
RunOnStatusUpdateThread(
"incoming-payload-callbacks",
[this, client, endpoint_id, payload_header, offset_bytes,
status]() RUN_ON_PAYLOAD_STATUS_UPDATE_THREAD() {
[this, client, endpoint_id, payload_header, offset_bytes, status,
operation_result_code]() RUN_ON_PAYLOAD_STATUS_UPDATE_THREAD() {
// Make sure we're still tracking this payload.
PendingPayloadHandle pending_payload = GetPayload(payload_header.id());
if (!pending_payload) {
@@ -852,9 +894,7 @@ void PayloadManager::SendClientCallbacksForFinishedIncomingPayload(
// Analyze
client->GetAnalyticsRecorder().OnIncomingPayloadDone(
endpoint_id, payload_header.id(), status,
location::nearby::proto::connections::OperationResultCode::
DETAIL_UNKNOWN);
endpoint_id, payload_header.id(), status, operation_result_code);
});
}
@@ -927,10 +967,10 @@ bool PayloadManager::WaitForReceivedAck(
// Local payload cancellation
if (latest_pending_payload->IsLocallyCanceled()) {
HandleFinishedOutgoingPayload(client, {endpoint_id}, payload_header,
payload_chunk_offset,
location::nearby::proto::connections::
PayloadStatus::LOCAL_CANCELLATION);
HandleFinishedOutgoingPayload(
client, {endpoint_id}, payload_header, payload_chunk_offset,
OperationResultCode::CLIENT_CANCELLATION_LOCAL_CANCEL_PAYLOAD,
PayloadStatus::LOCAL_CANCELLATION);
LOG(INFO) << "[safe-to-disconnect] short-circuiting local "
"payload cancellation for "
<< payload_header.id() << ", stop wait ack.";
@@ -941,6 +981,7 @@ bool PayloadManager::WaitForReceivedAck(
endpoint_info->status.Get())) {
HandleFinishedOutgoingPayload(
client, {endpoint_id}, payload_header, payload_chunk_offset,
OperationResultCode::CLIENT_CANCELLATION_REMOTE_CANCEL_PAYLOAD,
EndpointInfoStatusToPayloadStatus(endpoint_info->status.Get()));
LOG(INFO) << "[safe-to-disconnect] short-circuiting remote "
"payload cancellation for "
@@ -1002,20 +1043,19 @@ void PayloadManager::HandleFinishedOutgoingPayload(
ClientProxy* client, const EndpointIds& finished_endpoint_ids,
const PayloadTransferFrame::PayloadHeader& payload_header,
std::int64_t num_bytes_successfully_transferred,
location::nearby::proto::connections::PayloadStatus status) {
OperationResultCode operation_result_code, PayloadStatus status) {
// This call will destroy a pending payload.
SendClientCallbacksForFinishedOutgoingPayload(
client, finished_endpoint_ids, payload_header,
num_bytes_successfully_transferred, status);
num_bytes_successfully_transferred, status, operation_result_code);
switch (status) {
case location::nearby::proto::connections::PayloadStatus::LOCAL_ERROR:
case PayloadStatus::LOCAL_ERROR:
SendControlMessage(finished_endpoint_ids, payload_header,
num_bytes_successfully_transferred,
PayloadTransferFrame::ControlMessage::PAYLOAD_ERROR);
break;
case location::nearby::proto::connections::PayloadStatus::
LOCAL_CANCELLATION:
case PayloadStatus::LOCAL_CANCELLATION:
LOG(INFO) << "Sending PAYLOAD_CANCEL to receiver side; payload_id="
<< payload_header.id();
SendControlMessage(
@@ -1023,7 +1063,7 @@ void PayloadManager::HandleFinishedOutgoingPayload(
num_bytes_successfully_transferred,
PayloadTransferFrame::ControlMessage::PAYLOAD_CANCELED);
break;
case location::nearby::proto::connections::PayloadStatus::ENDPOINT_IO_ERROR:
case PayloadStatus::ENDPOINT_IO_ERROR:
// Unregister these endpoints, since we had an IO error on the physical
// connection.
for (const auto& endpoint_id : finished_endpoint_ids) {
@@ -1031,9 +1071,8 @@ void PayloadManager::HandleFinishedOutgoingPayload(
DisconnectionReason::IO_ERROR);
}
break;
case location::nearby::proto::connections::PayloadStatus::REMOTE_ERROR:
case location::nearby::proto::connections::PayloadStatus::
REMOTE_CANCELLATION:
case PayloadStatus::REMOTE_ERROR:
case PayloadStatus::REMOTE_CANCELLATION:
// No special handling needed for these.
break;
default:
@@ -1047,18 +1086,18 @@ void PayloadManager::HandleFinishedOutgoingPayload(
void PayloadManager::HandleFinishedIncomingPayload(
ClientProxy* client, const std::string& endpoint_id,
const PayloadTransferFrame::PayloadHeader& payload_header,
std::int64_t offset_bytes,
location::nearby::proto::connections::PayloadStatus status) {
SendClientCallbacksForFinishedIncomingPayload(
client, endpoint_id, payload_header, offset_bytes, status);
std::int64_t offset_bytes, PayloadStatus status,
OperationResultCode operation_result_code) {
SendClientCallbacksForFinishedIncomingPayload(client, endpoint_id,
payload_header, offset_bytes,
status, operation_result_code);
switch (status) {
case location::nearby::proto::connections::PayloadStatus::LOCAL_ERROR:
case PayloadStatus::LOCAL_ERROR:
SendControlMessage({endpoint_id}, payload_header, offset_bytes,
PayloadTransferFrame::ControlMessage::PAYLOAD_ERROR);
break;
case location::nearby::proto::connections::PayloadStatus::
LOCAL_CANCELLATION:
case PayloadStatus::LOCAL_CANCELLATION:
SendControlMessage(
{endpoint_id}, payload_header, offset_bytes,
PayloadTransferFrame::ControlMessage::PAYLOAD_CANCELED);
@@ -1147,10 +1186,8 @@ void PayloadManager::HandleSuccessfulOutgoingChunk(
if (is_last_chunk) {
client->GetAnalyticsRecorder().OnOutgoingPayloadDone(
endpoint_id, payload_header.id(),
location::nearby::proto::connections::SUCCESS,
location::nearby::proto::connections::OperationResultCode::
DETAIL_UNKNOWN);
endpoint_id, payload_header.id(), PayloadStatus::SUCCESS,
OperationResultCode::DETAIL_SUCCESS);
// Stop tracking this endpoint.
pending_payload->RemoveEndpoints({endpoint_id});
@@ -1246,10 +1283,8 @@ void PayloadManager::HandleSuccessfulIncomingChunk(
if (is_last_chunk) {
DestroyPendingPayload(payload_header.id());
client->GetAnalyticsRecorder().OnIncomingPayloadDone(
endpoint_id, payload_header.id(),
location::nearby::proto::connections::SUCCESS,
location::nearby::proto::connections::OperationResultCode::
DETAIL_UNKNOWN);
endpoint_id, payload_header.id(), PayloadStatus::SUCCESS,
OperationResultCode::DETAIL_SUCCESS);
} else {
client->GetAnalyticsRecorder().OnPayloadChunkReceived(
endpoint_id, payload_header.id(), payload_chunk_body_size);
@@ -1299,13 +1334,26 @@ void PayloadManager::ProcessDataPacket(
payload_header.total_size());
});
pending_payload =
std::pair<PendingPayloadHandle, OperationResultCode> result =
CreateIncomingPayload(payload_transfer_frame, from_endpoint_id);
pending_payload = std::move(result.first);
OperationResultCode operation_result_code = result.second;
if (!pending_payload) {
LOG(WARNING) << "PayloadManager failed to create InternalPayload from "
"PayloadTransferFrame with payload_id="
<< payload_header.id() << " and type "
<< payload_header.type() << ", aborting receipt.";
// Analyticize.
RunOnStatusUpdateThread(
"process-data-packet",
[to_client, from_endpoint_id, payload_header, operation_result_code]()
RUN_ON_PAYLOAD_STATUS_UPDATE_THREAD() {
to_client->GetAnalyticsRecorder().OnIncomingPayloadDone(
from_endpoint_id, payload_header.id(),
PayloadStatus::LOCAL_ERROR, operation_result_code);
});
// Send the error to the remote endpoint.
SendControlMessage({from_endpoint_id}, payload_header,
payload_chunk.offset(),
@@ -1340,10 +1388,10 @@ void PayloadManager::ProcessDataPacket(
// do all the cleanup. See go/nc-cancel-payload
LOG(INFO) << "ProcessDataPacket: [cancel] endpoint_id=" << from_endpoint_id
<< "; payload_id=" << pending_payload->GetId();
HandleFinishedIncomingPayload(to_client, from_endpoint_id, payload_header,
payload_chunk.offset(),
location::nearby::proto::connections::
PayloadStatus::LOCAL_CANCELLATION);
HandleFinishedIncomingPayload(
to_client, from_endpoint_id, payload_header, payload_chunk.offset(),
PayloadStatus::LOCAL_CANCELLATION,
OperationResultCode::CLIENT_CANCELLATION_LOCAL_CANCEL_PAYLOAD);
return;
}
@@ -1367,7 +1415,7 @@ void PayloadManager::ProcessDataPacket(
<< "; payload_id=" << pending_payload->GetId();
HandleFinishedIncomingPayload(
to_client, from_endpoint_id, payload_header, payload_chunk.offset(),
location::nearby::proto::connections::PayloadStatus::LOCAL_ERROR);
PayloadStatus::LOCAL_ERROR, OperationResultCode::IO_FILE_WRITING_ERROR);
return;
}
packet_meta_data.StopFileIo();
@@ -1417,7 +1465,8 @@ void PayloadManager::ProcessControlPacket(
HandleFinishedIncomingPayload(
to_client, from_endpoint_id, payload_header,
control_message.offset(),
ControlMessageEventToPayloadStatus(control_message.event()));
ControlMessageEventToPayloadStatus(control_message.event()),
ControlMessageEventToOperationResultCode(control_message.event()));
} else {
LOG(INFO) << "Outgoing PAYLOAD_CANCELED: from endpoint_id="
<< from_endpoint_id << "; self=" << this;
@@ -1436,7 +1485,8 @@ void PayloadManager::ProcessControlPacket(
HandleFinishedIncomingPayload(
to_client, from_endpoint_id, payload_header,
control_message.offset(),
ControlMessageEventToPayloadStatus(control_message.event()));
ControlMessageEventToPayloadStatus(control_message.event()),
ControlMessageEventToOperationResultCode(control_message.event()));
} else {
pending_payload->SetEndpointStatusFromControlMessage(from_endpoint_id,
control_message);
@@ -1493,16 +1543,14 @@ void PayloadManager::RecordPayloadStartedAnalytics(
void PayloadManager::RecordInvalidPayloadAnalytics(
ClientProxy* client, const EndpointIds& endpoint_ids,
std::int64_t payload_id, PayloadType payload_type, std::int64_t offset,
std::int64_t total_size) {
std::int64_t total_size, OperationResultCode operation_result_code) {
RecordPayloadStartedAnalytics(client, endpoint_ids, payload_id, payload_type,
offset, total_size);
for (const auto& endpoint_id : endpoint_ids) {
client->GetAnalyticsRecorder().OnOutgoingPayloadDone(
endpoint_id, payload_id,
location::nearby::proto::connections::LOCAL_ERROR,
location::nearby::proto::connections::OperationResultCode::
DETAIL_UNKNOWN);
endpoint_id, payload_id, PayloadStatus::LOCAL_ERROR,
operation_result_code);
}
}
@@ -1680,8 +1728,7 @@ void PayloadManager::RunOnStatusUpdateThread(
payload_status_update_executor_.Execute(name, std::move(runnable));
}
/////////////////////////////// PendingPayloads
//////////////////////////////////
/////////////////////////////// PendingPayloads ////////////////////////////////
void PayloadManager::PendingPayloads::StartTrackingPayload(
Payload::Id payload_id, std::unique_ptr<PendingPayload> pending_payload) {
+92 -54
View File
@@ -17,7 +17,6 @@
#include <cstddef>
#include <cstdint>
#include <functional>
#include <memory>
#include <string>
#include <utility>
@@ -45,7 +44,6 @@
namespace nearby {
namespace connections {
using ::location::nearby::connections::PayloadTransferFrame;
// Annotations for methods that need to run on PayloadStatusUpdateThread.
// Use only in PayloadManager
@@ -55,8 +53,7 @@ using ::location::nearby::connections::PayloadTransferFrame;
class PayloadManager : public EndpointManager::FrameProcessor {
public:
using EndpointIds = std::vector<std::string>;
constexpr static const absl::Duration kWaitCloseTimeout =
absl::Milliseconds(5000);
static constexpr absl::Duration kWaitCloseTimeout = absl::Milliseconds(5000);
explicit PayloadManager(EndpointManager& endpoint_manager);
~PayloadManager() override;
@@ -73,10 +70,11 @@ class PayloadManager : public EndpointManager::FrameProcessor {
analytics::PacketMetaData& packet_meta_data) override;
// @EndpointManagerThread
void OnEndpointDisconnect(ClientProxy* client, const std::string& service_id,
const std::string& endpoint_id,
CountDownLatch barrier,
DisconnectionReason reason) override;
void OnEndpointDisconnect(
ClientProxy* client, const std::string& service_id,
const std::string& endpoint_id, CountDownLatch barrier,
location::nearby::proto::connections::DisconnectionReason reason)
override;
void DisconnectFromEndpointManager();
@@ -94,10 +92,12 @@ class PayloadManager : public EndpointManager::FrameProcessor {
};
void SetStatusFromControlMessage(
const PayloadTransferFrame::ControlMessage& control_message);
const location::nearby::connections::PayloadTransferFrame::
ControlMessage& control_message);
static Status ControlMessageEventToEndpointInfoStatus(
PayloadTransferFrame::ControlMessage::EventType event);
location::nearby::connections::PayloadTransferFrame::ControlMessage::
EventType event);
void MarkReceivedAckFromEndpoint();
bool IsEndpointAvailable(ClientProxy* clientProxy,
EndpointInfo::Status status);
@@ -153,8 +153,8 @@ class PayloadManager : public EndpointManager::FrameProcessor {
// Sets the status for a particular endpoint.
void SetEndpointStatusFromControlMessage(
const std::string& endpoint_id,
const PayloadTransferFrame::ControlMessage& control_message)
ABSL_LOCKS_EXCLUDED(mutex_);
const location::nearby::connections::PayloadTransferFrame::
ControlMessage& control_message) ABSL_LOCKS_EXCLUDED(mutex_);
// Sets the offset for a particular endpoint.
void SetOffsetForEndpoint(const std::string& endpoint_id,
@@ -272,47 +272,63 @@ class PayloadManager : public EndpointManager::FrameProcessor {
// Returns list of endpoint ids.
static EndpointIds EndpointsToEndpointIds(const Endpoints& endpoints);
bool SendPayloadLoop(ClientProxy* client, PendingPayload& pending_payload,
PayloadTransferFrame::PayloadHeader& payload_header,
std::int64_t& next_chunk_offset, size_t resume_offset,
int index);
bool SendPayloadLoop(
ClientProxy* client, PendingPayload& pending_payload,
location::nearby::connections::PayloadTransferFrame::PayloadHeader&
payload_header,
std::int64_t& next_chunk_offset, size_t resume_offset, int index);
void SendClientCallbacksForFinishedIncomingPayloadRunnable(
ClientProxy* client, const std::string& endpoint_id,
const PayloadTransferFrame::PayloadHeader& payload_header,
const location::nearby::connections::PayloadTransferFrame::PayloadHeader&
payload_header,
std::int64_t offset_bytes,
location::nearby::proto::connections::PayloadStatus status);
location::nearby::proto::connections::PayloadStatus status,
location::nearby::proto::connections::OperationResultCode
operation_result_code);
// Converts the status of an endpoint that's been set out-of-band via a
// remote ControlMessage to the PayloadStatus for handling of that
// endpoint-payload pair.
static location::nearby::proto::connections::PayloadStatus
EndpointInfoStatusToPayloadStatus(EndpointInfo::Status status);
static location::nearby::proto::connections::OperationResultCode
EndpointInfoStatusToOperationResultCode(EndpointInfo::Status status);
// Converts a ControlMessage::EventType for a particular payload to a
// PayloadStatus. Called when we've received a ControlMessage with this
// event from a remote endpoint; thus the PayloadStatuses are REMOTE_*.
static location::nearby::proto::connections::PayloadStatus
ControlMessageEventToPayloadStatus(
PayloadTransferFrame::ControlMessage::EventType event);
location::nearby::connections::PayloadTransferFrame::ControlMessage::
EventType event);
static location::nearby::proto::connections::OperationResultCode
ControlMessageEventToOperationResultCode(
location::nearby::connections::PayloadTransferFrame::ControlMessage::
EventType event);
static PayloadProgressInfo::Status PayloadStatusToTransferUpdateStatus(
location::nearby::proto::connections::PayloadStatus status);
int GetOptimalChunkSize(EndpointIds endpoint_ids);
PayloadTransferFrame::PayloadHeader CreatePayloadHeader(
const InternalPayload& internal_payload, size_t offset,
const std::string& parent_folder, const std::string& file_name);
location::nearby::connections::PayloadTransferFrame::PayloadHeader
CreatePayloadHeader(const InternalPayload& internal_payload, size_t offset,
const std::string& parent_folder,
const std::string& file_name);
PayloadTransferFrame::PayloadChunk CreatePayloadChunk(std::int64_t offset,
ByteArray body,
int index);
bool IsLastChunk(PayloadTransferFrame::PayloadChunk payload_chunk) {
location::nearby::connections::PayloadTransferFrame::PayloadChunk
CreatePayloadChunk(std::int64_t offset, ByteArray body, int index);
bool IsLastChunk(
location::nearby::connections::PayloadTransferFrame::PayloadChunk
payload_chunk) {
return ((payload_chunk.flags() &
PayloadTransferFrame::PayloadChunk::LAST_CHUNK) != 0);
location::nearby::connections::PayloadTransferFrame::PayloadChunk::
LAST_CHUNK) != 0);
}
PendingPayloadHandle CreateIncomingPayload(const PayloadTransferFrame& frame,
const std::string& endpoint_id)
ABSL_LOCKS_EXCLUDED(mutex_);
std::pair<PendingPayloadHandle,
location::nearby::proto::connections::OperationResultCode>
CreateIncomingPayload(
const location::nearby::connections::PayloadTransferFrame& frame,
const std::string& endpoint_id) ABSL_LOCKS_EXCLUDED(mutex_);
Payload::Id CreateOutgoingPayload(Payload payload,
const EndpointIds& endpoint_ids)
@@ -320,20 +336,28 @@ class PayloadManager : public EndpointManager::FrameProcessor {
void SendClientCallbacksForFinishedOutgoingPayload(
ClientProxy* client, const EndpointIds& finished_endpoint_ids,
const PayloadTransferFrame::PayloadHeader& payload_header,
const location::nearby::connections::PayloadTransferFrame::PayloadHeader&
payload_header,
std::int64_t num_bytes_successfully_transferred,
location::nearby::proto::connections::PayloadStatus status);
location::nearby::proto::connections::PayloadStatus status,
location::nearby::proto::connections::OperationResultCode
operation_result_code);
void SendClientCallbacksForFinishedIncomingPayload(
ClientProxy* client, const std::string& endpoint_id,
const PayloadTransferFrame::PayloadHeader& payload_header,
const location::nearby::connections::PayloadTransferFrame::PayloadHeader&
payload_header,
std::int64_t offset_bytes,
location::nearby::proto::connections::PayloadStatus status);
location::nearby::proto::connections::PayloadStatus status,
location::nearby::proto::connections::OperationResultCode
operation_result_code);
void SendControlMessage(
const EndpointIds& endpoint_ids,
const PayloadTransferFrame::PayloadHeader& payload_header,
const location::nearby::connections::PayloadTransferFrame::PayloadHeader&
payload_header,
std::int64_t num_bytes_successfully_transferred,
PayloadTransferFrame::ControlMessage::EventType event_type);
location::nearby::connections::PayloadTransferFrame::ControlMessage::
EventType event_type);
void SendPayloadReceivedAck(ClientProxy* client,
PendingPayload& pending_payload,
@@ -343,7 +367,8 @@ class PayloadManager : public EndpointManager::FrameProcessor {
bool WaitForReceivedAck(
ClientProxy* client, const std::string& endpoint_id,
PendingPayload& pending_payload,
const PayloadTransferFrame::PayloadHeader& payload_header,
const location::nearby::connections::PayloadTransferFrame::PayloadHeader&
payload_header,
std::int64_t payload_chunk_offset, bool is_last_chunk);
bool IsPayloadReceivedAckEnabled(ClientProxy* client,
const std::string& endpoint_id,
@@ -353,37 +378,49 @@ class PayloadManager : public EndpointManager::FrameProcessor {
// statuses except for SUCCESS are handled here.
void HandleFinishedOutgoingPayload(
ClientProxy* client, const EndpointIds& finished_endpoint_ids,
const PayloadTransferFrame::PayloadHeader& payload_header,
const location::nearby::connections::PayloadTransferFrame::PayloadHeader&
payload_header,
std::int64_t num_bytes_successfully_transferred,
location::nearby::proto::connections::OperationResultCode
operation_result_code,
location::nearby::proto::connections::PayloadStatus status = location::
nearby::proto::connections::PayloadStatus::UNKNOWN_PAYLOAD_STATUS);
void HandleFinishedIncomingPayload(
ClientProxy* client, const std::string& endpoint_id,
const PayloadTransferFrame::PayloadHeader& payload_header,
const location::nearby::connections::PayloadTransferFrame::PayloadHeader&
payload_header,
std::int64_t offset_bytes,
location::nearby::proto::connections::PayloadStatus status);
location::nearby::proto::connections::PayloadStatus status,
location::nearby::proto::connections::OperationResultCode
operation_result_code);
void HandleSuccessfulOutgoingChunk(
ClientProxy* client, const std::string& endpoint_id,
const PayloadTransferFrame::PayloadHeader& payload_header,
const location::nearby::connections::PayloadTransferFrame::PayloadHeader&
payload_header,
std::int32_t payload_chunk_flags, std::int64_t payload_chunk_offset,
std::int64_t payload_chunk_body_size);
void HandleSuccessfulIncomingChunk(
ClientProxy* client, const std::string& endpoint_id,
const PayloadTransferFrame::PayloadHeader& payload_header,
const location::nearby::connections::PayloadTransferFrame::PayloadHeader&
payload_header,
std::int32_t payload_chunk_flags, std::int64_t payload_chunk_offset,
std::int64_t payload_chunk_body_size);
void ProcessDataPacket(ClientProxy* to_client,
const std::string& from_endpoint_id,
PayloadTransferFrame& payload_transfer_frame,
Medium medium,
location::nearby::connections::PayloadTransferFrame&
payload_transfer_frame,
location::nearby::proto::connections::Medium medium,
analytics::PacketMetaData& packet_meta_data);
void ProcessControlPacket(ClientProxy* to_client,
const std::string& from_endpoint_id,
PayloadTransferFrame& payload_transfer_frame);
void ProcessPayloadAckPacket(const std::string& from_endpoint_id,
PayloadTransferFrame& payload_transfer_frame);
location::nearby::connections::PayloadTransferFrame&
payload_transfer_frame);
void ProcessPayloadAckPacket(
const std::string& from_endpoint_id,
location::nearby::connections::PayloadTransferFrame&
payload_transfer_frame);
void NotifyClientOfIncomingPayloadProgressInfo(
ClientProxy* client, const std::string& endpoint_id,
@@ -407,15 +444,16 @@ class PayloadManager : public EndpointManager::FrameProcessor {
PayloadType payload_type,
std::int64_t offset,
std::int64_t total_size);
void RecordInvalidPayloadAnalytics(ClientProxy* client,
const EndpointIds& endpoint_ids,
std::int64_t payload_id,
PayloadType payload_type,
std::int64_t offset,
std::int64_t total_size);
void RecordInvalidPayloadAnalytics(
ClientProxy* client, const EndpointIds& endpoint_ids,
std::int64_t payload_id, PayloadType payload_type, std::int64_t offset,
std::int64_t total_size,
location::nearby::proto::connections::OperationResultCode
operation_result_code);
PayloadType FramePayloadTypeToPayloadType(
PayloadTransferFrame::PayloadHeader::PayloadType type);
location::nearby::connections::PayloadTransferFrame::PayloadHeader::
PayloadType type);
void OnPendingPayloadDestroy(const PendingPayload* payload);
mutable Mutex mutex_;
@@ -22,17 +22,16 @@
#include "absl/strings/string_view.h"
#include "absl/time/time.h"
#include "connections/implementation/analytics/packet_meta_data.h"
#include "connections/implementation/flags/nearby_connections_feature_flags.h"
#include "connections/implementation/offline_frames.h"
#include "connections/implementation/simulation_user.h"
#include "connections/listeners.h"
#include "connections/medium_selector.h"
#include "connections/payload.h"
#include "connections/status.h"
#include "internal/flags/nearby_flags.h"
#include "internal/platform/byte_array.h"
#include "internal/platform/count_down_latch.h"
#include "internal/platform/exception.h"
#include "internal/platform/input_stream.h"
#include "internal/platform/logging.h"
#include "internal/platform/medium_environment.h"
#include "internal/platform/pipe.h"
@@ -41,6 +40,7 @@ namespace nearby {
namespace connections {
namespace {
using ::location::nearby::connections::OfflineFrame;
using ::location::nearby::connections::PayloadTransferFrame;
using ::nearby::analytics::PacketMetaData;
using ::location::nearby::proto::connections::Medium;