analytics: Add operation result code for Payload analytics, I

PiperOrigin-RevId: 704248204
This commit is contained in:
Edwin Wu
2024-12-09 05:35:06 -08:00
committed by Copybara-Service
parent 438267a63a
commit c5fed0c0ed
4 changed files with 141 additions and 29 deletions
@@ -654,9 +654,9 @@ void AnalyticsRecorder::OnPayloadChunkReceived(const std::string &endpoint_id,
logical_connection->ChunkReceived(payload_id, chunk_size_bytes);
}
void AnalyticsRecorder::OnIncomingPayloadDone(const std::string &endpoint_id,
std::int64_t payload_id,
PayloadStatus status) {
void AnalyticsRecorder::OnIncomingPayloadDone(
const std::string &endpoint_id, std::int64_t payload_id,
PayloadStatus status, OperationResultCode operation_result_code) {
MutexLock lock(&mutex_);
if (!CanRecordAnalyticsLocked("OnIncomingPayloadDone")) {
return;
@@ -666,7 +666,8 @@ void AnalyticsRecorder::OnIncomingPayloadDone(const std::string &endpoint_id,
return;
}
const std::unique_ptr<LogicalConnection> &logical_connection = it->second;
logical_connection->IncomingPayloadDone(payload_id, status);
logical_connection->IncomingPayloadDone(payload_id, status,
operation_result_code);
}
void AnalyticsRecorder::OnOutgoingPayloadStarted(
@@ -702,9 +703,9 @@ void AnalyticsRecorder::OnPayloadChunkSent(const std::string &endpoint_id,
logical_connection->ChunkSent(payload_id, chunk_size_bytes);
}
void AnalyticsRecorder::OnOutgoingPayloadDone(const std::string &endpoint_id,
std::int64_t payload_id,
PayloadStatus status) {
void AnalyticsRecorder::OnOutgoingPayloadDone(
const std::string &endpoint_id, std::int64_t payload_id,
PayloadStatus status, OperationResultCode operation_result_code) {
MutexLock lock(&mutex_);
if (!CanRecordAnalyticsLocked("OnOutgoingPayloadDone")) {
return;
@@ -713,8 +714,10 @@ void AnalyticsRecorder::OnOutgoingPayloadDone(const std::string &endpoint_id,
if (it == active_connections_.end()) {
return;
}
const std::unique_ptr<LogicalConnection> &logical_connection = it->second;
logical_connection->OutgoingPayloadDone(payload_id, status);
logical_connection->OutgoingPayloadDone(payload_id, status,
operation_result_code);
}
void AnalyticsRecorder::OnBandwidthUpgradeStarted(
@@ -1310,6 +1313,13 @@ ConnectionsLog::Payload AnalyticsRecorder::PendingPayload::GetProtoPayload(
payload.set_num_chunks(num_chunks_);
payload.set_status(status);
auto operation_result_proto =
std::make_unique<ConnectionsLog::OperationResult>();
operation_result_proto->set_result_code(operation_result_code_);
operation_result_proto->set_result_category(
GetOperationResultCateory(operation_result_code_));
payload.set_allocated_operation_result(operation_result_proto.release());
return payload;
}
@@ -1327,6 +1337,14 @@ void AnalyticsRecorder::LogicalConnection::PhysicalConnectionEstablished(
established_connection->set_duration_millis(
absl::ToUnixMillis(SystemClock::ElapsedRealtime()));
established_connection->set_connection_token(connection_token);
auto operation_result_proto =
std::make_unique<ConnectionsLog::OperationResult>();
operation_result_proto->set_result_code(OperationResultCode::DETAIL_SUCCESS);
operation_result_proto->set_result_category(
OperationResultCategory::CATEGORY_SUCCESS);
established_connection->set_allocated_operation_result(
operation_result_proto.release());
physical_connections_.insert({medium, std::move(established_connection)});
current_medium_ = medium;
}
@@ -1428,7 +1446,8 @@ void AnalyticsRecorder::LogicalConnection::ChunkReceived(
}
void AnalyticsRecorder::LogicalConnection::IncomingPayloadDone(
std::int64_t payload_id, PayloadStatus status) {
std::int64_t payload_id, PayloadStatus status,
OperationResultCode operation_result_code) {
if (current_medium_ == UNKNOWN_MEDIUM) {
NEARBY_LOGS(WARNING) << "Unexpected call to incomingPayloadDone() while "
"AnalyticsRecorder has no active current medium.";
@@ -1440,6 +1459,7 @@ void AnalyticsRecorder::LogicalConnection::IncomingPayloadDone(
&established_connection = it->second;
auto it = incoming_payloads_.find(payload_id);
if (it != incoming_payloads_.end()) {
it->second->SetOperationResultCode(operation_result_code);
*established_connection->add_received_payload() =
it->second->GetProtoPayload(status);
incoming_payloads_.erase(it);
@@ -1464,7 +1484,8 @@ void AnalyticsRecorder::LogicalConnection::ChunkSent(std::int64_t payload_id,
}
void AnalyticsRecorder::LogicalConnection::OutgoingPayloadDone(
std::int64_t payload_id, PayloadStatus status) {
std::int64_t payload_id, PayloadStatus status,
OperationResultCode operation_result_code) {
if (current_medium_ == UNKNOWN_MEDIUM) {
NEARBY_LOGS(WARNING) << "Unexpected call to outgoingPayloadDone() while "
"AnalyticsRecorder has no active current medium.";
@@ -1476,6 +1497,7 @@ void AnalyticsRecorder::LogicalConnection::OutgoingPayloadDone(
&established_connection = it->second;
auto it = outgoing_payloads_.find(payload_id);
if (it != outgoing_payloads_.end()) {
it->second->SetOperationResultCode(operation_result_code);
*established_connection->add_sent_payload() =
it->second->GetProtoPayload(status);
outgoing_payloads_.erase(it);
@@ -1515,16 +1537,21 @@ AnalyticsRecorder::LogicalConnection::ResolvePendingPayloads(
upgraded_payloads;
PayloadStatus status =
reason == UPGRADED ? MOVED_TO_NEW_MEDIUM : CONNECTION_CLOSED;
OperationResultCode operation_result_code =
GetPendingPayloadResultCodeFromReason(reason);
for (const auto &item : pending_payloads) {
const std::unique_ptr<PendingPayload> &pending_payload = item.second;
pending_payload->SetOperationResultCode(operation_result_code);
ConnectionsLog::Payload proto_payload =
pending_payload->GetProtoPayload(status);
completed_payloads.push_back(proto_payload);
if (reason == UPGRADED) {
upgraded_payloads.insert(
{item.first,
std::make_unique<PendingPayload>(
pending_payload->type(), pending_payload->total_size_bytes())});
std::make_unique<PendingPayload>(pending_payload->type(),
pending_payload->total_size_bytes(),
operation_result_code)});
}
}
pending_payloads.clear();
@@ -1539,6 +1566,21 @@ AnalyticsRecorder::LogicalConnection::ResolvePendingPayloads(
return completed_payloads;
}
OperationResultCode
AnalyticsRecorder::LogicalConnection::GetPendingPayloadResultCodeFromReason(
DisconnectionReason reason) {
switch (reason) {
case UPGRADED:
return OperationResultCode::MISCELLEANEOUS_MOVE_TO_NEW_MEDIUM;
case DisconnectionReason::LOCAL_DISCONNECTION:
return OperationResultCode::CLIENT_CANCELLATION_LOCAL_DISCONNECT;
case DisconnectionReason::REMOTE_DISCONNECTION:
return OperationResultCode::CLIENT_CANCELLATION_REMOTE_DISCONNECT;
default:
return OperationResultCode::NEARBY_GENERIC_CONNECTION_CLOSED;
}
}
void AnalyticsRecorder::Sync() {
CountDownLatch latch(1);
serial_executor_.Execute([&]() { latch.CountDown(); });
@@ -154,8 +154,9 @@ class AnalyticsRecorder {
ABSL_LOCKS_EXCLUDED(mutex_);
void OnIncomingPayloadDone(
const std::string &endpoint_id, std::int64_t payload_id,
location::nearby::proto::connections::PayloadStatus status)
ABSL_LOCKS_EXCLUDED(mutex_);
location::nearby::proto::connections::PayloadStatus status,
location::nearby::proto::connections::OperationResultCode
operation_result_code) ABSL_LOCKS_EXCLUDED(mutex_);
void OnOutgoingPayloadStarted(const std::vector<std::string> &endpoint_ids,
std::int64_t payload_id,
connections::PayloadType type,
@@ -167,8 +168,9 @@ class AnalyticsRecorder {
ABSL_LOCKS_EXCLUDED(mutex_);
void OnOutgoingPayloadDone(
const std::string &endpoint_id, std::int64_t payload_id,
location::nearby::proto::connections::PayloadStatus status)
ABSL_LOCKS_EXCLUDED(mutex_);
location::nearby::proto::connections::PayloadStatus status,
location::nearby::proto::connections::OperationResultCode
operation_result_code) ABSL_LOCKS_EXCLUDED(mutex_);
// BandwidthUpgrade
void OnBandwidthUpgradeStarted(
@@ -212,11 +214,19 @@ class AnalyticsRecorder {
public:
PendingPayload(location::nearby::proto::connections::PayloadType type,
std::int64_t total_size_bytes)
: PendingPayload(type, total_size_bytes,
location::nearby::proto::connections::
OperationResultCode::DETAIL_UNKNOWN) {}
PendingPayload(location::nearby::proto::connections::PayloadType type,
std::int64_t total_size_bytes,
location::nearby::proto::connections::OperationResultCode
operation_result_code)
: start_time_(SystemClock::ElapsedRealtime()),
type_(type),
total_size_bytes_(total_size_bytes),
num_bytes_transferred_(0),
num_chunks_(0) {}
num_chunks_(0),
operation_result_code_(operation_result_code) {}
~PendingPayload() = default;
void AddChunk(std::int64_t chunk_size_bytes);
@@ -230,12 +240,21 @@ class AnalyticsRecorder {
std::int64_t total_size_bytes() const { return total_size_bytes_; }
void SetOperationResultCode(
location::nearby::proto::connections::OperationResultCode
operation_result_code) {
operation_result_code_ = operation_result_code;
}
private:
absl::Time start_time_;
location::nearby::proto::connections::PayloadType type_;
std::int64_t total_size_bytes_;
std::int64_t num_bytes_transferred_;
int num_chunks_;
location::nearby::proto::connections::OperationResultCode
operation_result_code_ = location::nearby::proto::connections::
OperationResultCode::DETAIL_UNKNOWN;
};
class LogicalConnection {
@@ -272,7 +291,9 @@ class AnalyticsRecorder {
void ChunkReceived(std::int64_t payload_id, std::int64_t size_bytes);
void IncomingPayloadDone(
std::int64_t payload_id,
location::nearby::proto::connections::PayloadStatus status);
location::nearby::proto::connections::PayloadStatus status,
location::nearby::proto::connections::OperationResultCode
operation_result_code);
void OutgoingPayloadStarted(
std::int64_t payload_id,
location::nearby::proto::connections::PayloadType type,
@@ -280,7 +301,9 @@ class AnalyticsRecorder {
void ChunkSent(std::int64_t payload_id, std::int64_t size_bytes);
void OutgoingPayloadDone(
std::int64_t payload_id,
location::nearby::proto::connections::PayloadStatus status);
location::nearby::proto::connections::PayloadStatus status,
location::nearby::proto::connections::OperationResultCode
operation_result_code);
std::vector<location::nearby::analytics::proto::ConnectionsLog::
EstablishedConnection>
@@ -298,6 +321,9 @@ class AnalyticsRecorder {
absl::btree_map<std::int64_t, std::unique_ptr<PendingPayload>>
&pending_payloads,
location::nearby::proto::connections::DisconnectionReason reason);
location::nearby::proto::connections::OperationResultCode
GetPendingPayloadResultCodeFromReason(
location::nearby::proto::connections::DisconnectionReason reason);
location::nearby::proto::connections::Medium current_medium_ =
location::nearby::proto::connections::UNKNOWN_MEDIUM;
@@ -780,11 +780,21 @@ TEST(AnalyticsRecorderTest, UnfinishedEstablishedConnectionsAddedAsUnfinished) {
medium: BLUETOOTH
disconnection_reason: UPGRADED
connection_token: "connection_token"
safe_disconnection_result: UNKNOWN_SAFE_DISCONNECTION_RESULT
operation_result <
result_category: CATEGORY_SUCCESS
result_code: DETAIL_SUCCESS
>
>
established_connection <
medium: WIFI_LAN
disconnection_reason: UNFINISHED
connection_token: "connection_token"
safe_disconnection_result: SAFE_DISCONNECTION
operation_result {
result_category: CATEGORY_SUCCESS
result_code: DETAIL_SUCCESS
}
>
>)pb");
@@ -817,7 +827,8 @@ TEST(AnalyticsRecorderTest, OutgoingPayloadUpgraded) {
analytics_recorder.OnPayloadChunkSent(endpoint_id, payload_id, 10);
analytics_recorder.OnPayloadChunkSent(endpoint_id, payload_id, 10);
analytics_recorder.OnPayloadChunkSent(endpoint_id, payload_id, 10);
analytics_recorder.OnOutgoingPayloadDone(endpoint_id, payload_id, SUCCESS);
analytics_recorder.OnOutgoingPayloadDone(endpoint_id, payload_id, SUCCESS,
OperationResultCode::DETAIL_SUCCESS);
analytics_recorder.OnConnectionClosed(
endpoint_id, WIFI_LAN, LOCAL_DISCONNECTION,
ConnectionsLog::EstablishedConnection::SAFE_DISCONNECTION);
@@ -847,9 +858,18 @@ TEST(AnalyticsRecorderTest, OutgoingPayloadUpgraded) {
num_bytes_transferred: 20
num_chunks: 2
status: MOVED_TO_NEW_MEDIUM
operation_result <
result_category: CATEGORY_MISCELLANEOUS
result_code: MISCELLEANEOUS_MOVE_TO_NEW_MEDIUM
>
>
disconnection_reason: UPGRADED
connection_token: "connection_token"
safe_disconnection_result: SAFE_DISCONNECTION
operation_result <
result_category: CATEGORY_SUCCESS
result_code: DETAIL_SUCCESS
>
>
established_connection <
medium: WIFI_LAN
@@ -859,9 +879,18 @@ TEST(AnalyticsRecorderTest, OutgoingPayloadUpgraded) {
num_bytes_transferred: 30
num_chunks: 3
status: SUCCESS
operation_result <
result_category: CATEGORY_SUCCESS
result_code: DETAIL_SUCCESS
>
>
disconnection_reason: LOCAL_DISCONNECTION
connection_token: "connection_token"
safe_disconnection_result: SAFE_DISCONNECTION
operation_result <
result_category: CATEGORY_SUCCESS
result_code: DETAIL_SUCCESS
>
>
>)pb");
@@ -1688,8 +1717,8 @@ TEST(AnalyticsRecorderTest,
analytics_recorder.LogSession();
ASSERT_TRUE(client_session_done_latch.Await(kDefaultTimeout).result());
//// The same strategy session shouldn't be logged again with the same client
//// session.
// The same strategy session shouldn't be logged again with the same client
// session.
EXPECT_THAT(event_logger.GetLoggedEventTypes(),
Contains(START_STRATEGY_SESSION).Times(1));
+22 -7
View File
@@ -67,6 +67,7 @@ constexpr absl::Duration kMinTransferUpdateInterval = absl::Milliseconds(50);
// 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,
@@ -622,10 +623,14 @@ void PayloadManager::OnEndpointDisconnect(ClientProxy* client,
if (pending_payload->IsIncoming()) {
client->GetAnalyticsRecorder().OnIncomingPayloadDone(
endpoint_id, pending_payload->GetId(), payload_status);
endpoint_id, pending_payload->GetId(), payload_status,
location::nearby::proto::connections::OperationResultCode::
DETAIL_UNKNOWN);
} else {
client->GetAnalyticsRecorder().OnOutgoingPayloadDone(
endpoint_id, pending_payload->GetId(), payload_status);
endpoint_id, pending_payload->GetId(), payload_status,
location::nearby::proto::connections::OperationResultCode::
DETAIL_UNKNOWN);
}
});
@@ -805,7 +810,9 @@ void PayloadManager::SendClientCallbacksForFinishedOutgoingPayload(
// Mark this payload as done for analytics.
client->GetAnalyticsRecorder().OnOutgoingPayloadDone(
endpoint_id, payload_header.id(), status);
endpoint_id, payload_header.id(), status,
location::nearby::proto::connections::OperationResultCode::
DETAIL_UNKNOWN);
}
// Remove these endpoints from our tracking list for this payload.
@@ -845,7 +852,9 @@ void PayloadManager::SendClientCallbacksForFinishedIncomingPayload(
// Analyze
client->GetAnalyticsRecorder().OnIncomingPayloadDone(
endpoint_id, payload_header.id(), status);
endpoint_id, payload_header.id(), status,
location::nearby::proto::connections::OperationResultCode::
DETAIL_UNKNOWN);
});
}
@@ -1139,7 +1148,9 @@ void PayloadManager::HandleSuccessfulOutgoingChunk(
if (is_last_chunk) {
client->GetAnalyticsRecorder().OnOutgoingPayloadDone(
endpoint_id, payload_header.id(),
location::nearby::proto::connections::SUCCESS);
location::nearby::proto::connections::SUCCESS,
location::nearby::proto::connections::OperationResultCode::
DETAIL_UNKNOWN);
// Stop tracking this endpoint.
pending_payload->RemoveEndpoints({endpoint_id});
@@ -1236,7 +1247,9 @@ void PayloadManager::HandleSuccessfulIncomingChunk(
DestroyPendingPayload(payload_header.id());
client->GetAnalyticsRecorder().OnIncomingPayloadDone(
endpoint_id, payload_header.id(),
location::nearby::proto::connections::SUCCESS);
location::nearby::proto::connections::SUCCESS,
location::nearby::proto::connections::OperationResultCode::
DETAIL_UNKNOWN);
} else {
client->GetAnalyticsRecorder().OnPayloadChunkReceived(
endpoint_id, payload_header.id(), payload_chunk_body_size);
@@ -1487,7 +1500,9 @@ void PayloadManager::RecordInvalidPayloadAnalytics(
for (const auto& endpoint_id : endpoint_ids) {
client->GetAnalyticsRecorder().OnOutgoingPayloadDone(
endpoint_id, payload_id,
location::nearby::proto::connections::LOCAL_ERROR);
location::nearby::proto::connections::LOCAL_ERROR,
location::nearby::proto::connections::OperationResultCode::
DETAIL_UNKNOWN);
}
}