feat(config): add standalone remote configuration service

This commit is contained in:
linbo 2026-08-25 17:37:53 +08:00
parent 02cff4483f
commit ffef60e9b8
6 changed files with 668 additions and 0 deletions

View File

@ -157,6 +157,17 @@ target_link_libraries(cmvr_es PRIVATE
) )
install(TARGETS cmvr_es RUNTIME DESTINATION bin) 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 [[ install(CODE [[
file(REMOVE_RECURSE file(REMOVE_RECURSE
"${CMAKE_INSTALL_PREFIX}/bin/config" "${CMAKE_INSTALL_PREFIX}/bin/config"

View File

@ -0,0 +1,461 @@
#include "config_server/config_file_service.h"
#include <algorithm>
#include <atomic>
#include <cerrno>
#include <chrono>
#include <cstdint>
#include <cstring>
#include <fcntl.h>
#include <fstream>
#include <iomanip>
#include <memory>
#include <sstream>
#include <stdexcept>
#include <string_view>
#include <sys/stat.h>
#include <unistd.h>
#include <vector>
#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<std::chrono::seconds>(now);
const auto nanos = std::chrono::duration_cast<std::chrono::nanoseconds>(now - seconds);
timestamp->set_seconds(seconds.count());
timestamp->set_nanos(static_cast<int>(nanos.count()));
}
template <typename Response>
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 <typename Response>
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<google::protobuf::Message> 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<char>(input), std::istreambuf_iterator<char>());
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<int64_t>(info.st_mtim.tv_sec) * 1000LL +
static_cast<int64_t>(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<std::size_t>(written);
}
return true;
}
bool atomicWriteFile(
const std::filesystem::path& target,
const std::string& content,
std::string* error)
{
static std::atomic<std::uint64_t> 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<std::filesystem::path> 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, &current_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

View File

@ -0,0 +1,45 @@
#pragma once
#include <filesystem>
#include <mutex>
#include <string>
#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

View File

@ -0,0 +1,82 @@
#include <csignal>
#include <filesystem>
#include <iostream>
#include <memory>
#include <pthread.h>
#include <string>
#include <grpcpp/grpcpp.h>
#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<grpc::Server> 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;
}
}

View File

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

View File

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