From c5fed0c0ed1f2c6a60a606db0b477acc6d2bdb08 Mon Sep 17 00:00:00 2001 From: Edwin Wu Date: Mon, 9 Dec 2024 05:30:03 -0800 Subject: [PATCH] analytics: Add operation result code for Payload analytics, I PiperOrigin-RevId: 704248204 --- .../analytics/analytics_recorder.cc | 66 +++++++++++++++---- .../analytics/analytics_recorder.h | 40 +++++++++-- .../analytics/analytics_recorder_test.cc | 35 +++++++++- connections/implementation/payload_manager.cc | 29 ++++++-- 4 files changed, 141 insertions(+), 29 deletions(-) diff --git a/connections/implementation/analytics/analytics_recorder.cc b/connections/implementation/analytics/analytics_recorder.cc index 566d5827..b1fb0a52 100644 --- a/connections/implementation/analytics/analytics_recorder.cc +++ b/connections/implementation/analytics/analytics_recorder.cc @@ -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 &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 &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(); + 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(); + 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 &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( - pending_payload->type(), pending_payload->total_size_bytes())}); + std::make_unique(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(); }); diff --git a/connections/implementation/analytics/analytics_recorder.h b/connections/implementation/analytics/analytics_recorder.h index 0302835e..22b0cfba 100644 --- a/connections/implementation/analytics/analytics_recorder.h +++ b/connections/implementation/analytics/analytics_recorder.h @@ -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 &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 @@ -298,6 +321,9 @@ class AnalyticsRecorder { absl::btree_map> &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; diff --git a/connections/implementation/analytics/analytics_recorder_test.cc b/connections/implementation/analytics/analytics_recorder_test.cc index 4bcc692a..d72dfacf 100644 --- a/connections/implementation/analytics/analytics_recorder_test.cc +++ b/connections/implementation/analytics/analytics_recorder_test.cc @@ -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)); diff --git a/connections/implementation/payload_manager.cc b/connections/implementation/payload_manager.cc index 12d333e0..6307965f 100644 --- a/connections/implementation/payload_manager.cc +++ b/connections/implementation/payload_manager.cc @@ -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); } }