From de764de60777fa9c9d2facb13d14b098793db26d Mon Sep 17 00:00:00 2001 From: xtkuang <87661715@qq.com> Date: Fri, 14 Aug 2026 08:38:00 +0800 Subject: [PATCH] feat(grpc): persist RPC failures in edge logs --- cmvr-es/common/base/logging/CMakeLists.txt | 23 + cmvr-es/common/base/logging/logger.cpp | 12 +- cmvr-es/common/base/logging/logger.h | 1 + .../common/base/logging/tests/logger_test.cpp | 93 ++ cmvr-es/config/logger/logger.pb.txt | 6 +- cmvr-es/service/CMakeLists.txt | 26 + .../include/grpc_error_logging_interceptor.h | 37 + .../service/grpc/src/grpc_camera_service.cpp | 1 - .../service/grpc/src/grpc_dexhand_service.cpp | 1 - .../src/grpc_error_logging_interceptor.cpp | 504 ++++++++++ .../service/grpc/src/grpc_head_service.cpp | 4 - .../grpc/src/grpc_microphone_service.cpp | 1 - .../service/grpc/src/grpc_speaker_service.cpp | 1 - .../grpc_error_logging_interceptor_test.cpp | 864 ++++++++++++++++++ .../grpc_server_task/src/grpc_server_task.cpp | 8 + 15 files changed, 1570 insertions(+), 12 deletions(-) create mode 100644 cmvr-es/common/base/logging/tests/logger_test.cpp create mode 100644 cmvr-es/service/grpc/include/grpc_error_logging_interceptor.h create mode 100644 cmvr-es/service/grpc/src/grpc_error_logging_interceptor.cpp create mode 100644 cmvr-es/service/grpc/tests/grpc_error_logging_interceptor_test.cpp diff --git a/cmvr-es/common/base/logging/CMakeLists.txt b/cmvr-es/common/base/logging/CMakeLists.txt index c6c31e5c..5322a2e5 100644 --- a/cmvr-es/common/base/logging/CMakeLists.txt +++ b/cmvr-es/common/base/logging/CMakeLists.txt @@ -13,3 +13,26 @@ target_link_libraries(logging PUBLIC add_library(cmvr_es::logging ALIAS logging) install(TARGETS logging ARCHIVE DESTINATION lib) + +if(BUILD_TESTING) + add_executable(logger_test + tests/logger_test.cpp + ) + target_link_libraries(logger_test PRIVATE + cmvr_es::logging + gtest + gtest_main + pthread + ) + add_test(NAME logger_test COMMAND logger_test) + set(_logger_test_environment + "LD_LIBRARY_PATH=${CMVR_TEST_EXTERNAL_LIBRARY_PATH}") + if(CMVR_TEST_SYSTEM_LIBSTDCXX) + list(APPEND _logger_test_environment + "LD_PRELOAD=${CMVR_TEST_SYSTEM_LIBSTDCXX}") + endif() + set_tests_properties(logger_test PROPERTIES + TIMEOUT 10 + ENVIRONMENT "${_logger_test_environment}" + ) +endif() diff --git a/cmvr-es/common/base/logging/logger.cpp b/cmvr-es/common/base/logging/logger.cpp index 0a6355ee..594b95c1 100644 --- a/cmvr-es/common/base/logging/logger.cpp +++ b/cmvr-es/common/base/logging/logger.cpp @@ -167,6 +167,15 @@ void Logger::shutdown() initialized_ = false; } +void Logger::flush() +{ + std::lock_guard lock(mutex_); + if (log_file_.is_open()) { + log_file_.flush(); + last_flush_ = std::chrono::steady_clock::now(); + } +} + bool Logger::enabled(const Level level) const { std::lock_guard lock(mutex_); @@ -200,7 +209,8 @@ void Logger::write(const Level level, rotateIfNeeded_(); log_file_ << line << '\n'; const auto now = std::chrono::steady_clock::now(); - if (level == Level::ERROR || level == Level::FATAL || now - last_flush_ >= flush_interval_) { + if (level == Level::ERROR || level == Level::FATAL || + now - last_flush_ >= flush_interval_) { log_file_.flush(); last_flush_ = now; } diff --git a/cmvr-es/common/base/logging/logger.h b/cmvr-es/common/base/logging/logger.h index a7ee3a74..d6bd17c5 100644 --- a/cmvr-es/common/base/logging/logger.h +++ b/cmvr-es/common/base/logging/logger.h @@ -43,6 +43,7 @@ public: const std::string& application_name, const std::filesystem::path& executable_directory); void shutdown(); + void flush(); bool enabled(Level level) const; void write(Level level, const char* source_file, int source_line, const std::string& message); diff --git a/cmvr-es/common/base/logging/tests/logger_test.cpp b/cmvr-es/common/base/logging/tests/logger_test.cpp new file mode 100644 index 00000000..cc0aef19 --- /dev/null +++ b/cmvr-es/common/base/logging/tests/logger_test.cpp @@ -0,0 +1,93 @@ +#include "common/base/logging/logger.h" + +#include +#include +#include +#include +#include +#include + +#include +#include + +namespace cmvr::logging { +namespace { + +class LoggerTest : public testing::Test { +protected: + void SetUp() override + { + std::array pattern{}; + const std::string value = "/tmp/cmvr-logger-test-XXXXXX"; + std::copy(value.begin(), value.end(), pattern.begin()); + if (char* created = ::mkdtemp(pattern.data())) { + directory_ = created; + } + ASSERT_FALSE(directory_.empty()); + } + + void TearDown() override + { + shutdownLogging(); + std::error_code error; + std::filesystem::remove_all(directory_, error); + } + + config::LoggerConfig warningFileConfig() const + { + config::LoggerConfig config; + config.set_minimum_level(config::LOG_LEVEL_DEBUG); + config.set_directory(directory_.string()); + config.set_flush_interval_seconds(3600); + auto* route = config.add_routes(); + route->set_level(config::LOG_LEVEL_WARNING); + route->set_terminal(false); + route->set_file(true); + return config; + } + + std::string fileContents() const + { + std::ifstream input(directory_ / "logger_test.log"); + return {std::istreambuf_iterator(input), + std::istreambuf_iterator()}; + } + + std::filesystem::path directory_; +}; + +TEST_F(LoggerTest, ExplicitFlushMakesWarningVisibleInFile) +{ + ASSERT_TRUE(initLogging( + warningFileConfig(), "logger_test", directory_)); + + Logger::instance().write( + Level::WARNING, __FILE__, __LINE__, "warning sentinel"); + Logger::instance().flush(); + + EXPECT_NE(fileContents().find("warning sentinel"), std::string::npos); +} + +TEST_F(LoggerTest, ErrorIsVisibleInFileImmediately) +{ + auto config = warningFileConfig(); + config.mutable_routes(0)->set_level(config::LOG_LEVEL_ERROR); + ASSERT_TRUE(initLogging(config, "logger_test", directory_)); + + Logger::instance().write( + Level::ERROR, __FILE__, __LINE__, "error sentinel"); + + EXPECT_NE(fileContents().find("error sentinel"), std::string::npos); +} + +TEST_F(LoggerTest, FlushIsSafeOutsideInitializedLifetime) +{ + Logger::instance().flush(); + ASSERT_TRUE(initLogging( + warningFileConfig(), "logger_test", directory_)); + shutdownLogging(); + Logger::instance().flush(); +} + +} // namespace +} // namespace cmvr::logging diff --git a/cmvr-es/config/logger/logger.pb.txt b/cmvr-es/config/logger/logger.pb.txt index 85f8aac0..54182979 100644 --- a/cmvr-es/config/logger/logger.pb.txt +++ b/cmvr-es/config/logger/logger.pb.txt @@ -13,17 +13,17 @@ logger { routes { level: LOG_LEVEL_WARNING terminal: true - file: false + file: true } routes { level: LOG_LEVEL_ERROR terminal: true - file: false + file: true } routes { level: LOG_LEVEL_FATAL terminal: true - file: false + file: true } directory: "../log" diff --git a/cmvr-es/service/CMakeLists.txt b/cmvr-es/service/CMakeLists.txt index ef742c48..1244e7cb 100644 --- a/cmvr-es/service/CMakeLists.txt +++ b/cmvr-es/service/CMakeLists.txt @@ -6,6 +6,7 @@ add_library(service grpc/src/media_activity_coordinator.cpp grpc/src/motor_activity_coordinator.cpp grpc/src/grpc_camera_service.cpp + grpc/src/grpc_error_logging_interceptor.cpp grpc/src/grpc_system_service.cpp grpc/src/grpc_speaker_service.cpp grpc/src/grpc_microphone_service.cpp @@ -216,6 +217,31 @@ if(BUILD_TESTING) ENVIRONMENT "${_grpc_system_test_environment}" ) + add_executable(grpc_error_logging_interceptor_test + grpc/tests/grpc_error_logging_interceptor_test.cpp + grpc/src/grpc_error_logging_interceptor.cpp + ) + target_include_directories(grpc_error_logging_interceptor_test + PRIVATE + ${CMAKE_SOURCE_DIR}/cmvr-es + ) + target_link_libraries(grpc_error_logging_interceptor_test + PRIVATE + cmvr_es::logging + cmvr_es::proto + gtest + gtest_main + pthread + ) + add_test( + NAME grpc_error_logging_interceptor_test + COMMAND grpc_error_logging_interceptor_test + ) + set_tests_properties(grpc_error_logging_interceptor_test PROPERTIES + TIMEOUT 10 + ENVIRONMENT "${_grpc_system_test_environment}" + ) + add_executable(grpc_arm_service_test grpc/tests/grpc_arm_service_test.cpp ) diff --git a/cmvr-es/service/grpc/include/grpc_error_logging_interceptor.h b/cmvr-es/service/grpc/include/grpc_error_logging_interceptor.h new file mode 100644 index 00000000..f780aa5b --- /dev/null +++ b/cmvr-es/service/grpc/include/grpc_error_logging_interceptor.h @@ -0,0 +1,37 @@ +#pragma once + +#include +#include +#include + +#include + +namespace cmvr::service { + +enum class GrpcFailureKind { + APPLICATION, + GRPC_STATUS, + MALFORMED_RESPONSE, +}; + +enum class GrpcFailureSeverity { + INFO, + WARNING, + ERROR, +}; + +struct GrpcFailureRecord { + GrpcFailureKind kind{GrpcFailureKind::APPLICATION}; + std::string method; + std::string peer; + grpc::StatusCode status_code{grpc::StatusCode::OK}; + GrpcFailureSeverity severity{GrpcFailureSeverity::ERROR}; + std::string detail; +}; + +using GrpcFailureSink = std::function; + +std::unique_ptr +makeGrpcErrorLoggingInterceptorFactory(GrpcFailureSink sink = {}); + +} // namespace cmvr::service diff --git a/cmvr-es/service/grpc/src/grpc_camera_service.cpp b/cmvr-es/service/grpc/src/grpc_camera_service.cpp index f0a4dc64..1f7712b3 100644 --- a/cmvr-es/service/grpc/src/grpc_camera_service.cpp +++ b/cmvr-es/service/grpc/src/grpc_camera_service.cpp @@ -21,7 +21,6 @@ using namespace cmvr::device; namespace { template grpc::Status failResponse(ResponseT* response, const std::string& message) { - CMVR_LOG(ERROR) << "[gRPCCameraServiceImpl] " << message; response->mutable_header()->set_success(false); response->mutable_header()->set_error_message(message); setCurrentTimestamp(response->mutable_header()->mutable_timestamp()); diff --git a/cmvr-es/service/grpc/src/grpc_dexhand_service.cpp b/cmvr-es/service/grpc/src/grpc_dexhand_service.cpp index 6149aa6f..d0fcecc1 100644 --- a/cmvr-es/service/grpc/src/grpc_dexhand_service.cpp +++ b/cmvr-es/service/grpc/src/grpc_dexhand_service.cpp @@ -120,7 +120,6 @@ bool applyFreedomValues(const FreedomCollection& freedoms, template grpc::Status failResponse(ResponseT* response, const std::string& message) { - CMVR_LOG(ERROR) << "[gRPCDexHandServiceImpl] " << message; response->mutable_header()->set_success(false); response->mutable_header()->set_error_message(message); setCurrentTimestamp(response->mutable_header()->mutable_timestamp()); diff --git a/cmvr-es/service/grpc/src/grpc_error_logging_interceptor.cpp b/cmvr-es/service/grpc/src/grpc_error_logging_interceptor.cpp new file mode 100644 index 00000000..ecbe3da9 --- /dev/null +++ b/cmvr-es/service/grpc/src/grpc_error_logging_interceptor.cpp @@ -0,0 +1,504 @@ +#include "service/grpc/include/grpc_error_logging_interceptor.h" + +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include +#include +#include +#include +#include +#include +#include + +#include "cmvr/api/common.pb.h" +#include "common/base/logging/logger.h" + +namespace cmvr::service { + +namespace { + +using Hook = grpc::experimental::InterceptionHookPoints; +using CodedInputStream = google::protobuf::io::CodedInputStream; +using WireFormatLite = google::protobuf::internal::WireFormatLite; + +constexpr std::string_view kFeedbackType = + "cmvr.api.CommandHeader.Feedback"; +constexpr std::size_t kMaxLogDetailBytes = 1024; +constexpr int kMaxFeedbackBytes = 64 * 1024; + +enum class ResponseShape { + NONE, + DIRECT_FEEDBACK, + ENVELOPE, +}; + +struct ApplicationFeedback { + bool present{false}; + bool parsed{false}; + bool success{false}; + std::string error_message; +}; + +ApplicationFeedback feedbackFromProto( + const api::CommandHeader_Feedback& message) +{ + return ApplicationFeedback{ + true, true, message.success(), message.error_message()}; +} + +std::string boundedDetail(const std::string& detail) +{ + if (detail.size() <= kMaxLogDetailBytes) { + return detail; + } + + constexpr std::string_view suffix = "... [truncated]"; + std::size_t length = kMaxLogDetailBytes - suffix.size(); + while (length > 0 && + (static_cast(detail[length]) & 0xc0U) == 0x80U) { + --length; + } + return detail.substr(0, length) + std::string(suffix); +} + +const char* failureKindName(const GrpcFailureKind kind) +{ + switch (kind) { + case GrpcFailureKind::APPLICATION: + return "application"; + case GrpcFailureKind::GRPC_STATUS: + return "grpc_status"; + case GrpcFailureKind::MALFORMED_RESPONSE: + return "malformed_response"; + } + return "unknown"; +} + +const char* statusCodeName(const grpc::StatusCode code) +{ + switch (code) { + case grpc::StatusCode::OK: return "OK"; + case grpc::StatusCode::CANCELLED: return "CANCELLED"; + case grpc::StatusCode::UNKNOWN: return "UNKNOWN"; + case grpc::StatusCode::INVALID_ARGUMENT: return "INVALID_ARGUMENT"; + case grpc::StatusCode::DEADLINE_EXCEEDED: return "DEADLINE_EXCEEDED"; + case grpc::StatusCode::NOT_FOUND: return "NOT_FOUND"; + case grpc::StatusCode::ALREADY_EXISTS: return "ALREADY_EXISTS"; + case grpc::StatusCode::PERMISSION_DENIED: return "PERMISSION_DENIED"; + case grpc::StatusCode::RESOURCE_EXHAUSTED: return "RESOURCE_EXHAUSTED"; + case grpc::StatusCode::FAILED_PRECONDITION: return "FAILED_PRECONDITION"; + case grpc::StatusCode::ABORTED: return "ABORTED"; + case grpc::StatusCode::OUT_OF_RANGE: return "OUT_OF_RANGE"; + case grpc::StatusCode::UNIMPLEMENTED: return "UNIMPLEMENTED"; + case grpc::StatusCode::INTERNAL: return "INTERNAL"; + case grpc::StatusCode::UNAVAILABLE: return "UNAVAILABLE"; + case grpc::StatusCode::DATA_LOSS: return "DATA_LOSS"; + case grpc::StatusCode::UNAUTHENTICATED: return "UNAUTHENTICATED"; + case grpc::StatusCode::DO_NOT_USE: break; + } + return "UNKNOWN_CODE"; +} + +void logFailure(const GrpcFailureRecord& record) +{ + std::ostringstream message; + message << "[gRPC] request failed, method=" << record.method + << ", peer=" << (record.peer.empty() ? "unknown" : record.peer) + << ", kind=" << failureKindName(record.kind) + << ", code=" << static_cast(record.status_code) + << '(' << statusCodeName(record.status_code) << ')' + << ", detail=" << boundedDetail(record.detail); + switch (record.severity) { + case GrpcFailureSeverity::INFO: + CMVR_LOG(INFO) << message.str(); + break; + case GrpcFailureSeverity::WARNING: + CMVR_LOG(WARNING) << message.str(); + logging::Logger::instance().flush(); + break; + case GrpcFailureSeverity::ERROR: + CMVR_LOG(ERROR) << message.str(); + break; + } +} + +GrpcFailureSeverity severityForStatus(const grpc::StatusCode code) +{ + switch (code) { + case grpc::StatusCode::CANCELLED: + case grpc::StatusCode::INVALID_ARGUMENT: + case grpc::StatusCode::DEADLINE_EXCEEDED: + case grpc::StatusCode::NOT_FOUND: + case grpc::StatusCode::ALREADY_EXISTS: + case grpc::StatusCode::PERMISSION_DENIED: + case grpc::StatusCode::RESOURCE_EXHAUSTED: + case grpc::StatusCode::FAILED_PRECONDITION: + case grpc::StatusCode::ABORTED: + case grpc::StatusCode::OUT_OF_RANGE: + case grpc::StatusCode::UNIMPLEMENTED: + case grpc::StatusCode::UNAVAILABLE: + case grpc::StatusCode::UNAUTHENTICATED: + return GrpcFailureSeverity::WARNING; + case grpc::StatusCode::OK: + case grpc::StatusCode::UNKNOWN: + case grpc::StatusCode::INTERNAL: + case grpc::StatusCode::DATA_LOSS: + case grpc::StatusCode::DO_NOT_USE: + return GrpcFailureSeverity::ERROR; + } + return GrpcFailureSeverity::ERROR; +} + +bool splitMethodName(const std::string_view full_method, + std::string_view& service_name, + std::string_view& method_name) +{ + if (full_method.empty()) { + return false; + } + const std::size_t service_begin = full_method.front() == '/' ? 1 : 0; + const std::size_t separator = full_method.find('/', service_begin); + if (separator == std::string_view::npos || + separator == service_begin || separator + 1 >= full_method.size()) { + return false; + } + service_name = full_method.substr(service_begin, separator - service_begin); + method_name = full_method.substr(separator + 1); + return true; +} + +ResponseShape responseShapeForMethod(const std::string_view full_method) +{ + std::string_view service_name; + std::string_view method_name; + if (!splitMethodName(full_method, service_name, method_name)) { + return ResponseShape::NONE; + } + + const auto* pool = google::protobuf::DescriptorPool::generated_pool(); + const auto* service = pool->FindServiceByName(std::string(service_name)); + if (service == nullptr) { + return ResponseShape::NONE; + } + const auto* method = service->FindMethodByName(std::string(method_name)); + if (method == nullptr || method->output_type() == nullptr) { + return ResponseShape::NONE; + } + + const auto* output = method->output_type(); + if (output->full_name() == kFeedbackType) { + return ResponseShape::DIRECT_FEEDBACK; + } + + const auto* header = output->FindFieldByNumber(1); + if (header == nullptr || header->name() != "header" || + header->is_repeated() || + header->cpp_type() != google::protobuf::FieldDescriptor::CPPTYPE_MESSAGE || + header->message_type() == nullptr || + header->message_type()->full_name() != kFeedbackType) { + return ResponseShape::NONE; + } + return ResponseShape::ENVELOPE; +} + +bool mergeFeedback(CodedInputStream& input, + api::CommandHeader_Feedback& feedback) +{ + return feedback.MergePartialFromCodedStream(&input) && + input.ConsumedEntireMessage(); +} + +bool mergeBoundedFeedback(CodedInputStream& input, + api::CommandHeader_Feedback& feedback) +{ + std::uint32_t length = 0; + if (!input.ReadVarint32(&length) || + length > static_cast(kMaxFeedbackBytes) || + length > static_cast(std::numeric_limits::max())) { + return false; + } + + std::string payload; + if (!input.ReadString(&payload, static_cast(length))) { + return false; + } + return feedback.MergeFromString(payload); +} + +ApplicationFeedback parseDirectFeedback(grpc::ByteBuffer& buffer) +{ + ApplicationFeedback feedback; + grpc::ProtoBufferReader reader(&buffer); + if (!reader.status().ok()) { + return feedback; + } + CodedInputStream input(&reader); + input.SetTotalBytesLimit(kMaxFeedbackBytes); + api::CommandHeader_Feedback message; + if (!mergeFeedback(input, message)) { + feedback.present = true; + return feedback; + } + return feedbackFromProto(message); +} + +ApplicationFeedback parseEnvelopeFeedback(grpc::ByteBuffer& buffer) +{ + ApplicationFeedback feedback; + grpc::ProtoBufferReader reader(&buffer); + if (!reader.status().ok()) { + return feedback; + } + CodedInputStream input(&reader); + const std::uint32_t tag = input.ReadTag(); + if (tag == 0 || WireFormatLite::GetTagFieldNumber(tag) != 1) { + feedback.parsed = true; + return feedback; + } + + feedback.present = true; + if (WireFormatLite::GetTagWireType(tag) != + WireFormatLite::WIRETYPE_LENGTH_DELIMITED) { + return feedback; + } + + api::CommandHeader_Feedback message; + if (!mergeBoundedFeedback(input, message)) { + return feedback; + } + return feedbackFromProto(message); +} + +class GrpcErrorLoggingInterceptor final + : public grpc::experimental::Interceptor { +public: + GrpcErrorLoggingInterceptor(std::string method, + grpc::ServerContextBase* context, + const ResponseShape response_shape, + const bool server_streaming, + const bool client_streaming, + GrpcFailureSink sink) + : method_(std::move(method)), + context_(context), + peer_(context == nullptr ? std::string{} : context->peer()), + response_shape_(response_shape), + server_streaming_(server_streaming), + client_streaming_(client_streaming), + sink_(std::move(sink)) + { + } + + void Intercept(grpc::experimental::InterceptorBatchMethods* methods) override + { + if (methods->QueryInterceptionHookPoint(Hook::PRE_SEND_CANCEL)) { + // gRPC forbids delaying this hook. Only publish the signal here; + // the final status hook performs any logging. + server_cancel_requested_.store(true, std::memory_order_release); + return; + } + if (methods->QueryInterceptionHookPoint(Hook::PRE_SEND_MESSAGE)) { + inspectResponse(methods->GetSerializedSendMessage()); + } + if (methods->QueryInterceptionHookPoint(Hook::PRE_SEND_STATUS)) { + inspectStatus(methods->GetSendStatus()); + } + methods->Proceed(); + } + +private: + void emit(GrpcFailureRecord record) + { + record.method = method_; + record.peer = peer_; + record.detail = boundedDetail(record.detail); + if (sink_) { + std::lock_guard lock(sink_mutex_); + sink_(record); + } else { + logFailure(record); + } + } + + void inspectResponse(grpc::ByteBuffer* buffer) + { + if (response_shape_ == ResponseShape::NONE || + buffer == nullptr || !buffer->Valid()) { + return; + } + + ApplicationFeedback feedback = + response_shape_ == ResponseShape::DIRECT_FEEDBACK + ? parseDirectFeedback(*buffer) + : parseEnvelopeFeedback(*buffer); + std::optional failure; + if (!feedback.present) { + failure = GrpcFailureRecord{ + GrpcFailureKind::MALFORMED_RESPONSE, + {}, {}, grpc::StatusCode::INTERNAL, + GrpcFailureSeverity::ERROR, + "response is missing CommandHeader.Feedback"}; + } else if (!feedback.parsed) { + failure = GrpcFailureRecord{ + GrpcFailureKind::MALFORMED_RESPONSE, + {}, {}, grpc::StatusCode::INTERNAL, + GrpcFailureSeverity::ERROR, + "response contains an invalid CommandHeader.Feedback"}; + } else if (!feedback.success) { + failure = GrpcFailureRecord{ + GrpcFailureKind::APPLICATION, + {}, {}, grpc::StatusCode::OK, + GrpcFailureSeverity::ERROR, + feedback.error_message.empty() + ? "operation failed without an error message" + : feedback.error_message}; + } + + if (!failure.has_value()) { + return; + } + if (server_streaming_) { + emit(std::move(*failure)); + return; + } + pending_failure_ = std::move(*failure); + } + + void inspectStatus(const grpc::Status& status) + { + const grpc::StatusCode status_code = status.error_code(); + std::string status_detail = status.error_message(); + if (!status.ok() && status_detail.empty()) { + status_detail = "RPC completed with a non-OK status"; + } + + if (!status.ok()) { + // A non-OK unary/client-streaming status suppresses the response + // body, so only the status is visible to the platform. + pending_failure_.reset(); + if (status_code == grpc::StatusCode::CANCELLED) { + emitCancellationOnce(status_code, std::move(status_detail)); + return; + } + emit({GrpcFailureKind::GRPC_STATUS, + {}, {}, status_code, + severityForStatus(status_code), + status_detail}); + return; + } + + const bool deadline_expired = + context_ != nullptr && + context_->deadline() <= std::chrono::system_clock::now(); + // IsCancelled() waits for the client close on the synchronous API. + // Unary requests have already sent that close, but querying it here + // could block an early client-streaming/bidi response indefinitely. + if (deadline_expired || + (context_ != nullptr && !client_streaming_ && + context_->IsCancelled())) { + pending_failure_.reset(); + const grpc::StatusCode code = deadline_expired + ? grpc::StatusCode::DEADLINE_EXCEEDED + : grpc::StatusCode::CANCELLED; + std::string detail; + if (deadline_expired) { + detail = "RPC deadline expired before completion"; + } else if (server_cancel_requested_.load( + std::memory_order_acquire)) { + detail = + "RPC cancellation was requested by the edge server before completion"; + } else { + detail = + "RPC was cancelled by the client or transport before completion"; + } + emitCancellationOnce(code, std::move(detail)); + return; + } + + if (server_cancel_requested_.load(std::memory_order_acquire)) { + pending_failure_.reset(); + emitCancellationOnce( + grpc::StatusCode::CANCELLED, + "RPC cancellation was requested by the edge server before completion"); + return; + } + + if (pending_failure_.has_value()) { + emit(std::move(*pending_failure_)); + pending_failure_.reset(); + } + } + + void emitCancellationOnce(const grpc::StatusCode code, std::string detail) + { + if (cancellation_recorded_.exchange(true, std::memory_order_acq_rel)) { + return; + } + emit({GrpcFailureKind::GRPC_STATUS, + {}, {}, code, severityForStatus(code), std::move(detail)}); + } + + std::string method_; + grpc::ServerContextBase* context_{nullptr}; + std::string peer_; + ResponseShape response_shape_{ResponseShape::NONE}; + bool server_streaming_{false}; + bool client_streaming_{false}; + GrpcFailureSink sink_; + std::optional pending_failure_; + std::mutex sink_mutex_; + std::atomic server_cancel_requested_{false}; + std::atomic cancellation_recorded_{false}; +}; + +class GrpcErrorLoggingInterceptorFactory final + : public grpc::experimental::ServerInterceptorFactoryInterface { +public: + explicit GrpcErrorLoggingInterceptorFactory(GrpcFailureSink sink) + : sink_(std::move(sink)) + { + } + + grpc::experimental::Interceptor* CreateServerInterceptor( + grpc::experimental::ServerRpcInfo* info) override + { + const std::string method = + info != nullptr && info->method() != nullptr + ? info->method() + : ""; + const bool server_streaming = + info != nullptr && + (info->type() == grpc::experimental::ServerRpcInfo::Type::SERVER_STREAMING || + info->type() == grpc::experimental::ServerRpcInfo::Type::BIDI_STREAMING); + const bool client_streaming = + info != nullptr && + (info->type() == grpc::experimental::ServerRpcInfo::Type::CLIENT_STREAMING || + info->type() == grpc::experimental::ServerRpcInfo::Type::BIDI_STREAMING); + return new GrpcErrorLoggingInterceptor( + method, + info == nullptr ? nullptr : info->server_context(), + responseShapeForMethod(method), server_streaming, + client_streaming, sink_); + } + +private: + GrpcFailureSink sink_; +}; + +} // namespace + +std::unique_ptr +makeGrpcErrorLoggingInterceptorFactory(GrpcFailureSink sink) +{ + return std::make_unique( + std::move(sink)); +} + +} // namespace cmvr::service diff --git a/cmvr-es/service/grpc/src/grpc_head_service.cpp b/cmvr-es/service/grpc/src/grpc_head_service.cpp index 9ba03019..8e6934e7 100644 --- a/cmvr-es/service/grpc/src/grpc_head_service.cpp +++ b/cmvr-es/service/grpc/src/grpc_head_service.cpp @@ -20,7 +20,6 @@ using namespace cmvr::api; namespace { template grpc::Status failResponse(ResponseT* response, const std::string& message) { - CMVR_LOG(ERROR) << "[gRPCMBioHeadServiceImpl] " << message; response->mutable_header()->set_success(false); response->mutable_header()->set_error_message(message); setCurrentTimestamp(response->mutable_header()->mutable_timestamp()); @@ -150,7 +149,6 @@ grpc::Status gRPCMBioHeadServiceImpl::StreamExpression( if (first_message) { dev_id = request_msg.header().device_id(); if (dev_id.empty()) { - CMVR_LOG(ERROR) << "[gRPCMBioHeadServiceImpl] Device ID is empty in first message"; feedback_msg.mutable_header()->set_success(false); feedback_msg.mutable_header()->set_error_message("Device ID is empty in first message"); setCurrentTimestamp(feedback_msg.mutable_header()->mutable_timestamp()); @@ -160,7 +158,6 @@ grpc::Status gRPCMBioHeadServiceImpl::StreamExpression( robot = dmgr_.getDevice(dev_id); if (!robot) { const std::string message = "Biohead device not found: " + dev_id; - CMVR_LOG(ERROR) << "[gRPCMBioHeadServiceImpl] " << message; feedback_msg.mutable_header()->set_success(false); feedback_msg.mutable_header()->set_error_message(message); setCurrentTimestamp(feedback_msg.mutable_header()->mutable_timestamp()); @@ -259,7 +256,6 @@ grpc::Status gRPCMBioHeadServiceImpl::StreamExpression( CMVR_LOG(DEBUG) << "[gRPCMBioHeadServiceImpl] (StreamExpression): finished, id=" << dev_id; return grpc::Status::OK; } catch (const std::exception& e) { - CMVR_LOG(ERROR) << "StreamExpression error: " << e.what(); feedback_msg.mutable_header()->set_success(false); feedback_msg.mutable_header()->set_error_message(e.what()); setCurrentTimestamp(feedback_msg.mutable_header()->mutable_timestamp()); diff --git a/cmvr-es/service/grpc/src/grpc_microphone_service.cpp b/cmvr-es/service/grpc/src/grpc_microphone_service.cpp index c44a0f57..e2cb0272 100644 --- a/cmvr-es/service/grpc/src/grpc_microphone_service.cpp +++ b/cmvr-es/service/grpc/src/grpc_microphone_service.cpp @@ -18,7 +18,6 @@ using namespace cmvr::service; namespace { template grpc::Status failResponse(ResponseT* response, const std::string& message) { - CMVR_LOG(ERROR) << "[gRPCMicroPhoneServiceImpl] " << message; response->mutable_header()->set_success(false); response->mutable_header()->set_error_message(message); setCurrentTimestamp(response->mutable_header()->mutable_timestamp()); diff --git a/cmvr-es/service/grpc/src/grpc_speaker_service.cpp b/cmvr-es/service/grpc/src/grpc_speaker_service.cpp index 4eaf03dd..f80a257f 100644 --- a/cmvr-es/service/grpc/src/grpc_speaker_service.cpp +++ b/cmvr-es/service/grpc/src/grpc_speaker_service.cpp @@ -14,7 +14,6 @@ using namespace cmvr::service; namespace { template grpc::Status failResponse(ResponseT* response, const std::string& message) { - CMVR_LOG(ERROR) << "[gRPCSpeakerServiceImpl] " << message; response->mutable_header()->set_success(false); response->mutable_header()->set_error_message(message); setCurrentTimestamp(response->mutable_header()->mutable_timestamp()); diff --git a/cmvr-es/service/grpc/tests/grpc_error_logging_interceptor_test.cpp b/cmvr-es/service/grpc/tests/grpc_error_logging_interceptor_test.cpp new file mode 100644 index 00000000..02b7adca --- /dev/null +++ b/cmvr-es/service/grpc/tests/grpc_error_logging_interceptor_test.cpp @@ -0,0 +1,864 @@ +#include "service/grpc/include/grpc_error_logging_interceptor.h" + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include +#include +#include +#include +#include +#include + +#include "cmvr/api/agv_service.grpc.pb.h" +#include "cmvr/api/camera_service.grpc.pb.h" +#include "cmvr/api/test_service.grpc.pb.h" +#include "common/base/logging/logger.h" + +namespace cmvr::service { +namespace { + +std::atomic g_socket_sequence{0}; + +void setRpcDeadline(grpc::ClientContext& context) +{ + context.set_deadline( + std::chrono::system_clock::now() + std::chrono::seconds(3)); +} + +grpc::ByteBuffer byteBuffer(const std::string& payload) +{ + const grpc::Slice slice(payload); + return grpc::ByteBuffer(&slice, 1); +} + +std::string byteBufferString(const grpc::ByteBuffer& buffer) +{ + std::vector slices; + EXPECT_TRUE(buffer.Dump(&slices).ok()); + std::string payload; + for (const auto& slice : slices) { + payload.append( + reinterpret_cast(slice.begin()), slice.size()); + } + return payload; +} + +grpc::Status rawUnaryCall(const std::shared_ptr& channel, + const std::string& method, + const std::string& request_payload, + grpc::ByteBuffer& response) +{ + grpc::ClientContext context; + setRpcDeadline(context); + const auto request = byteBuffer(request_payload); + const grpc::internal::RpcMethod rpc_method( + method.c_str(), grpc::internal::RpcMethod::NORMAL_RPC); + return grpc::internal::BlockingUnaryCall( + channel.get(), rpc_method, &context, request, &response); +} + +class ErrorLoggingTestService final : public api::TestService::Service { +public: + grpc::Status Call(grpc::ServerContext* context, + const api::TestReqeust* request, + api::TestResponse* response) override + { + if (request->data() == "grpc-error") { + return grpc::Status( + grpc::StatusCode::INTERNAL, "transport sentinel"); + } + if (request->data() == "cancelled") { + return grpc::Status( + grpc::StatusCode::CANCELLED, "client closed stream"); + } + if (request->data() == "server-cancel") { + context->TryCancel(); + return grpc::Status::OK; + } + if (request->data() == "client-cancel") { + { + std::lock_guard lock(cancel_mutex_); + cancel_started_ = true; + } + cancel_started_condition_.notify_all(); + const auto timeout = + std::chrono::steady_clock::now() + std::chrono::seconds(3); + while (!context->IsCancelled() && + std::chrono::steady_clock::now() < timeout) { + std::this_thread::sleep_for(std::chrono::milliseconds(1)); + } + return grpc::Status::OK; + } + response->set_data("ok"); + return grpc::Status::OK; + } + + bool waitForCancelableCall() + { + std::unique_lock lock(cancel_mutex_); + return cancel_started_condition_.wait_for( + lock, std::chrono::seconds(3), + [this] { return cancel_started_; }); + } + +private: + std::mutex cancel_mutex_; + std::condition_variable cancel_started_condition_; + bool cancel_started_{false}; +}; + +class ErrorLoggingAgvService final : public api::AgvService::Service { +public: + grpc::Status emergencyStop( + grpc::ServerContext*, + const api::CommandHeader_Request* request, + api::CommandHeader_Feedback* response) override + { + if (request->device_id() == "missing") { + response->set_success(false); + response->set_error_message("missing sentinel"); + return grpc::Status( + grpc::StatusCode::NOT_FOUND, "missing sentinel"); + } + response->set_success(false); + response->set_error_message("direct sentinel"); + return grpc::Status::OK; + } + + grpc::Status streamMap( + grpc::ServerContext*, + const api::AgvMapStreamCommand_Request* request, + grpc::ServerWriter* writer) override + { + const auto write_failure = [writer](const std::string& detail) { + api::AgvMapStreamCommand_Feedback response; + response.mutable_header()->set_success(false); + response.mutable_header()->set_error_message(detail); + return writer->Write(response); + }; + + if (request->resume_token() == "multiple") { + write_failure("first stream sentinel"); + write_failure("second stream sentinel"); + return grpc::Status::OK; + } + if (request->resume_token() == "hold") { + write_failure("immediate stream sentinel"); + std::unique_lock lock(stream_mutex_); + stream_released_.wait_for( + lock, std::chrono::seconds(3), + [this] { return release_stream_; }); + return grpc::Status::OK; + } + if (request->resume_token() == "different") { + write_failure("application stream sentinel"); + return grpc::Status( + grpc::StatusCode::INTERNAL, "status stream sentinel"); + } + if (request->resume_token() == "long-different") { + const std::string common_prefix(1100, 'x'); + write_failure(common_prefix + " application"); + return grpc::Status( + grpc::StatusCode::INTERNAL, + common_prefix + " status"); + } + + write_failure("stream sentinel"); + return grpc::Status( + grpc::StatusCode::INTERNAL, "stream sentinel"); + } + + void releaseHeldStream() + { + { + std::lock_guard lock(stream_mutex_); + release_stream_ = true; + } + stream_released_.notify_all(); + } + +private: + std::mutex stream_mutex_; + std::condition_variable stream_released_; + bool release_stream_{false}; +}; + +class ErrorLoggingCameraService final : public api::CameraService::Service { +public: + grpc::Status StartCamera( + grpc::ServerContext*, + const api::StartCameraCommand_Request*, + api::StartCameraCommand_Feedback* response) override + { + response->mutable_header()->set_success(false); + response->mutable_header()->set_error_message("nested sentinel"); + return grpc::Status::OK; + } + + grpc::Status GetRGBImage( + grpc::ServerContext*, + const api::GetRGBImageCommand_Request*, + api::GetRGBImageCommand_Feedback* response) override + { + response->mutable_header()->set_success(true); + response->mutable_color_frame()->set_data( + std::string(2 * 1024 * 1024, 'x')); + return grpc::Status::OK; + } +}; + +class RawCameraService final + : public api::CameraService::WithRawCallbackMethod_StartCamera< + api::CameraService::Service> { +public: + explicit RawCameraService(std::string response) + : response_(std::move(response)) + { + } + + grpc::ServerUnaryReactor* StartCamera( + grpc::CallbackServerContext* context, + const grpc::ByteBuffer*, + grpc::ByteBuffer* response) override + { + *response = byteBuffer(response_); + auto* reactor = context->DefaultReactor(); + reactor->Finish(grpc::Status::OK); + return reactor; + } + +private: + std::string response_; +}; + +struct RawCallResult { + grpc::Status status; + std::vector records; + std::string response; +}; + +class ScopedServer final { +public: + ScopedServer(std::unique_ptr server, + std::string socket_path = {}) + : server_(std::move(server)), socket_path_(std::move(socket_path)) + { + } + + ~ScopedServer() + { + if (server_) { + server_->Shutdown(); + server_->Wait(); + } + if (!socket_path_.empty()) { + std::remove(socket_path_.c_str()); + } + } + + grpc::Server* get() const { return server_.get(); } + +private: + std::unique_ptr server_; + std::string socket_path_; +}; + +class ScopedSocket final { +public: + explicit ScopedSocket(std::string path) : path_(std::move(path)) + { + std::remove(path_.c_str()); + } + + ~ScopedSocket() { std::remove(path_.c_str()); } + +private: + std::string path_; +}; + +RawCallResult callRawStartCamera(const std::string& response_payload) +{ + std::mutex records_mutex; + std::vector records; + RawCameraService service(response_payload); + grpc::ServerBuilder builder; + const std::string address = + "unix:/tmp/cmvr_grpc_error_logging_raw_" + + std::to_string(static_cast(::getpid())) + "_" + + std::to_string(g_socket_sequence.fetch_add(1U)) + ".sock"; + const std::string socket_path = address.substr(5); + ScopedSocket socket(socket_path); + builder.AddListeningPort(address, grpc::InsecureServerCredentials()); + builder.RegisterService(&service); + std::vector> factories; + factories.emplace_back(makeGrpcErrorLoggingInterceptorFactory( + [&records, &records_mutex](const GrpcFailureRecord& record) { + std::lock_guard lock(records_mutex); + records.push_back(record); + })); + builder.experimental().SetInterceptorCreators(std::move(factories)); + ScopedServer server(builder.BuildAndStart()); + EXPECT_NE(server.get(), nullptr); + if (server.get() == nullptr) { + return {}; + } + + const auto channel = grpc::CreateChannel( + address, grpc::InsecureChannelCredentials()); + api::StartCameraCommand_Request request; + grpc::ByteBuffer response; + const auto status = rawUnaryCall( + channel, "/cmvr.api.CameraService/StartCamera", + request.SerializeAsString(), response); + std::lock_guard lock(records_mutex); + return {status, records, byteBufferString(response)}; +} + +class GrpcErrorLoggingInterceptorTest : public testing::Test { +protected: + void SetUp() override + { + grpc::ServerBuilder builder; + socket_path_ = + "/tmp/cmvr_grpc_error_logging_interceptor_test_" + + std::to_string(static_cast(::getpid())) + "_" + + std::to_string(g_socket_sequence.fetch_add(1U)) + ".sock"; + std::remove(socket_path_.c_str()); + const std::string address = "unix:" + socket_path_; + builder.AddListeningPort( + address, grpc::InsecureServerCredentials()); + builder.RegisterService(&test_service_); + builder.RegisterService(&agv_service_); + builder.RegisterService(&camera_service_); + + std::vector> factories; + factories.emplace_back(makeGrpcErrorLoggingInterceptorFactory( + [this](const GrpcFailureRecord& record) { + { + std::lock_guard lock(records_mutex_); + records_.push_back(record); + } + records_changed_.notify_all(); + })); + builder.experimental().SetInterceptorCreators(std::move(factories)); + server_ = builder.BuildAndStart(); + ASSERT_NE(server_, nullptr); + + const auto channel = grpc::CreateChannel( + address, + grpc::InsecureChannelCredentials()); + test_stub_ = api::TestService::NewStub(channel); + agv_stub_ = api::AgvService::NewStub(channel); + camera_stub_ = api::CameraService::NewStub(channel); + } + + void TearDown() override + { + if (server_) { + server_->Shutdown(); + server_->Wait(); + } + test_stub_.reset(); + agv_stub_.reset(); + camera_stub_.reset(); + if (!socket_path_.empty()) { + std::remove(socket_path_.c_str()); + } + } + + std::vector records() + { + std::lock_guard lock(records_mutex_); + return records_; + } + + std::vector waitForRecords( + const std::size_t count, + const std::chrono::milliseconds timeout = std::chrono::seconds(3)) + { + std::unique_lock lock(records_mutex_); + records_changed_.wait_for( + lock, timeout, [this, count] { return records_.size() >= count; }); + return records_; + } + + ErrorLoggingTestService test_service_; + ErrorLoggingAgvService agv_service_; + ErrorLoggingCameraService camera_service_; + std::unique_ptr server_; + std::unique_ptr test_stub_; + std::unique_ptr agv_stub_; + std::unique_ptr camera_stub_; + std::string socket_path_; + std::mutex records_mutex_; + std::condition_variable records_changed_; + std::vector records_; +}; + +TEST_F(GrpcErrorLoggingInterceptorTest, RecordsNonOkGrpcStatus) +{ + grpc::ClientContext context; + setRpcDeadline(context); + api::TestReqeust request; + api::TestResponse response; + request.set_data("grpc-error"); + + const grpc::Status status = test_stub_->Call(&context, request, &response); + ASSERT_EQ(status.error_code(), grpc::StatusCode::INTERNAL); + + const auto captured = waitForRecords(1U); + ASSERT_EQ(captured.size(), 1U); + EXPECT_EQ(captured.front().kind, GrpcFailureKind::GRPC_STATUS); + EXPECT_EQ(captured.front().method, "/cmvr.api.TestService/Call"); + EXPECT_EQ(captured.front().status_code, grpc::StatusCode::INTERNAL); + EXPECT_EQ(captured.front().severity, GrpcFailureSeverity::ERROR); + EXPECT_EQ(captured.front().detail, "transport sentinel"); +} + +TEST_F(GrpcErrorLoggingInterceptorTest, RecordsDirectFeedbackFailure) +{ + grpc::ClientContext context; + setRpcDeadline(context); + api::CommandHeader_Request request; + api::CommandHeader_Feedback response; + + ASSERT_TRUE(agv_stub_->emergencyStop(&context, request, &response).ok()); + ASSERT_FALSE(response.success()); + + const auto captured = waitForRecords(1U); + ASSERT_EQ(captured.size(), 1U); + EXPECT_EQ(captured.front().kind, GrpcFailureKind::APPLICATION); + EXPECT_EQ(captured.front().method, "/cmvr.api.AgvService/emergencyStop"); + EXPECT_EQ(captured.front().status_code, grpc::StatusCode::OK); + EXPECT_EQ(captured.front().severity, GrpcFailureSeverity::ERROR); + EXPECT_EQ(captured.front().detail, "direct sentinel"); +} + +TEST_F(GrpcErrorLoggingInterceptorTest, RecordsCancellationAsWarning) +{ + grpc::ClientContext context; + setRpcDeadline(context); + api::TestReqeust request; + api::TestResponse response; + request.set_data("cancelled"); + + const grpc::Status status = test_stub_->Call(&context, request, &response); + ASSERT_EQ(status.error_code(), grpc::StatusCode::CANCELLED); + + const auto captured = waitForRecords(1U); + ASSERT_EQ(captured.size(), 1U); + EXPECT_EQ(captured.front().kind, GrpcFailureKind::GRPC_STATUS); + EXPECT_EQ(captured.front().severity, GrpcFailureSeverity::WARNING); + EXPECT_EQ(captured.front().detail, "client closed stream"); +} + +TEST(GrpcErrorLoggingInterceptorFileTest, WarningStatusIsFlushedToFile) +{ + std::array directory_pattern{}; + const std::string pattern = "/tmp/cmvr-grpc-log-test-XXXXXX"; + std::copy(pattern.begin(), pattern.end(), directory_pattern.begin()); + const char* created = ::mkdtemp(directory_pattern.data()); + ASSERT_NE(created, nullptr); + const std::filesystem::path directory(created); + + struct Cleanup { + ~Cleanup() + { + logging::shutdownLogging(); + if (!socket_path.empty()) { + std::remove(socket_path.c_str()); + } + std::error_code error; + std::filesystem::remove_all(directory, error); + } + + std::filesystem::path directory; + std::string socket_path; + } cleanup{directory, {}}; + + config::LoggerConfig config; + config.set_minimum_level(config::LOG_LEVEL_WARNING); + config.set_directory(directory.string()); + config.set_flush_interval_seconds(3600); + auto* warning_route = config.add_routes(); + warning_route->set_level(config::LOG_LEVEL_WARNING); + warning_route->set_file(true); + ASSERT_TRUE(logging::initLogging(config, "grpc_log_test", directory)); + + ErrorLoggingTestService service; + grpc::ServerBuilder builder; + const std::string socket_path = + "/tmp/cmvr_grpc_error_logging_file_test_" + + std::to_string(static_cast(::getpid())) + "_" + + std::to_string(g_socket_sequence.fetch_add(1U)) + ".sock"; + cleanup.socket_path = socket_path; + ScopedSocket socket(socket_path); + const std::string address = "unix:" + socket_path; + builder.AddListeningPort(address, grpc::InsecureServerCredentials()); + builder.RegisterService(&service); + std::vector> factories; + factories.emplace_back(makeGrpcErrorLoggingInterceptorFactory()); + builder.experimental().SetInterceptorCreators(std::move(factories)); + ScopedServer server(builder.BuildAndStart()); + ASSERT_NE(server.get(), nullptr); + + const auto channel = grpc::CreateChannel( + address, grpc::InsecureChannelCredentials()); + auto stub = api::TestService::NewStub(channel); + grpc::ClientContext context; + setRpcDeadline(context); + api::TestReqeust request; + api::TestResponse response; + request.set_data("cancelled"); + const grpc::Status status = stub->Call(&context, request, &response); + ASSERT_EQ(status.error_code(), grpc::StatusCode::CANCELLED); + + std::ifstream input(directory / "grpc_log_test.log"); + const std::string contents{ + std::istreambuf_iterator(input), + std::istreambuf_iterator()}; + EXPECT_NE(contents.find("/cmvr.api.TestService/Call"), std::string::npos); + EXPECT_NE(contents.find("code=1(CANCELLED)"), std::string::npos); + EXPECT_NE(contents.find("client closed stream"), std::string::npos); + +} + +TEST_F(GrpcErrorLoggingInterceptorTest, RecordsServerCancellationReturnedAsOk) +{ + grpc::ClientContext context; + setRpcDeadline(context); + api::TestReqeust request; + api::TestResponse response; + request.set_data("server-cancel"); + + const grpc::Status status = test_stub_->Call(&context, request, &response); + ASSERT_EQ(status.error_code(), grpc::StatusCode::CANCELLED); + + // TryCancel is best effort: the final status may be intercepted before + // the cancellation hook. Either way, the interceptor must not recurse or + // emit duplicate records. + const auto captured = waitForRecords(1U, std::chrono::milliseconds(50)); + ASSERT_LE(captured.size(), 1U); + if (!captured.empty()) { + EXPECT_EQ(captured.front().kind, GrpcFailureKind::GRPC_STATUS); + EXPECT_EQ(captured.front().status_code, grpc::StatusCode::CANCELLED); + EXPECT_EQ(captured.front().severity, GrpcFailureSeverity::WARNING); + EXPECT_EQ( + captured.front().detail, + "RPC cancellation was requested by the edge server before completion"); + } +} + +TEST_F(GrpcErrorLoggingInterceptorTest, RecordsClientCancellation) +{ + grpc::ClientContext context; + setRpcDeadline(context); + api::TestReqeust request; + api::TestResponse response; + request.set_data("client-cancel"); + + grpc::Status status; + std::thread call([&] { + status = test_stub_->Call(&context, request, &response); + }); + const bool call_started = test_service_.waitForCancelableCall(); + context.TryCancel(); + call.join(); + + ASSERT_TRUE(call_started); + ASSERT_EQ(status.error_code(), grpc::StatusCode::CANCELLED); + const auto captured = waitForRecords(1U); + ASSERT_EQ(captured.size(), 1U); + EXPECT_EQ(captured.front().kind, GrpcFailureKind::GRPC_STATUS); + EXPECT_EQ(captured.front().status_code, grpc::StatusCode::CANCELLED); + EXPECT_EQ(captured.front().severity, GrpcFailureSeverity::WARNING); + EXPECT_EQ( + captured.front().detail, + "RPC was cancelled by the client or transport before completion"); +} + +TEST_F(GrpcErrorLoggingInterceptorTest, + RecordsStatusWhenNonOkSuppressesUnaryResponse) +{ + grpc::ClientContext context; + setRpcDeadline(context); + api::CommandHeader_Request request; + api::CommandHeader_Feedback response; + request.set_device_id("missing"); + + const grpc::Status status = + agv_stub_->emergencyStop(&context, request, &response); + ASSERT_EQ(status.error_code(), grpc::StatusCode::NOT_FOUND); + + const auto captured = waitForRecords(1U); + ASSERT_EQ(captured.size(), 1U); + EXPECT_EQ(captured.front().kind, GrpcFailureKind::GRPC_STATUS); + EXPECT_EQ(captured.front().status_code, grpc::StatusCode::NOT_FOUND); + EXPECT_EQ(captured.front().severity, GrpcFailureSeverity::WARNING); + EXPECT_EQ(captured.front().detail, "missing sentinel"); +} + +TEST_F(GrpcErrorLoggingInterceptorTest, RecordsNestedFeedbackFailure) +{ + grpc::ClientContext context; + setRpcDeadline(context); + api::StartCameraCommand_Request request; + api::StartCameraCommand_Feedback response; + + ASSERT_TRUE(camera_stub_->StartCamera(&context, request, &response).ok()); + ASSERT_FALSE(response.header().success()); + + const auto captured = waitForRecords(1U); + ASSERT_EQ(captured.size(), 1U); + EXPECT_EQ(captured.front().kind, GrpcFailureKind::APPLICATION); + EXPECT_EQ(captured.front().method, "/cmvr.api.CameraService/StartCamera"); + EXPECT_EQ(captured.front().detail, "nested sentinel"); +} + +TEST_F(GrpcErrorLoggingInterceptorTest, + RecordsStreamingResponseAndFinalNonOkStatus) +{ + grpc::ClientContext context; + setRpcDeadline(context); + api::AgvMapStreamCommand_Request request; + auto reader = agv_stub_->streamMap(&context, request); + api::AgvMapStreamCommand_Feedback response; + + ASSERT_TRUE(reader->Read(&response)); + ASSERT_FALSE(response.header().success()); + const grpc::Status status = reader->Finish(); + ASSERT_EQ(status.error_code(), grpc::StatusCode::INTERNAL); + + const auto captured = waitForRecords(2U); + ASSERT_EQ(captured.size(), 2U); + EXPECT_EQ(captured[0].kind, GrpcFailureKind::APPLICATION); + EXPECT_EQ(captured[0].method, "/cmvr.api.AgvService/streamMap"); + EXPECT_EQ(captured[0].status_code, grpc::StatusCode::OK); + EXPECT_EQ(captured[0].severity, GrpcFailureSeverity::ERROR); + EXPECT_EQ(captured[0].detail, "stream sentinel"); + EXPECT_EQ(captured[1].kind, GrpcFailureKind::GRPC_STATUS); + EXPECT_EQ(captured[1].status_code, grpc::StatusCode::INTERNAL); + EXPECT_EQ(captured[1].severity, GrpcFailureSeverity::ERROR); + EXPECT_EQ(captured[1].detail, "stream sentinel"); +} + +TEST_F(GrpcErrorLoggingInterceptorTest, + RecordsStreamingFailureBeforeTheStreamFinishes) +{ + grpc::ClientContext context; + setRpcDeadline(context); + api::AgvMapStreamCommand_Request request; + request.set_resume_token("hold"); + auto reader = agv_stub_->streamMap(&context, request); + api::AgvMapStreamCommand_Feedback response; + + const bool read = reader->Read(&response); + const auto captured_before_finish = waitForRecords(1U); + agv_service_.releaseHeldStream(); + const grpc::Status status = reader->Finish(); + + ASSERT_TRUE(read); + ASSERT_TRUE(status.ok()); + ASSERT_EQ(captured_before_finish.size(), 1U); + EXPECT_EQ( + captured_before_finish.front().detail, + "immediate stream sentinel"); +} + +TEST_F(GrpcErrorLoggingInterceptorTest, + RecordsEveryStreamingFailureResponse) +{ + grpc::ClientContext context; + setRpcDeadline(context); + api::AgvMapStreamCommand_Request request; + request.set_resume_token("multiple"); + auto reader = agv_stub_->streamMap(&context, request); + api::AgvMapStreamCommand_Feedback response; + + ASSERT_TRUE(reader->Read(&response)); + ASSERT_TRUE(reader->Read(&response)); + ASSERT_TRUE(reader->Finish().ok()); + + const auto captured = waitForRecords(2U); + ASSERT_EQ(captured.size(), 2U); + EXPECT_EQ(captured[0].detail, "first stream sentinel"); + EXPECT_EQ(captured[1].detail, "second stream sentinel"); +} + +TEST_F(GrpcErrorLoggingInterceptorTest, + PreservesDifferentResponseAndStatusErrors) +{ + grpc::ClientContext context; + setRpcDeadline(context); + api::AgvMapStreamCommand_Request request; + request.set_resume_token("different"); + auto reader = agv_stub_->streamMap(&context, request); + api::AgvMapStreamCommand_Feedback response; + + ASSERT_TRUE(reader->Read(&response)); + const grpc::Status status = reader->Finish(); + ASSERT_EQ(status.error_code(), grpc::StatusCode::INTERNAL); + + const auto captured = waitForRecords(2U); + ASSERT_EQ(captured.size(), 2U); + EXPECT_EQ(captured[0].kind, GrpcFailureKind::APPLICATION); + EXPECT_EQ(captured[0].detail, "application stream sentinel"); + EXPECT_EQ(captured[1].kind, GrpcFailureKind::GRPC_STATUS); + EXPECT_EQ(captured[1].status_code, grpc::StatusCode::INTERNAL); + EXPECT_EQ(captured[1].detail, "status stream sentinel"); +} + +TEST_F(GrpcErrorLoggingInterceptorTest, + DoesNotDeduplicateLongErrorsByTruncatedPrefix) +{ + grpc::ClientContext context; + setRpcDeadline(context); + api::AgvMapStreamCommand_Request request; + request.set_resume_token("long-different"); + auto reader = agv_stub_->streamMap(&context, request); + api::AgvMapStreamCommand_Feedback response; + + ASSERT_TRUE(reader->Read(&response)); + const grpc::Status status = reader->Finish(); + ASSERT_EQ(status.error_code(), grpc::StatusCode::INTERNAL); + + const auto captured = waitForRecords(2U); + ASSERT_EQ(captured.size(), 2U); + EXPECT_EQ(captured[0].kind, GrpcFailureKind::APPLICATION); + EXPECT_EQ(captured[1].kind, GrpcFailureKind::GRPC_STATUS); +} + +TEST(GrpcErrorLoggingInterceptorWireTest, RequiresHeaderAsFirstField) +{ + api::CommandHeader_Feedback header; + header.set_success(false); + header.set_error_message("out-of-order sentinel"); + std::string payload; + payload.push_back(static_cast(0x12)); + payload.push_back(static_cast(0x01)); + payload.push_back('x'); + payload.push_back(static_cast(0x0a)); + payload.push_back(static_cast(header.ByteSizeLong())); + payload.append(header.SerializeAsString()); + + const auto result = callRawStartCamera(payload); + + ASSERT_TRUE(result.status.ok()); + ASSERT_EQ(result.response, payload); + ASSERT_EQ(result.records.size(), 1U); + EXPECT_EQ( + result.records.front().kind, + GrpcFailureKind::MALFORMED_RESPONSE); +} + +TEST(GrpcErrorLoggingInterceptorWireTest, DoesNotScanPayloadAfterHeader) +{ + api::CommandHeader_Feedback header; + header.set_success(false); + header.set_error_message("leading header sentinel"); + std::string payload; + payload.push_back(static_cast(0x0a)); + payload.push_back(static_cast(header.ByteSizeLong())); + payload.append(header.SerializeAsString()); + payload.push_back(static_cast(0x12)); + payload.push_back(static_cast(0xff)); + payload.append(2U * 1024U * 1024U, 'x'); + + const auto result = callRawStartCamera(payload); + + ASSERT_TRUE(result.status.ok()); + ASSERT_EQ(result.records.size(), 1U); + EXPECT_EQ(result.records.front().kind, GrpcFailureKind::APPLICATION); + EXPECT_EQ(result.records.front().detail, "leading header sentinel"); +} + +TEST(GrpcErrorLoggingInterceptorWireTest, TreatsExplicitEmptyHeaderAsFailure) +{ + const std::string payload{ + static_cast(0x0a), static_cast(0x00)}; + + const auto result = callRawStartCamera(payload); + + ASSERT_TRUE(result.status.ok()); + ASSERT_EQ(result.records.size(), 1U); + EXPECT_EQ(result.records.front().kind, GrpcFailureKind::APPLICATION); + EXPECT_EQ( + result.records.front().detail, + "operation failed without an error message"); +} + +TEST(GrpcErrorLoggingInterceptorWireTest, RejectsInvalidHeaderWireType) +{ + const std::string payload{ + static_cast(0x08), static_cast(0x01)}; + + const auto result = callRawStartCamera(payload); + + ASSERT_TRUE(result.status.ok()); + ASSERT_EQ(result.records.size(), 1U); + EXPECT_EQ( + result.records.front().kind, + GrpcFailureKind::MALFORMED_RESPONSE); +} + +TEST(GrpcErrorLoggingInterceptorWireTest, RejectsInvalidUtf8ErrorMessage) +{ + const std::string payload{ + static_cast(0x0a), static_cast(0x04), + static_cast(0x12), static_cast(0x02), + static_cast(0xc3), static_cast(0x28)}; + + const auto result = callRawStartCamera(payload); + + ASSERT_TRUE(result.status.ok()); + ASSERT_EQ(result.records.size(), 1U); + EXPECT_EQ( + result.records.front().kind, + GrpcFailureKind::MALFORMED_RESPONSE); +} + +TEST_F(GrpcErrorLoggingInterceptorTest, DoesNotRecordSuccessfulLargeResponse) +{ + grpc::ClientContext context; + setRpcDeadline(context); + api::GetRGBImageCommand_Request request; + api::GetRGBImageCommand_Feedback response; + + ASSERT_TRUE(camera_stub_->GetRGBImage(&context, request, &response).ok()); + ASSERT_TRUE(response.header().success()); + ASSERT_EQ(response.color_frame().data().size(), 2U * 1024U * 1024U); + EXPECT_TRUE(records().empty()); +} + +TEST_F(GrpcErrorLoggingInterceptorTest, IgnoresSuccessfulResponseWithoutHeader) +{ + grpc::ClientContext context; + setRpcDeadline(context); + api::TestReqeust request; + api::TestResponse response; + request.set_data("ok"); + + ASSERT_TRUE(test_stub_->Call(&context, request, &response).ok()); + EXPECT_EQ(response.data(), "ok"); + EXPECT_TRUE(records().empty()); +} + +} // namespace +} // namespace cmvr::service diff --git a/cmvr-es/task/grpc_server_task/src/grpc_server_task.cpp b/cmvr-es/task/grpc_server_task/src/grpc_server_task.cpp index 34325bd9..c7e46a9f 100644 --- a/cmvr-es/task/grpc_server_task/src/grpc_server_task.cpp +++ b/cmvr-es/task/grpc_server_task/src/grpc_server_task.cpp @@ -16,6 +16,7 @@ #include "service/grpc/include/grpc_robot_arm_teleop_backend.h" #include "service/grpc/include/grpc_camera_service.h" #include "service/grpc/include/grpc_dexhand_service.h" +#include "service/grpc/include/grpc_error_logging_interceptor.h" #include "service/grpc/include/grpc_head_service.h" #include "service/grpc/include/grpc_hlc_service.h" #include "service/grpc/include/grpc_microphone_service.h" @@ -120,6 +121,13 @@ bool GrpcServerTask::start() grpc::ServerBuilder builder; builder.AddListeningPort(local_address, grpc::InsecureServerCredentials()); + std::vector> + interceptor_factories; + interceptor_factories.emplace_back( + service::makeGrpcErrorLoggingInterceptorFactory()); + builder.experimental().SetInterceptorCreators( + std::move(interceptor_factories)); builder.RegisterService(camera_service_.get()); builder.RegisterService(system_service_.get()); builder.RegisterService(speaker_service_.get());