cmvr-es/cmvr-es/task/grpc_server_task/src/grpc_server_task.cpp

267 lines
7.4 KiB
C++

#include "task/grpc_server_task/include/grpc_server_task.h"
#include <exception>
#include <grpcpp/ext/proto_server_reflection_plugin.h>
#include "cmvr/config/grpc_server_config/grpc_server_config.pb.h"
#include "cmvr/config/task_manager_config/task_manager_config.pb.h"
#include "common/base/logging/logger.h"
#include "common/config/config_files.h"
#include "service/grpc/include/grpc_arm_service.h"
#include "service/grpc/include/grpc_camera_service.h"
#include "service/grpc/include/grpc_dexhand_service.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"
#include "service/grpc/include/grpc_speaker_service.h"
#include "service/grpc/include/grpc_system_service.h"
#include "task/task_factory.h"
namespace cmvr::task {
namespace {
std::shared_ptr<Task> createGrpcServerTask(const config::TaskConfigEntry& entry)
{
if (entry.id().empty()) {
CMVR_LOG(ERROR) << "[GrpcServerTask] Empty gRPC task ID";
return nullptr;
}
if (entry.config_file().empty()) {
CMVR_LOG(ERROR) << "[GrpcServerTask] Empty gRPC config_file for task ID: " << entry.id();
return nullptr;
}
config::GRPCServerRootConfig root_cfg;
if (!ConfigHelper::loadConfigFile(entry.config_file(), root_cfg)) {
CMVR_LOG(ERROR) << "[GrpcServerTask] Failed to load gRPC config: " << entry.config_file();
return nullptr;
}
const auto& cfg = root_cfg.grpc_server();
if (cfg.id().empty()) {
CMVR_LOG(ERROR) << "[GrpcServerTask] gRPC config missing id: " << entry.config_file();
return nullptr;
}
if (cfg.id() != entry.id()) {
CMVR_LOG(ERROR) << "[GrpcServerTask] gRPC task ID mismatch: manager id=" << entry.id()
<< ", config id=" << cfg.id();
return nullptr;
}
return std::make_shared<GrpcServerTask>(cfg);
}
} // namespace
GrpcServerTask::GrpcServerTask(const config::GRPCServerConfig& cfg)
: cfg_(cfg),
id_(cfg.id())
{
state_ = TaskState::UNINITIALIZED;
}
GrpcServerTask::~GrpcServerTask()
{
stop();
}
bool GrpcServerTask::start()
{
std::lock_guard lock(mutex_);
if (state_ == TaskState::RUNNING) {
return true;
}
if (state_ != TaskState::IDLE && state_ != TaskState::STOPPED) {
last_error_ = "gRPC task is not initialized";
state_ = TaskState::FAILED;
return false;
}
if (id_.empty()) {
last_error_ = "gRPC task id is empty";
state_ = TaskState::FAILED;
return false;
}
if (cfg_.enable_reflection()) {
static std::once_flag reflection_once;
std::call_once(reflection_once, []() {
grpc::reflection::InitProtoReflectionServerBuilderPlugin();
});
}
const std::string host = cfg_.host().empty() ? "0.0.0.0" : cfg_.host();
const std::string port = cfg_.port().empty() ? "50051" : cfg_.port();
const std::string local_address = host + ":" + port;
camera_service_ = std::make_unique<service::gRPCCameraServiceImpl>();
system_service_ = std::make_unique<service::gRPCSystemServiceImpl>();
speaker_service_ = std::make_unique<service::gRPCSpeakerServiceImpl>();
microphone_service_ = std::make_unique<service::gRPCMicroPhoneServiceImpl>();
dexhand_service_ = std::make_unique<service::gRPCDexHandServiceImpl>();
biohand_service_ = std::make_unique<service::gRPCMBioHeadServiceImpl>();
arm_service_ = std::make_unique<service::gRPCArmServiceImpl>();
hlc_service_ = std::make_unique<service::gRPCHlcServiceImpl>();
grpc::ServerBuilder builder;
builder.AddListeningPort(local_address, grpc::InsecureServerCredentials());
builder.RegisterService(camera_service_.get());
builder.RegisterService(system_service_.get());
builder.RegisterService(speaker_service_.get());
builder.RegisterService(microphone_service_.get());
builder.RegisterService(dexhand_service_.get());
builder.RegisterService(biohand_service_.get());
builder.RegisterService(arm_service_.get());
builder.RegisterService(hlc_service_.get());
server_ = builder.BuildAndStart();
if (!server_) {
clearServices();
last_error_ = "Failed to build and start gRPC server at " + local_address;
CMVR_LOG(ERROR) << "[GrpcServerTask] " << last_error_;
state_ = TaskState::FAILED;
return false;
}
address_ = local_address;
state_ = TaskState::RUNNING;
CMVR_LOG(INFO) << "[GrpcServerTask] gRPC server started, address=" << address_;
wait_thread_ = std::thread(&GrpcServerTask::waitLoop, this);
return true;
}
bool GrpcServerTask::init()
{
std::lock_guard lock(mutex_);
if (id_.empty()) {
last_error_ = "gRPC task id is empty";
state_ = TaskState::FAILED;
return false;
}
if (cfg_.port().empty()) {
last_error_ = "gRPC task port is empty";
state_ = TaskState::FAILED;
return false;
}
last_error_.clear();
state_ = TaskState::IDLE;
return true;
}
bool GrpcServerTask::step(const double dt)
{
(void)dt;
return true;
}
void GrpcServerTask::stop()
{
{
std::lock_guard lock(mutex_);
if (server_) {
server_->Shutdown();
}
}
if (wait_thread_.joinable()) {
wait_thread_.join();
}
std::lock_guard lock(mutex_);
server_.reset();
clearServices();
if (state_ == TaskState::RUNNING ||
state_ == TaskState::IDLE ||
state_ == TaskState::UNINITIALIZED) {
state_ = TaskState::STOPPED;
}
}
TaskState GrpcServerTask::state() const
{
std::lock_guard lock(mutex_);
return state_;
}
bool GrpcServerTask::isBusy() const
{
return state() == TaskState::RUNNING;
}
bool GrpcServerTask::isFinished() const
{
return state() == TaskState::STOPPED;
}
bool GrpcServerTask::isFailed() const
{
return state() == TaskState::FAILED;
}
std::string GrpcServerTask::stateString() const
{
return taskStateToString(state());
}
std::string GrpcServerTask::detailStatusString() const
{
std::lock_guard lock(mutex_);
if (last_error_.empty()) {
return std::string(taskStateToString(state_)) + " " + address_;
}
return std::string(taskStateToString(state_)) + " " + last_error_;
}
std::string GrpcServerTask::address() const
{
std::lock_guard lock(mutex_);
return address_;
}
void GrpcServerTask::waitLoop()
{
try {
grpc::Server* server = nullptr;
{
std::lock_guard lock(mutex_);
server = server_.get();
}
if (server) {
server->Wait();
}
std::lock_guard lock(mutex_);
if (state_ == TaskState::RUNNING) {
state_ = TaskState::STOPPED;
}
} catch (const std::exception& e) {
std::lock_guard lock(mutex_);
last_error_ = e.what();
state_ = TaskState::FAILED;
} catch (...) {
std::lock_guard lock(mutex_);
last_error_ = "Unknown gRPC server error";
state_ = TaskState::FAILED;
}
}
void GrpcServerTask::clearServices()
{
hlc_service_.reset();
arm_service_.reset();
biohand_service_.reset();
dexhand_service_.reset();
microphone_service_.reset();
speaker_service_.reset();
system_service_.reset();
camera_service_.reset();
}
void registerGrpcServerTaskFactory()
{
TaskFactory::registerCreator(
config::TaskConfigEntry::TASK_TYPE_GRPC_SERVER,
createGrpcServerTask);
}
} // namespace cmvr::task