feat(grpc): persist RPC failures in edge logs

This commit is contained in:
xtkuang 2026-08-14 08:38:00 +08:00
parent 4c8b320b4b
commit de764de607
15 changed files with 1570 additions and 12 deletions

View File

@ -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()

View File

@ -167,6 +167,15 @@ void Logger::shutdown()
initialized_ = false;
}
void Logger::flush()
{
std::lock_guard<std::mutex> 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<std::mutex> 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;
}

View File

@ -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);

View File

@ -0,0 +1,93 @@
#include "common/base/logging/logger.h"
#include <algorithm>
#include <array>
#include <filesystem>
#include <fstream>
#include <iterator>
#include <string>
#include <gtest/gtest.h>
#include <unistd.h>
namespace cmvr::logging {
namespace {
class LoggerTest : public testing::Test {
protected:
void SetUp() override
{
std::array<char, 64> 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<char>(input),
std::istreambuf_iterator<char>()};
}
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

View File

@ -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"

View File

@ -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
)

View File

@ -0,0 +1,37 @@
#pragma once
#include <functional>
#include <memory>
#include <string>
#include <grpcpp/support/server_interceptor.h>
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<void(const GrpcFailureRecord&)>;
std::unique_ptr<grpc::experimental::ServerInterceptorFactoryInterface>
makeGrpcErrorLoggingInterceptorFactory(GrpcFailureSink sink = {});
} // namespace cmvr::service

View File

@ -21,7 +21,6 @@ using namespace cmvr::device;
namespace {
template <typename ResponseT>
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());

View File

@ -120,7 +120,6 @@ bool applyFreedomValues(const FreedomCollection& freedoms,
template <typename ResponseT>
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());

View File

@ -0,0 +1,504 @@
#include "service/grpc/include/grpc_error_logging_interceptor.h"
#include <atomic>
#include <chrono>
#include <cstdint>
#include <limits>
#include <mutex>
#include <optional>
#include <sstream>
#include <string_view>
#include <utility>
#include <google/protobuf/descriptor.h>
#include <google/protobuf/io/coded_stream.h>
#include <google/protobuf/wire_format_lite.h>
#include <grpcpp/server_context.h>
#include <grpcpp/support/interceptor.h>
#include <grpcpp/support/proto_buffer_reader.h>
#include <grpcpp/support/status.h>
#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<unsigned char>(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<int>(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<std::uint32_t>(kMaxFeedbackBytes) ||
length > static_cast<std::uint32_t>(std::numeric_limits<int>::max())) {
return false;
}
std::string payload;
if (!input.ReadString(&payload, static_cast<int>(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<GrpcFailureRecord> 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<GrpcFailureRecord> pending_failure_;
std::mutex sink_mutex_;
std::atomic<bool> server_cancel_requested_{false};
std::atomic<bool> 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()
: "<unknown>";
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<grpc::experimental::ServerInterceptorFactoryInterface>
makeGrpcErrorLoggingInterceptorFactory(GrpcFailureSink sink)
{
return std::make_unique<GrpcErrorLoggingInterceptorFactory>(
std::move(sink));
}
} // namespace cmvr::service

View File

@ -20,7 +20,6 @@ using namespace cmvr::api;
namespace {
template <typename ResponseT>
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<AbstractBiohead>(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());

View File

@ -18,7 +18,6 @@ using namespace cmvr::service;
namespace {
template <typename ResponseT>
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());

View File

@ -14,7 +14,6 @@ using namespace cmvr::service;
namespace {
template <typename ResponseT>
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());

View File

@ -0,0 +1,864 @@
#include "service/grpc/include/grpc_error_logging_interceptor.h"
#include <atomic>
#include <algorithm>
#include <array>
#include <chrono>
#include <condition_variable>
#include <cstdio>
#include <filesystem>
#include <fstream>
#include <iterator>
#include <memory>
#include <mutex>
#include <string>
#include <thread>
#include <utility>
#include <vector>
#include <grpcpp/grpcpp.h>
#include <grpcpp/generic/generic_stub.h>
#include <grpcpp/impl/client_unary_call.h>
#include <grpcpp/impl/rpc_method.h>
#include <gtest/gtest.h>
#include <unistd.h>
#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<std::uint64_t> 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<grpc::Slice> slices;
EXPECT_TRUE(buffer.Dump(&slices).ok());
std::string payload;
for (const auto& slice : slices) {
payload.append(
reinterpret_cast<const char*>(slice.begin()), slice.size());
}
return payload;
}
grpc::Status rawUnaryCall(const std::shared_ptr<grpc::Channel>& 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<api::AgvMapStreamCommand_Feedback>* 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<GrpcFailureRecord> records;
std::string response;
};
class ScopedServer final {
public:
ScopedServer(std::unique_ptr<grpc::Server> 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<grpc::Server> 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<GrpcFailureRecord> records;
RawCameraService service(response_payload);
grpc::ServerBuilder builder;
const std::string address =
"unix:/tmp/cmvr_grpc_error_logging_raw_" +
std::to_string(static_cast<long long>(::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<std::unique_ptr<
grpc::experimental::ServerInterceptorFactoryInterface>> 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<long long>(::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<std::unique_ptr<
grpc::experimental::ServerInterceptorFactoryInterface>> 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<GrpcFailureRecord> records()
{
std::lock_guard lock(records_mutex_);
return records_;
}
std::vector<GrpcFailureRecord> 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<grpc::Server> server_;
std::unique_ptr<api::TestService::Stub> test_stub_;
std::unique_ptr<api::AgvService::Stub> agv_stub_;
std::unique_ptr<api::CameraService::Stub> camera_stub_;
std::string socket_path_;
std::mutex records_mutex_;
std::condition_variable records_changed_;
std::vector<GrpcFailureRecord> 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<char, 80> 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<long long>(::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<std::unique_ptr<
grpc::experimental::ServerInterceptorFactoryInterface>> 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<char>(input),
std::istreambuf_iterator<char>()};
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<char>(0x12));
payload.push_back(static_cast<char>(0x01));
payload.push_back('x');
payload.push_back(static_cast<char>(0x0a));
payload.push_back(static_cast<char>(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<char>(0x0a));
payload.push_back(static_cast<char>(header.ByteSizeLong()));
payload.append(header.SerializeAsString());
payload.push_back(static_cast<char>(0x12));
payload.push_back(static_cast<char>(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<char>(0x0a), static_cast<char>(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<char>(0x08), static_cast<char>(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<char>(0x0a), static_cast<char>(0x04),
static_cast<char>(0x12), static_cast<char>(0x02),
static_cast<char>(0xc3), static_cast<char>(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

View File

@ -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<std::unique_ptr<
grpc::experimental::ServerInterceptorFactoryInterface>>
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());