From ffef60e9b80817551a548ecb1ce30120ec03b53b Mon Sep 17 00:00:00 2001 From: linbo <1034003879@qq.com> Date: Tue, 25 Aug 2026 17:37:53 +0800 Subject: [PATCH] feat(config): add standalone remote configuration service --- CMakeLists.txt | 11 + cmvr-es/config_server/config_file_service.cpp | 461 ++++++++++++++++++ cmvr-es/config_server/config_file_service.h | 45 ++ cmvr-es/config_server/config_server_main.cpp | 82 ++++ protos/cmvr/api/configuration_command.proto | 57 +++ protos/cmvr/api/configuration_service.proto | 12 + 6 files changed, 668 insertions(+) create mode 100644 cmvr-es/config_server/config_file_service.cpp create mode 100644 cmvr-es/config_server/config_file_service.h create mode 100644 cmvr-es/config_server/config_server_main.cpp create mode 100644 protos/cmvr/api/configuration_command.proto create mode 100644 protos/cmvr/api/configuration_service.proto diff --git a/CMakeLists.txt b/CMakeLists.txt index 2d5a7b21..a1cb2158 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -157,6 +157,17 @@ target_link_libraries(cmvr_es PRIVATE ) install(TARGETS cmvr_es RUNTIME DESTINATION bin) + +add_executable(cmvr_config_server + cmvr-es/config_server/config_server_main.cpp + cmvr-es/config_server/config_file_service.cpp +) +target_link_libraries(cmvr_config_server PRIVATE + cmvr_es::proto + protobuf::libprotobuf + gRPC::grpc++ +) +install(TARGETS cmvr_config_server RUNTIME DESTINATION bin) install(CODE [[ file(REMOVE_RECURSE "${CMAKE_INSTALL_PREFIX}/bin/config" diff --git a/cmvr-es/config_server/config_file_service.cpp b/cmvr-es/config_server/config_file_service.cpp new file mode 100644 index 00000000..b5498393 --- /dev/null +++ b/cmvr-es/config_server/config_file_service.cpp @@ -0,0 +1,461 @@ +#include "config_server/config_file_service.h" + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include "cmvr/config/agv_config/agv_config.pb.h" +#include "cmvr/config/arm_config/arm_config.pb.h" +#include "cmvr/config/biohead_config/biohead_config.pb.h" +#include "cmvr/config/camera_config/camera_config.pb.h" +#include "cmvr/config/cmvr_es_config/cmvr_es_config.pb.h" +#include "cmvr/config/device_manager_config/device_manager_config.pb.h" +#include "cmvr/config/dexhand_config/dexhand_config.pb.h" +#include "cmvr/config/grpc_server_config/grpc_server_config.pb.h" +#include "cmvr/config/logger_config/logger_config.pb.h" +#include "cmvr/config/microphone_config/microphone_config.pb.h" +#include "cmvr/config/motor_config/motor_config.pb.h" +#include "cmvr/config/mujoco_config/mujoco_world_config.pb.h" +#include "cmvr/config/quic_edge_config/quic_edge_config.pb.h" +#include "cmvr/config/self_collision_task_config/self_collision_task_config.pb.h" +#include "cmvr/config/speaker_config/speaker_config.pb.h" +#include "cmvr/config/task_manager_config/task_manager_config.pb.h" +#include "cmvr/config/touch_screen_task_config/touch_screen_task_config.pb.h" +#include "google/protobuf/io/tokenizer.h" +#include "google/protobuf/text_format.h" + +namespace cmvr::config_server { +namespace { + +constexpr std::uintmax_t kMaximumConfigBytes = 4U * 1024U * 1024U; + +void setTimestamp(google::protobuf::Timestamp* timestamp) +{ + if (!timestamp) return; + const auto now = std::chrono::system_clock::now().time_since_epoch(); + const auto seconds = std::chrono::duration_cast(now); + const auto nanos = std::chrono::duration_cast(now - seconds); + timestamp->set_seconds(seconds.count()); + timestamp->set_nanos(static_cast(nanos.count())); +} + +template +grpc::Status fail(Response* response, const std::string& message) +{ + response->mutable_header()->set_success(false); + response->mutable_header()->set_error_message(message); + setTimestamp(response->mutable_header()->mutable_timestamp()); + return grpc::Status::OK; +} + +template +void succeed(Response* response) +{ + response->mutable_header()->set_success(true); + setTimestamp(response->mutable_header()->mutable_timestamp()); +} + +bool hasConfigExtension(const std::filesystem::path& path) +{ + const std::string name = path.filename().string(); + constexpr std::string_view suffix = ".pb.txt"; + return name.size() > suffix.size() && + name.compare(name.size() - suffix.size(), suffix.size(), suffix) == 0; +} + +bool startsWith(const std::string& value, const std::string& prefix) +{ + return value.compare(0, prefix.size(), prefix) == 0; +} + +const google::protobuf::Descriptor* descriptorForPath(const std::string& path) +{ + using namespace cmvr::config; + if (path == "cmvr_es.pb.txt") return CMVRESRootConfig::descriptor(); + if (path == "logger/logger.pb.txt") return LoggerRootConfig::descriptor(); + if (path == "manager/device_manager.pb.txt") return DeviceManagerRootConfig::descriptor(); + if (path == "manager/task_manager.pb.txt") return TaskManagerRootConfig::descriptor(); + if (startsWith(path, "devices/agv/")) return AGVRootConfig::descriptor(); + if (startsWith(path, "devices/arm/")) return ArmRootConfig::descriptor(); + if (startsWith(path, "devices/biohead/")) return BioHeadRobotRootConfig::descriptor(); + if (startsWith(path, "devices/camera/")) return CameraRootConfig::descriptor(); + if (startsWith(path, "devices/dexhand/")) return DexHandRootConfig::descriptor(); + if (startsWith(path, "devices/microphone/")) return MicroPhoneRootConfig::descriptor(); + if (startsWith(path, "devices/motor/")) return MotorRootConfig::descriptor(); + if (startsWith(path, "devices/mujoco/") && path.find("viewer") != std::string::npos) { + return MujocoViewerRootConfig::descriptor(); + } + if (startsWith(path, "devices/mujoco/")) return MujocoWorldRootConfig::descriptor(); + if (startsWith(path, "devices/speaker/")) return SpeakerRootConfig::descriptor(); + if (startsWith(path, "tasks/grpc_server_task/")) return GRPCServerRootConfig::descriptor(); + if (startsWith(path, "tasks/quic_edge_task/")) return QuicEdgeRootConfig::descriptor(); + if (startsWith(path, "tasks/self_collision_task/")) { + return SelfCollisionTaskRootConfig::descriptor(); + } + if (startsWith(path, "tasks/touch_screen_task/")) { + return TouchScreenTaskRootConfig::descriptor(); + } + return nullptr; +} + +class TextErrorCollector final : public google::protobuf::io::ErrorCollector { +public: + void AddError(int line, int column, const std::string& message) override + { + if (!errors_.empty()) errors_ += "; "; + errors_ += "line " + std::to_string(line + 1) + ", column " + + std::to_string(column + 1) + ": " + message; + } + + void AddWarning(int, int, const std::string&) override {} + + const std::string& errors() const noexcept { return errors_; } + +private: + std::string errors_; +}; + +bool validateTextProto( + const std::string& relative_path, + const std::string& content, + std::string* message_type, + std::string* error) +{ + const auto* descriptor = descriptorForPath(relative_path); + if (!descriptor) { + if (error) *error = "No protobuf schema is registered for: " + relative_path; + return false; + } + const auto* prototype = + google::protobuf::MessageFactory::generated_factory()->GetPrototype(descriptor); + if (!prototype) { + if (error) *error = "Cannot create protobuf message for: " + relative_path; + return false; + } + + std::unique_ptr message(prototype->New()); + TextErrorCollector collector; + google::protobuf::TextFormat::Parser parser; + parser.RecordErrorsTo(&collector); + parser.AllowUnknownField(false); + if (!parser.ParseFromString(content, message.get())) { + if (error) { + *error = collector.errors().empty() + ? "Invalid protobuf text in: " + relative_path + : collector.errors(); + } + return false; + } + if (message_type) *message_type = descriptor->full_name(); + return true; +} + +bool readFile(const std::filesystem::path& path, std::string* content, std::string* error) +{ + std::error_code ec; + const auto size = std::filesystem::file_size(path, ec); + if (ec) { + if (error) *error = "Cannot inspect config file: " + ec.message(); + return false; + } + if (size > kMaximumConfigBytes) { + if (error) *error = "Config file exceeds 4 MiB limit"; + return false; + } + std::ifstream input(path, std::ios::binary); + if (!input) { + if (error) *error = "Cannot open config file for reading"; + return false; + } + content->assign(std::istreambuf_iterator(input), std::istreambuf_iterator()); + if (!input.eof() && input.fail()) { + if (error) *error = "Failed while reading config file"; + return false; + } + return true; +} + +std::string revisionFor(const std::string& content) +{ + std::uint64_t hash = 14695981039346656037ULL; + for (const unsigned char byte : content) { + hash ^= byte; + hash *= 1099511628211ULL; + } + std::ostringstream output; + output << std::hex << std::setfill('0') << std::setw(16) << hash; + return output.str(); +} + +int64_t modifiedUnixMs(const std::filesystem::path& path) +{ + struct stat info {}; + if (::stat(path.c_str(), &info) != 0) return 0; + return static_cast(info.st_mtim.tv_sec) * 1000LL + + static_cast(info.st_mtim.tv_nsec / 1000000L); +} + +void fillFileInfo( + const std::filesystem::path& path, + const std::string& relative_path, + api::ConfigFileInfo* info) +{ + std::error_code ec; + info->set_relative_path(relative_path); + info->set_size_bytes(std::filesystem::file_size(path, ec)); + info->set_modified_unix_ms(modifiedUnixMs(path)); + if (const auto* descriptor = descriptorForPath(relative_path)) { + info->set_message_type(descriptor->full_name()); + } +} + +bool writeAll(int fd, const std::string& content) +{ + std::size_t offset = 0; + while (offset < content.size()) { + const ssize_t written = ::write(fd, content.data() + offset, content.size() - offset); + if (written < 0) { + if (errno == EINTR) continue; + return false; + } + offset += static_cast(written); + } + return true; +} + +bool atomicWriteFile( + const std::filesystem::path& target, + const std::string& content, + std::string* error) +{ + static std::atomic sequence{0}; + struct stat existing {}; + if (::stat(target.c_str(), &existing) != 0) { + if (error) *error = "Cannot inspect existing config: " + std::string(std::strerror(errno)); + return false; + } + + const auto temporary = target.string() + ".tmp." + std::to_string(::getpid()) + "." + + std::to_string(++sequence); + const int fd = ::open( + temporary.c_str(), O_WRONLY | O_CREAT | O_EXCL | O_CLOEXEC | O_NOFOLLOW, + existing.st_mode & 0777); + if (fd < 0) { + if (error) *error = "Cannot create temporary config: " + std::string(std::strerror(errno)); + return false; + } + + bool ok = writeAll(fd, content); + if (ok) ok = ::fsync(fd) == 0; + const int close_result = ::close(fd); + ok = ok && close_result == 0; + if (ok) ok = ::rename(temporary.c_str(), target.c_str()) == 0; + if (!ok) { + const int saved_errno = errno; + ::unlink(temporary.c_str()); + if (error) *error = "Cannot save config atomically: " + + std::string(std::strerror(saved_errno)); + return false; + } + + const int directory_fd = ::open(target.parent_path().c_str(), O_RDONLY | O_DIRECTORY | O_CLOEXEC); + if (directory_fd >= 0) { + ::fsync(directory_fd); + ::close(directory_fd); + } + return true; +} + +} // namespace + +ConfigFileService::ConfigFileService(std::filesystem::path config_root) +{ + const std::string requested_root = config_root.string(); + std::error_code ec; + config_root_ = std::filesystem::canonical(std::move(config_root), ec); + if (ec || !std::filesystem::is_directory(config_root_)) { + throw std::runtime_error("Invalid config root: " + requested_root); + } +} + +bool ConfigFileService::resolveExistingConfigPath( + const std::string& relative_path, + std::filesystem::path* resolved, + std::string* normalized_relative, + std::string* error) const +{ + if (relative_path.empty()) { + if (error) *error = "Config relative path is required"; + return false; + } + const std::filesystem::path input(relative_path); + if (input.is_absolute() || input.has_root_path()) { + if (error) *error = "Absolute config paths are not allowed"; + return false; + } + const auto normalized = input.lexically_normal(); + for (const auto& component : normalized) { + if (component == "..") { + if (error) *error = "Config path cannot leave the config root"; + return false; + } + } + if (!hasConfigExtension(normalized)) { + if (error) *error = "Only .pb.txt files are allowed"; + return false; + } + + std::error_code ec; + const auto unresolved = config_root_ / normalized; + const auto unresolved_status = std::filesystem::symlink_status(unresolved, ec); + if (ec || std::filesystem::is_symlink(unresolved_status)) { + if (error) *error = "Config file symlinks are not allowed"; + return false; + } + const auto candidate = std::filesystem::canonical(unresolved, ec); + if (ec || !std::filesystem::is_regular_file(candidate) || + !std::filesystem::is_regular_file(unresolved_status)) { + if (error) *error = "Config file does not exist or is not a regular file"; + return false; + } + const auto relative = std::filesystem::relative(candidate, config_root_, ec); + if (ec || relative.empty()) { + if (error) *error = "Config path cannot be resolved inside the config root"; + return false; + } + for (const auto& component : relative) { + if (component == "..") { + if (error) *error = "Config path resolves outside the config root"; + return false; + } + } + *resolved = candidate; + *normalized_relative = relative.generic_string(); + return true; +} + +grpc::Status ConfigFileService::ListConfigFiles( + grpc::ServerContext*, const api::ListConfigFilesCommand_Request*, + api::ListConfigFilesCommand_Feedback* response) +{ + std::lock_guard lock(mutex_); + std::vector paths; + std::error_code ec; + for (std::filesystem::recursive_directory_iterator iterator( + config_root_, std::filesystem::directory_options::skip_permission_denied, ec), end; + iterator != end; iterator.increment(ec)) { + if (ec) { + ec.clear(); + continue; + } + const auto status = iterator->symlink_status(ec); + if (ec || std::filesystem::is_symlink(status)) { + if (iterator->is_directory(ec)) iterator.disable_recursion_pending(); + ec.clear(); + continue; + } + if (std::filesystem::is_regular_file(status) && hasConfigExtension(iterator->path())) { + paths.push_back(iterator->path()); + } + } + std::sort(paths.begin(), paths.end()); + response->set_config_root(config_root_.string()); + for (const auto& path : paths) { + const auto relative = std::filesystem::relative(path, config_root_).generic_string(); + fillFileInfo(path, relative, response->add_files()); + } + succeed(response); + return grpc::Status::OK; +} + +grpc::Status ConfigFileService::ReadConfigFile( + grpc::ServerContext*, const api::ReadConfigFileCommand_Request* request, + api::ReadConfigFileCommand_Feedback* response) +{ + std::lock_guard lock(mutex_); + std::filesystem::path path; + std::string relative; + std::string error; + if (!resolveExistingConfigPath(request->relative_path(), &path, &relative, &error)) { + return fail(response, error); + } + std::string content; + if (!readFile(path, &content, &error)) return fail(response, error); + fillFileInfo(path, relative, response->mutable_file()); + response->set_content(content); + response->set_revision(revisionFor(content)); + succeed(response); + return grpc::Status::OK; +} + +grpc::Status ConfigFileService::ValidateConfigFile( + grpc::ServerContext*, const api::ValidateConfigFileCommand_Request* request, + api::ValidateConfigFileCommand_Feedback* response) +{ + std::lock_guard lock(mutex_); + std::filesystem::path path; + std::string relative; + std::string error; + if (!resolveExistingConfigPath(request->relative_path(), &path, &relative, &error)) { + return fail(response, error); + } + if (request->content().size() > kMaximumConfigBytes) { + return fail(response, "Config content exceeds 4 MiB limit"); + } + std::string message_type; + if (!validateTextProto(relative, request->content(), &message_type, &error)) { + return fail(response, error); + } + response->set_message_type(message_type); + succeed(response); + return grpc::Status::OK; +} + +grpc::Status ConfigFileService::SaveConfigFile( + grpc::ServerContext*, const api::SaveConfigFileCommand_Request* request, + api::SaveConfigFileCommand_Feedback* response) +{ + std::lock_guard lock(mutex_); + std::filesystem::path path; + std::string relative; + std::string error; + if (!resolveExistingConfigPath(request->relative_path(), &path, &relative, &error)) { + return fail(response, error); + } + if (request->content().size() > kMaximumConfigBytes) { + return fail(response, "Config content exceeds 4 MiB limit"); + } + + std::string current_content; + if (!readFile(path, ¤t_content, &error)) return fail(response, error); + if (request->expected_revision().empty()) { + return fail(response, "Expected revision is required"); + } + if (request->expected_revision() != revisionFor(current_content)) { + return fail(response, "Config file changed on the server; reload before saving"); + } + if (!validateTextProto(relative, request->content(), nullptr, &error)) { + return fail(response, error); + } + if (!atomicWriteFile(path, request->content(), &error)) return fail(response, error); + + response->set_revision(revisionFor(request->content())); + response->set_modified_unix_ms(modifiedUnixMs(path)); + succeed(response); + return grpc::Status::OK; +} + +} // namespace cmvr::config_server diff --git a/cmvr-es/config_server/config_file_service.h b/cmvr-es/config_server/config_file_service.h new file mode 100644 index 00000000..1ffa385e --- /dev/null +++ b/cmvr-es/config_server/config_file_service.h @@ -0,0 +1,45 @@ +#pragma once + +#include +#include +#include + +#include "cmvr/api/configuration_service.grpc.pb.h" + +namespace cmvr::config_server { + +class ConfigFileService final : public api::ConfigurationService::Service { +public: + explicit ConfigFileService(std::filesystem::path config_root); + + grpc::Status ListConfigFiles( + grpc::ServerContext* context, + const api::ListConfigFilesCommand_Request* request, + api::ListConfigFilesCommand_Feedback* response) override; + grpc::Status ReadConfigFile( + grpc::ServerContext* context, + const api::ReadConfigFileCommand_Request* request, + api::ReadConfigFileCommand_Feedback* response) override; + grpc::Status ValidateConfigFile( + grpc::ServerContext* context, + const api::ValidateConfigFileCommand_Request* request, + api::ValidateConfigFileCommand_Feedback* response) override; + grpc::Status SaveConfigFile( + grpc::ServerContext* context, + const api::SaveConfigFileCommand_Request* request, + api::SaveConfigFileCommand_Feedback* response) override; + + const std::filesystem::path& configRoot() const noexcept { return config_root_; } + +private: + bool resolveExistingConfigPath( + const std::string& relative_path, + std::filesystem::path* resolved, + std::string* normalized_relative, + std::string* error) const; + + std::filesystem::path config_root_; + mutable std::mutex mutex_; +}; + +} // namespace cmvr::config_server diff --git a/cmvr-es/config_server/config_server_main.cpp b/cmvr-es/config_server/config_server_main.cpp new file mode 100644 index 00000000..0027efad --- /dev/null +++ b/cmvr-es/config_server/config_server_main.cpp @@ -0,0 +1,82 @@ +#include +#include +#include +#include +#include +#include + +#include + +#include "config_server/config_file_service.h" + +namespace { + +std::filesystem::path executableDirectory() +{ + std::error_code ec; + const auto executable = std::filesystem::read_symlink("/proc/self/exe", ec); + return ec ? std::filesystem::current_path() : executable.parent_path(); +} + +void printUsage(const char* program) +{ + std::cerr << "Usage: " << program + << " [--listen ADDRESS] [--config-root DIRECTORY]\n"; +} + +} // namespace + +int main(int argc, char** argv) +{ + std::string listen_address = "0.0.0.0:50052"; + std::filesystem::path config_root = executableDirectory() / "config"; + for (int index = 1; index < argc; ++index) { + const std::string argument = argv[index]; + if (argument == "--listen" && index + 1 < argc) { + listen_address = argv[++index]; + } else if (argument == "--config-root" && index + 1 < argc) { + config_root = argv[++index]; + } else if (argument == "--help" || argument == "-h") { + printUsage(argv[0]); + return 0; + } else { + printUsage(argv[0]); + return 2; + } + } + + sigset_t shutdown_signals; + sigemptyset(&shutdown_signals); + sigaddset(&shutdown_signals, SIGINT); + sigaddset(&shutdown_signals, SIGTERM); + if (pthread_sigmask(SIG_BLOCK, &shutdown_signals, nullptr) != 0) { + std::cerr << "Failed to block shutdown signals\n"; + return 1; + } + + try { + cmvr::config_server::ConfigFileService service(config_root); + grpc::ServerBuilder builder; + builder.AddListeningPort(listen_address, grpc::InsecureServerCredentials()); + builder.RegisterService(&service); + std::unique_ptr server = builder.BuildAndStart(); + if (!server) { + std::cerr << "Failed to start configuration server at " << listen_address << '\n'; + return 1; + } + std::cout << "CMVR configuration server listening at " << listen_address + << ", root=" << service.configRoot() << std::endl; + int received_signal = 0; + if (sigwait(&shutdown_signals, &received_signal) != 0) { + std::cerr << "Failed while waiting for shutdown signal\n"; + server->Shutdown(); + return 1; + } + server->Shutdown(); + server->Wait(); + return 0; + } catch (const std::exception& error) { + std::cerr << "Configuration server error: " << error.what() << '\n'; + return 1; + } +} diff --git a/protos/cmvr/api/configuration_command.proto b/protos/cmvr/api/configuration_command.proto new file mode 100644 index 00000000..e696e829 --- /dev/null +++ b/protos/cmvr/api/configuration_command.proto @@ -0,0 +1,57 @@ +syntax = "proto3"; + +import "cmvr/api/common.proto"; + +package cmvr.api; + +message ConfigFileInfo { + string relative_path = 1; + uint64 size_bytes = 2; + int64 modified_unix_ms = 3; + string message_type = 4; +} + +message ListConfigFilesCommand { + message Request {} + message Feedback { + CommandHeader.Feedback header = 1; + string config_root = 2; + repeated ConfigFileInfo files = 3; + } +} + +message ReadConfigFileCommand { + message Request { + string relative_path = 1; + } + message Feedback { + CommandHeader.Feedback header = 1; + ConfigFileInfo file = 2; + string content = 3; + string revision = 4; + } +} + +message ValidateConfigFileCommand { + message Request { + string relative_path = 1; + string content = 2; + } + message Feedback { + CommandHeader.Feedback header = 1; + string message_type = 2; + } +} + +message SaveConfigFileCommand { + message Request { + string relative_path = 1; + string content = 2; + string expected_revision = 3; + } + message Feedback { + CommandHeader.Feedback header = 1; + string revision = 2; + int64 modified_unix_ms = 3; + } +} diff --git a/protos/cmvr/api/configuration_service.proto b/protos/cmvr/api/configuration_service.proto new file mode 100644 index 00000000..a45ed4f2 --- /dev/null +++ b/protos/cmvr/api/configuration_service.proto @@ -0,0 +1,12 @@ +syntax = "proto3"; + +import "cmvr/api/configuration_command.proto"; + +package cmvr.api; + +service ConfigurationService { + rpc ListConfigFiles(ListConfigFilesCommand.Request) returns (ListConfigFilesCommand.Feedback); + rpc ReadConfigFile(ReadConfigFileCommand.Request) returns (ReadConfigFileCommand.Feedback); + rpc ValidateConfigFile(ValidateConfigFileCommand.Request) returns (ValidateConfigFileCommand.Feedback); + rpc SaveConfigFile(SaveConfigFileCommand.Request) returns (SaveConfigFileCommand.Feedback); +}