feat(quic): support multi-platform heartbeats
This commit is contained in:
parent
55912abad0
commit
4f36cf4957
@ -73,9 +73,45 @@ cmvr_es.pb.txt
|
||||
- 新增 loader 对不认识的 enum 和未设置的 oneof 必须明确失败;当前个别历史路径仍有退化默认行为,不应复制;
|
||||
- 设备端口、坐标系、速度和单位写入注释;
|
||||
- `enable` 应由 manager 层控制,后端内部的 enable 字段不能替代 manager 开关;
|
||||
- QUIC 需要 TaskManager 与 `QuicEdgeConfig.enable` 同时开启;
|
||||
- QUIC 任务需要在 TaskManager 中显式开启;
|
||||
- QUIC 零媒体轨道是合法配置。
|
||||
|
||||
### QUIC 多平台
|
||||
|
||||
`QuicEdgeTask` 可以同时连接多个平台。原有顶层
|
||||
`server_host`、`server_port`、`tls` 继续表示主平台;每个 `platforms` 条目会与主平台
|
||||
并行运行。也可以不配置顶层目标,只使用一个或多个 `platforms` 条目:
|
||||
|
||||
```protobuf
|
||||
platforms {
|
||||
id: "operations"
|
||||
server_host: "192.168.0.222"
|
||||
server_port: 4433
|
||||
enable_media: false
|
||||
tls {
|
||||
ca_file: "certs/cmvr-quic-ca.crt"
|
||||
server_name: "192.168.0.222"
|
||||
}
|
||||
}
|
||||
platforms {
|
||||
id: "analytics"
|
||||
server_host: "192.168.0.223"
|
||||
server_port: 4433
|
||||
enable_media: true
|
||||
tls {
|
||||
ca_file: "certs/cmvr-quic-ca.crt"
|
||||
server_name: "192.168.0.223"
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
平台 ID 和 `host:port` 必须分别唯一;配置兼容主平台时,其平台 ID 使用任务
|
||||
`QuicEdgeConfig.id`,新增条目也不能与它重名。每个平台拥有独立的 QUIC 连接、
|
||||
注册会话、心跳序号、ACK 超时和重连退避,一个平台断线不会阻塞其他平台。
|
||||
`enable_media` 默认为 `false`,此时仍发送注册、心跳、网络接口和设备状态,但不会
|
||||
复制音视频;设为 `true` 才会把全局 `tracks` 转发到该平台。兼容的顶层主平台保持
|
||||
原有媒体行为。
|
||||
|
||||
### gRPC 相机实时流
|
||||
|
||||
[`tasks/grpc_server_task/grpc_server_task.pb.txt`](tasks/grpc_server_task/grpc_server_task.pb.txt)
|
||||
|
||||
@ -29,6 +29,21 @@ quic_edge {
|
||||
allow_insecure: false
|
||||
}
|
||||
|
||||
# Additional platforms run concurrently with the primary endpoint above.
|
||||
# Each one has an independent connection, registration, heartbeat ACK state
|
||||
# and reconnect loop. Media forwarding is opt-in for additional platforms.
|
||||
# platforms {
|
||||
# id: "operations_backup"
|
||||
# server_host: "192.168.0.223"
|
||||
# server_port: 4433
|
||||
# enable_media: false
|
||||
# tls {
|
||||
# ca_file: "certs/cmvr-quic-ca.crt"
|
||||
# server_name: "192.168.0.223"
|
||||
# allow_insecure: false
|
||||
# }
|
||||
# }
|
||||
|
||||
reconnect {
|
||||
initial_delay_ms: 500
|
||||
maximum_delay_ms: 30000
|
||||
|
||||
@ -1,9 +1,11 @@
|
||||
#ifndef CMVR_ES_QUIC_EDGE_TASK_H
|
||||
#define CMVR_ES_QUIC_EDGE_TASK_H
|
||||
|
||||
#include <functional>
|
||||
#include <memory>
|
||||
#include <mutex>
|
||||
#include <string>
|
||||
#include <vector>
|
||||
|
||||
#include "cmvr/config/quic_edge_config/quic_edge_config.pb.h"
|
||||
#include "service/quic_edge/include/quic_edge_service.h"
|
||||
@ -13,7 +15,12 @@ namespace cmvr::task {
|
||||
|
||||
class QuicEdgeTask final : public Task {
|
||||
public:
|
||||
using ServiceFactory = std::function<std::unique_ptr<quic_edge::QuicEdgeService>(
|
||||
config::QuicEdgeConfig)>;
|
||||
|
||||
explicit QuicEdgeTask(const config::QuicEdgeConfig& config);
|
||||
QuicEdgeTask(const config::QuicEdgeConfig& config,
|
||||
ServiceFactory service_factory);
|
||||
~QuicEdgeTask() override;
|
||||
|
||||
const std::string& id() const override { return id_; }
|
||||
@ -32,11 +39,19 @@ public:
|
||||
std::string detailStatusString() const override;
|
||||
|
||||
private:
|
||||
struct PlatformService {
|
||||
std::string id;
|
||||
bool media_enabled{false};
|
||||
config::QuicEdgeConfig config;
|
||||
std::unique_ptr<quic_edge::QuicEdgeService> service;
|
||||
};
|
||||
|
||||
TaskState mappedState() const;
|
||||
|
||||
config::QuicEdgeConfig config_;
|
||||
std::string id_;
|
||||
std::unique_ptr<quic_edge::QuicEdgeService> service_;
|
||||
ServiceFactory service_factory_;
|
||||
std::vector<PlatformService> services_;
|
||||
|
||||
mutable std::mutex mutex_;
|
||||
TaskState state_{TaskState::UNINITIALIZED};
|
||||
|
||||
@ -1,7 +1,10 @@
|
||||
#include "task/quic_edge_task/include/quic_edge_task.h"
|
||||
|
||||
#include <set>
|
||||
#include <sstream>
|
||||
#include <unordered_set>
|
||||
#include <utility>
|
||||
#include <vector>
|
||||
|
||||
#include "cmvr/config/task_manager_config/task_manager_config.pb.h"
|
||||
#include "common/base/logging/logger.h"
|
||||
@ -13,6 +16,107 @@
|
||||
namespace cmvr::task {
|
||||
namespace {
|
||||
|
||||
struct PlatformConfig {
|
||||
std::string id;
|
||||
bool media_enabled{false};
|
||||
config::QuicEdgeConfig config;
|
||||
};
|
||||
|
||||
void resolveTlsFiles(config::QuicEdgeTlsConfig* tls)
|
||||
{
|
||||
if (!tls) return;
|
||||
if (!tls->ca_file().empty()) {
|
||||
tls->set_ca_file(ConfigHelper::resolveConfigFile(tls->ca_file()));
|
||||
}
|
||||
if (!tls->certificate_file().empty()) {
|
||||
tls->set_certificate_file(
|
||||
ConfigHelper::resolveConfigFile(tls->certificate_file()));
|
||||
}
|
||||
if (!tls->private_key_file().empty()) {
|
||||
tls->set_private_key_file(
|
||||
ConfigHelper::resolveConfigFile(tls->private_key_file()));
|
||||
}
|
||||
}
|
||||
|
||||
bool appendPlatformConfig(
|
||||
const config::QuicEdgeConfig& root_config,
|
||||
const std::string& platform_id,
|
||||
const std::string& server_host,
|
||||
const std::uint32_t server_port,
|
||||
const config::QuicEdgeTlsConfig& tls,
|
||||
const bool media_enabled,
|
||||
std::unordered_set<std::string>* platform_ids,
|
||||
std::set<std::pair<std::string, std::uint32_t>>* endpoints,
|
||||
std::vector<PlatformConfig>* platforms,
|
||||
std::string* error)
|
||||
{
|
||||
if (platform_id.empty()) {
|
||||
if (error) *error = "QUIC edge platform id is empty";
|
||||
return false;
|
||||
}
|
||||
if (!platform_ids->insert(platform_id).second) {
|
||||
if (error) *error = "duplicate QUIC edge platform id: " + platform_id;
|
||||
return false;
|
||||
}
|
||||
if (!endpoints->emplace(server_host, server_port).second) {
|
||||
if (error) {
|
||||
*error = "duplicate QUIC edge platform endpoint: " + server_host +
|
||||
':' + std::to_string(server_port);
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
PlatformConfig platform;
|
||||
platform.id = platform_id;
|
||||
platform.media_enabled = media_enabled;
|
||||
platform.config = root_config;
|
||||
platform.config.set_server_host(server_host);
|
||||
platform.config.set_server_port(server_port);
|
||||
*platform.config.mutable_tls() = tls;
|
||||
platform.config.clear_platforms();
|
||||
if (!media_enabled) platform.config.clear_tracks();
|
||||
platforms->push_back(std::move(platform));
|
||||
return true;
|
||||
}
|
||||
|
||||
bool expandPlatformConfigs(const config::QuicEdgeConfig& config,
|
||||
std::vector<PlatformConfig>* platforms,
|
||||
std::string* error)
|
||||
{
|
||||
if (!platforms) {
|
||||
if (error) *error = "QUIC edge platform output is null";
|
||||
return false;
|
||||
}
|
||||
platforms->clear();
|
||||
std::unordered_set<std::string> platform_ids;
|
||||
std::set<std::pair<std::string, std::uint32_t>> endpoints;
|
||||
|
||||
const bool has_primary_endpoint = !config.server_host().empty() ||
|
||||
config.server_port() != 0U || config.has_tls();
|
||||
if (has_primary_endpoint &&
|
||||
!appendPlatformConfig(
|
||||
config, config.id(), config.server_host(), config.server_port(),
|
||||
config.tls(), true, &platform_ids, &endpoints, platforms, error)) {
|
||||
return false;
|
||||
}
|
||||
|
||||
for (const auto& configured_platform : config.platforms()) {
|
||||
if (!appendPlatformConfig(
|
||||
config, configured_platform.id(),
|
||||
configured_platform.server_host(),
|
||||
configured_platform.server_port(), configured_platform.tls(),
|
||||
configured_platform.enable_media(), &platform_ids, &endpoints,
|
||||
platforms, error)) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if (platforms->empty()) {
|
||||
if (error) *error = "QUIC edge has no platform endpoint";
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
std::shared_ptr<Task> createQuicEdgeTask(const config::TaskConfigEntry& entry)
|
||||
{
|
||||
if (entry.id().empty() || entry.config_file().empty()) {
|
||||
@ -31,17 +135,9 @@ std::shared_ptr<Task> createQuicEdgeTask(const config::TaskConfigEntry& entry)
|
||||
<< ", config=" << config.id();
|
||||
return nullptr;
|
||||
}
|
||||
if (!config.tls().ca_file().empty()) {
|
||||
config.mutable_tls()->set_ca_file(
|
||||
ConfigHelper::resolveConfigFile(config.tls().ca_file()));
|
||||
}
|
||||
if (!config.tls().certificate_file().empty()) {
|
||||
config.mutable_tls()->set_certificate_file(
|
||||
ConfigHelper::resolveConfigFile(config.tls().certificate_file()));
|
||||
}
|
||||
if (!config.tls().private_key_file().empty()) {
|
||||
config.mutable_tls()->set_private_key_file(
|
||||
ConfigHelper::resolveConfigFile(config.tls().private_key_file()));
|
||||
if (config.has_tls()) resolveTlsFiles(config.mutable_tls());
|
||||
for (auto& platform : *config.mutable_platforms()) {
|
||||
if (platform.has_tls()) resolveTlsFiles(platform.mutable_tls());
|
||||
}
|
||||
if (config.software_version().empty()) {
|
||||
config.set_software_version(device::DeviceManager::getInstance().version());
|
||||
@ -52,7 +148,19 @@ std::shared_ptr<Task> createQuicEdgeTask(const config::TaskConfigEntry& entry)
|
||||
} // namespace
|
||||
|
||||
QuicEdgeTask::QuicEdgeTask(const config::QuicEdgeConfig& config)
|
||||
: config_(config), id_(config.id())
|
||||
: QuicEdgeTask(
|
||||
config,
|
||||
[](config::QuicEdgeConfig platform_config) {
|
||||
return std::make_unique<quic_edge::QuicEdgeService>(
|
||||
std::move(platform_config));
|
||||
})
|
||||
{
|
||||
}
|
||||
|
||||
QuicEdgeTask::QuicEdgeTask(const config::QuicEdgeConfig& config,
|
||||
ServiceFactory service_factory)
|
||||
: config_(config), id_(config.id()),
|
||||
service_factory_(std::move(service_factory))
|
||||
{
|
||||
}
|
||||
|
||||
@ -69,13 +177,40 @@ bool QuicEdgeTask::init()
|
||||
state_ = TaskState::FAILED;
|
||||
return false;
|
||||
}
|
||||
service_ = std::make_unique<quic_edge::QuicEdgeService>(config_);
|
||||
|
||||
std::vector<PlatformConfig> platform_configs;
|
||||
std::string error;
|
||||
if (!service_->initialize(&error)) {
|
||||
if (!service_factory_) {
|
||||
last_error_ = "QUIC edge service factory is empty";
|
||||
state_ = TaskState::FAILED;
|
||||
return false;
|
||||
}
|
||||
if (!expandPlatformConfigs(config_, &platform_configs, &error)) {
|
||||
last_error_ = std::move(error);
|
||||
state_ = TaskState::FAILED;
|
||||
return false;
|
||||
}
|
||||
|
||||
std::vector<PlatformService> initialized_services;
|
||||
initialized_services.reserve(platform_configs.size());
|
||||
for (auto& platform_config : platform_configs) {
|
||||
auto service = service_factory_(platform_config.config);
|
||||
if (!service) {
|
||||
last_error_ = "failed to create QUIC edge service for platform " +
|
||||
platform_config.id;
|
||||
state_ = TaskState::FAILED;
|
||||
return false;
|
||||
}
|
||||
if (!service->initialize(&error)) {
|
||||
last_error_ = "platform " + platform_config.id + ": " + error;
|
||||
state_ = TaskState::FAILED;
|
||||
return false;
|
||||
}
|
||||
initialized_services.push_back(PlatformService{
|
||||
std::move(platform_config.id), platform_config.media_enabled,
|
||||
std::move(platform_config.config), std::move(service)});
|
||||
}
|
||||
services_ = std::move(initialized_services);
|
||||
last_error_.clear();
|
||||
state_ = TaskState::IDLE;
|
||||
return true;
|
||||
@ -101,20 +236,30 @@ bool QuicEdgeTask::start()
|
||||
admission.generation() == admission_generation;
|
||||
}
|
||||
if ((state_ != TaskState::IDLE && state_ != TaskState::STOPPED) ||
|
||||
!service_) {
|
||||
services_.empty()) {
|
||||
last_error_ = "QUIC edge task is not initialized";
|
||||
state_ = TaskState::FAILED;
|
||||
return false;
|
||||
}
|
||||
std::string error;
|
||||
if (!service_->start(&error)) {
|
||||
last_error_ = std::move(error);
|
||||
state_ = TaskState::FAILED;
|
||||
return false;
|
||||
|
||||
std::size_t started_services = 0U;
|
||||
for (auto& platform : services_) {
|
||||
std::string error;
|
||||
if (!platform.service->start(&error)) {
|
||||
for (std::size_t index = 0U;
|
||||
index < started_services; ++index) {
|
||||
services_[index].service->stop();
|
||||
}
|
||||
last_error_ = "platform " + platform.id + ": " + error;
|
||||
state_ = TaskState::FAILED;
|
||||
return false;
|
||||
}
|
||||
++started_services;
|
||||
}
|
||||
last_error_.clear();
|
||||
state_ = TaskState::RUNNING;
|
||||
CMVR_LOG(INFO) << "[QuicEdgeTask] Started, id=" << id_;
|
||||
CMVR_LOG(INFO) << "[QuicEdgeTask] Started, id=" << id_
|
||||
<< ", platforms=" << services_.size();
|
||||
}
|
||||
|
||||
bool admission_current = false;
|
||||
@ -142,26 +287,37 @@ void QuicEdgeTask::stop()
|
||||
std::lock_guard lock(mutex_);
|
||||
// stop() is the task lifecycle terminator used by TaskManager shutdown.
|
||||
// Unlike StopAll's stopActivity(), it intentionally tears down transport.
|
||||
if (service_) service_->stop();
|
||||
for (auto& platform : services_) platform.service->stop();
|
||||
if (state_ != TaskState::FAILED) state_ = TaskState::STOPPED;
|
||||
}
|
||||
|
||||
bool QuicEdgeTask::stopActivity()
|
||||
{
|
||||
std::lock_guard lock(mutex_);
|
||||
if (!service_ || state_ == TaskState::FAILED) {
|
||||
if (services_.empty() || state_ == TaskState::FAILED) {
|
||||
return false;
|
||||
}
|
||||
// System StopAll must preserve the QUIC presence channel. Only old media
|
||||
// subscriptions are fenced; registration and heartbeat stay online.
|
||||
return service_->interruptMediaActivities();
|
||||
bool all_stopped = true;
|
||||
for (auto& platform : services_) {
|
||||
if (!platform.service->interruptMediaActivities()) {
|
||||
all_stopped = false;
|
||||
}
|
||||
}
|
||||
return all_stopped;
|
||||
}
|
||||
|
||||
TaskState QuicEdgeTask::mappedState() const
|
||||
{
|
||||
if (!service_ || state_ != TaskState::RUNNING) return state_;
|
||||
return service_->state() == quic_edge::QuicEdgeServiceState::FAILED
|
||||
? TaskState::FAILED : state_;
|
||||
if (services_.empty() || state_ != TaskState::RUNNING) return state_;
|
||||
for (const auto& platform : services_) {
|
||||
if (platform.service->state() ==
|
||||
quic_edge::QuicEdgeServiceState::FAILED) {
|
||||
return TaskState::FAILED;
|
||||
}
|
||||
}
|
||||
return state_;
|
||||
}
|
||||
|
||||
TaskState QuicEdgeTask::state() const
|
||||
@ -195,31 +351,38 @@ std::string QuicEdgeTask::detailStatusString() const
|
||||
std::lock_guard lock(mutex_);
|
||||
std::ostringstream output;
|
||||
output << taskStateToString(mappedState());
|
||||
if (service_) {
|
||||
const auto stats = service_->stats();
|
||||
const auto service_status = service_->status();
|
||||
output << " service=" << quic_edge::toString(service_->state())
|
||||
<< " target=" << config_.server_host() << ':' << config_.server_port()
|
||||
<< " node_id=" << service_status.node_id
|
||||
<< " registered=" << service_status.registered
|
||||
<< " connections=" << stats.successful_connections
|
||||
<< " registrations=" << stats.registrations_accepted
|
||||
<< " heartbeat_sequence=" << service_status.heartbeat_sequence
|
||||
<< " heartbeat_acks=" << stats.heartbeats_acknowledged
|
||||
<< " media_tracks=" << service_status.active_media_tracks
|
||||
<< " frames=" << stats.frames_queued
|
||||
<< " datagrams=" << stats.datagrams_queued;
|
||||
if (!service_status.session_id.empty()) {
|
||||
output << " session=" << service_status.session_id;
|
||||
if (!services_.empty()) {
|
||||
output << " platforms=" << services_.size();
|
||||
for (const auto& platform : services_) {
|
||||
const auto stats = platform.service->stats();
|
||||
const auto service_status = platform.service->status();
|
||||
output << " platform[" << platform.id << "]={"
|
||||
<< "service=" << quic_edge::toString(platform.service->state())
|
||||
<< ",target=" << platform.config.server_host() << ':'
|
||||
<< platform.config.server_port()
|
||||
<< ",media_enabled=" << platform.media_enabled
|
||||
<< ",node_id=" << service_status.node_id
|
||||
<< ",registered=" << service_status.registered
|
||||
<< ",connections=" << stats.successful_connections
|
||||
<< ",registrations=" << stats.registrations_accepted
|
||||
<< ",heartbeat_sequence=" << service_status.heartbeat_sequence
|
||||
<< ",heartbeat_acks=" << stats.heartbeats_acknowledged
|
||||
<< ",media_tracks=" << service_status.active_media_tracks
|
||||
<< ",frames=" << stats.frames_queued
|
||||
<< ",datagrams=" << stats.datagrams_queued;
|
||||
if (!service_status.session_id.empty()) {
|
||||
output << ",session=" << service_status.session_id;
|
||||
}
|
||||
if (!service_status.observed_source_ip.empty()) {
|
||||
output << ",observed_ip=" << service_status.observed_source_ip;
|
||||
}
|
||||
if (!service_status.last_media_error.empty()) {
|
||||
output << ",media_error=" << service_status.last_media_error;
|
||||
}
|
||||
const std::string service_error = platform.service->lastError();
|
||||
if (!service_error.empty()) output << ",error=" << service_error;
|
||||
output << '}';
|
||||
}
|
||||
if (!service_status.observed_source_ip.empty()) {
|
||||
output << " observed_ip=" << service_status.observed_source_ip;
|
||||
}
|
||||
if (!service_status.last_media_error.empty()) {
|
||||
output << " media_error=" << service_status.last_media_error;
|
||||
}
|
||||
const std::string service_error = service_->lastError();
|
||||
if (!service_error.empty()) output << " error=" << service_error;
|
||||
} else if (!last_error_.empty()) {
|
||||
output << " error=" << last_error_;
|
||||
}
|
||||
|
||||
@ -1,10 +1,179 @@
|
||||
#include <atomic>
|
||||
#include <chrono>
|
||||
#include <condition_variable>
|
||||
#include <cstdint>
|
||||
#include <deque>
|
||||
#include <iostream>
|
||||
#include <memory>
|
||||
#include <mutex>
|
||||
#include <string>
|
||||
#include <thread>
|
||||
#include <utility>
|
||||
#include <vector>
|
||||
|
||||
#include "cmvr/quic_edge/v1/quic_edge.pb.h"
|
||||
#include "manager/media_source_hub/include/media_source_hub.h"
|
||||
#include "service/quic_edge/include/control_framing.h"
|
||||
#include "service/stop_all/include/stop_all_admission_gate.h"
|
||||
#include "task/quic_edge_task/include/quic_edge_task.h"
|
||||
|
||||
namespace {
|
||||
|
||||
using Clock = std::chrono::steady_clock;
|
||||
using QuicEdgeService = cmvr::quic_edge::QuicEdgeService;
|
||||
|
||||
class HeartbeatTransport final : public cmvr::quic_edge::QuicTransport {
|
||||
public:
|
||||
explicit HeartbeatTransport(std::string platform_id)
|
||||
: platform_id_(std::move(platform_id))
|
||||
{
|
||||
}
|
||||
|
||||
bool connect(const cmvr::config::QuicEdgeConfig&,
|
||||
std::chrono::milliseconds,
|
||||
std::string*) override
|
||||
{
|
||||
connected_.store(true);
|
||||
return true;
|
||||
}
|
||||
|
||||
void disconnect() override
|
||||
{
|
||||
connected_.store(false);
|
||||
condition_.notify_all();
|
||||
}
|
||||
|
||||
bool isConnected() const override { return connected_.load(); }
|
||||
std::size_t maximumDatagramBytes() const override { return 1200U; }
|
||||
std::size_t maximumDatagramBatchPackets() const override { return 32U; }
|
||||
|
||||
cmvr::quic_edge::TransportSendResult sendControl(
|
||||
std::vector<std::uint8_t> framed_message,
|
||||
std::string*) override
|
||||
{
|
||||
if (!connected_.load()) {
|
||||
return cmvr::quic_edge::TransportSendResult::DISCONNECTED;
|
||||
}
|
||||
|
||||
std::vector<std::vector<std::uint8_t>> frames;
|
||||
std::string error;
|
||||
cmvr::quic_edge::ControlFrameDecoder decoder(1024U * 1024U);
|
||||
if (!decoder.push(framed_message, &frames, &error) ||
|
||||
frames.size() != 1U) {
|
||||
return cmvr::quic_edge::TransportSendResult::ERROR;
|
||||
}
|
||||
cmvr::quic_edge::v1::EdgeControlEnvelope request;
|
||||
if (!request.ParseFromArray(
|
||||
frames.front().data(),
|
||||
static_cast<int>(frames.front().size()))) {
|
||||
return cmvr::quic_edge::TransportSendResult::ERROR;
|
||||
}
|
||||
|
||||
cmvr::quic_edge::v1::EdgeControlEnvelope response;
|
||||
response.set_protocol_version(cmvr::quic_edge::kProtocolVersion);
|
||||
{
|
||||
std::lock_guard lock(mutex_);
|
||||
response.set_message_sequence(server_message_sequence_++);
|
||||
if (request.has_node_register_request()) {
|
||||
++registrations_;
|
||||
auto* registration = response.mutable_node_register_response();
|
||||
registration->set_accepted(true);
|
||||
registration->set_session_id("session-" + platform_id_);
|
||||
registration->set_heartbeat_interval_ms(250U);
|
||||
} else if (request.has_node_heartbeat()) {
|
||||
const auto& heartbeat = request.node_heartbeat();
|
||||
++heartbeats_;
|
||||
last_heartbeat_sequence_.store(heartbeat.sequence());
|
||||
auto* ack = response.mutable_node_heartbeat_ack();
|
||||
ack->set_accepted(true);
|
||||
ack->set_session_id(heartbeat.session_id());
|
||||
ack->set_acknowledged_sequence(heartbeat.sequence());
|
||||
} else {
|
||||
return cmvr::quic_edge::TransportSendResult::ERROR;
|
||||
}
|
||||
if (!enqueueResponse(response)) {
|
||||
return cmvr::quic_edge::TransportSendResult::ERROR;
|
||||
}
|
||||
}
|
||||
condition_.notify_all();
|
||||
return cmvr::quic_edge::TransportSendResult::QUEUED;
|
||||
}
|
||||
|
||||
cmvr::quic_edge::TransportReceiveResult receiveControl(
|
||||
std::vector<std::uint8_t>* chunk,
|
||||
const std::chrono::milliseconds timeout,
|
||||
std::string*) override
|
||||
{
|
||||
if (!chunk) return cmvr::quic_edge::TransportReceiveResult::ERROR;
|
||||
std::unique_lock lock(mutex_);
|
||||
condition_.wait_for(lock, timeout, [this] {
|
||||
return !responses_.empty() || !connected_.load();
|
||||
});
|
||||
if (!responses_.empty()) {
|
||||
*chunk = std::move(responses_.front());
|
||||
responses_.pop_front();
|
||||
return cmvr::quic_edge::TransportReceiveResult::DATA;
|
||||
}
|
||||
return connected_.load()
|
||||
? cmvr::quic_edge::TransportReceiveResult::TIMEOUT
|
||||
: cmvr::quic_edge::TransportReceiveResult::DISCONNECTED;
|
||||
}
|
||||
|
||||
cmvr::quic_edge::TransportSendResult sendDatagramBatch(
|
||||
std::vector<cmvr::quic_edge::DatagramPacket>,
|
||||
std::string*) override
|
||||
{
|
||||
return connected_.load()
|
||||
? cmvr::quic_edge::TransportSendResult::QUEUED
|
||||
: cmvr::quic_edge::TransportSendResult::DISCONNECTED;
|
||||
}
|
||||
|
||||
std::uint64_t registrations() const { return registrations_.load(); }
|
||||
std::uint64_t heartbeats() const { return heartbeats_.load(); }
|
||||
std::uint64_t lastHeartbeatSequence() const
|
||||
{
|
||||
return last_heartbeat_sequence_.load();
|
||||
}
|
||||
|
||||
private:
|
||||
bool enqueueResponse(
|
||||
const cmvr::quic_edge::v1::EdgeControlEnvelope& response)
|
||||
{
|
||||
std::string serialized;
|
||||
if (!response.SerializeToString(&serialized)) return false;
|
||||
std::vector<std::uint8_t> framed;
|
||||
std::string error;
|
||||
if (!cmvr::quic_edge::ControlFrameEncoder::encode(
|
||||
reinterpret_cast<const std::uint8_t*>(serialized.data()),
|
||||
serialized.size(), 1024U * 1024U, &framed, &error)) {
|
||||
return false;
|
||||
}
|
||||
responses_.push_back(std::move(framed));
|
||||
return true;
|
||||
}
|
||||
|
||||
std::string platform_id_;
|
||||
std::atomic<bool> connected_{false};
|
||||
std::atomic<std::uint64_t> registrations_{0U};
|
||||
std::atomic<std::uint64_t> heartbeats_{0U};
|
||||
std::atomic<std::uint64_t> last_heartbeat_sequence_{0U};
|
||||
mutable std::mutex mutex_;
|
||||
std::condition_variable condition_;
|
||||
std::deque<std::vector<std::uint8_t>> responses_;
|
||||
std::uint64_t server_message_sequence_{0U};
|
||||
};
|
||||
|
||||
template <typename Predicate>
|
||||
bool waitUntil(const std::chrono::milliseconds timeout, Predicate predicate)
|
||||
{
|
||||
const auto deadline = Clock::now() + timeout;
|
||||
while (Clock::now() < deadline) {
|
||||
if (predicate()) return true;
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds(5));
|
||||
}
|
||||
return predicate();
|
||||
}
|
||||
|
||||
cmvr::config::QuicEdgeConfig validConfig()
|
||||
{
|
||||
cmvr::config::QuicEdgeConfig config;
|
||||
@ -32,6 +201,36 @@ cmvr::config::QuicEdgeConfig validConfig()
|
||||
return config;
|
||||
}
|
||||
|
||||
cmvr::config::QuicEdgeConfig multiPlatformConfig()
|
||||
{
|
||||
auto config = validConfig();
|
||||
config.set_id("quic-multi-platform-test");
|
||||
config.clear_server_host();
|
||||
config.clear_server_port();
|
||||
config.clear_tls();
|
||||
|
||||
auto* platform_a = config.add_platforms();
|
||||
platform_a->set_id("platform-a");
|
||||
platform_a->set_server_host("192.0.2.10");
|
||||
platform_a->set_server_port(4433U);
|
||||
platform_a->mutable_tls()->set_allow_insecure(true);
|
||||
|
||||
auto* platform_b = config.add_platforms();
|
||||
platform_b->set_id("platform-b");
|
||||
platform_b->set_server_host("198.51.100.20");
|
||||
platform_b->set_server_port(4434U);
|
||||
platform_b->mutable_tls()->set_allow_insecure(true);
|
||||
platform_b->set_enable_media(true);
|
||||
|
||||
auto* disabled_track = config.add_tracks();
|
||||
disabled_track->set_track_id(1U);
|
||||
disabled_track->set_source_kind(
|
||||
cmvr::config::QuicEdgeTrackConfig::SOURCE_KIND_CAMERA);
|
||||
disabled_track->set_device_id("disabled-test-camera");
|
||||
disabled_track->set_enable(false);
|
||||
return config;
|
||||
}
|
||||
|
||||
} // namespace
|
||||
|
||||
int main()
|
||||
@ -97,7 +296,93 @@ int main()
|
||||
std::cerr << "QUIC service did not remain available after StopAll\n";
|
||||
return 1;
|
||||
}
|
||||
admission_task.stop();
|
||||
admission.clearForTesting();
|
||||
|
||||
cmvr::media::MediaSourceHub media_hub;
|
||||
std::vector<HeartbeatTransport*> transports;
|
||||
std::vector<QuicEdgeService*> services;
|
||||
std::vector<cmvr::config::QuicEdgeConfig> expanded_configs;
|
||||
cmvr::task::QuicEdgeTask multi_platform_task(
|
||||
multiPlatformConfig(),
|
||||
[&](cmvr::config::QuicEdgeConfig platform_config) {
|
||||
expanded_configs.push_back(platform_config);
|
||||
auto transport = std::make_unique<HeartbeatTransport>(
|
||||
platform_config.server_host());
|
||||
transports.push_back(transport.get());
|
||||
auto service = std::make_unique<QuicEdgeService>(
|
||||
std::move(platform_config), std::move(transport), media_hub,
|
||||
[] { return cmvr::device::DeviceManagerSnapshot{}; });
|
||||
services.push_back(service.get());
|
||||
return service;
|
||||
});
|
||||
if (!multi_platform_task.init() || expanded_configs.size() != 2U ||
|
||||
transports.size() != 2U || services.size() != 2U) {
|
||||
std::cerr << "multi-platform QUIC task did not create two services: "
|
||||
<< multi_platform_task.detailStatusString() << '\n';
|
||||
return 1;
|
||||
}
|
||||
if (expanded_configs[0].platforms_size() != 0 ||
|
||||
expanded_configs[1].platforms_size() != 0 ||
|
||||
expanded_configs[0].server_host() != "192.0.2.10" ||
|
||||
expanded_configs[1].server_host() != "198.51.100.20" ||
|
||||
expanded_configs[0].tracks_size() != 0 ||
|
||||
expanded_configs[1].tracks_size() != 1) {
|
||||
std::cerr << "multi-platform QUIC task expanded invalid service configs\n";
|
||||
return 1;
|
||||
}
|
||||
if (!multi_platform_task.start()) {
|
||||
std::cerr << "multi-platform QUIC task did not start: "
|
||||
<< multi_platform_task.detailStatusString() << '\n';
|
||||
return 1;
|
||||
}
|
||||
const bool both_platforms_online = waitUntil(
|
||||
std::chrono::seconds(2), [&] {
|
||||
return transports[0]->registrations() >= 1U &&
|
||||
transports[1]->registrations() >= 1U &&
|
||||
transports[0]->heartbeats() >= 1U &&
|
||||
transports[1]->heartbeats() >= 1U &&
|
||||
services[0]->stats().heartbeats_acknowledged >= 1U &&
|
||||
services[1]->stats().heartbeats_acknowledged >= 1U;
|
||||
});
|
||||
if (!both_platforms_online ||
|
||||
transports[0]->lastHeartbeatSequence() != 1U ||
|
||||
transports[1]->lastHeartbeatSequence() != 1U) {
|
||||
std::cerr << "both platforms did not independently ACK a heartbeat: "
|
||||
<< multi_platform_task.detailStatusString() << '\n';
|
||||
return 1;
|
||||
}
|
||||
const std::string multi_status = multi_platform_task.detailStatusString();
|
||||
if (multi_status.find("platforms=2") == std::string::npos ||
|
||||
multi_status.find("platform[platform-a]") == std::string::npos ||
|
||||
multi_status.find("platform[platform-b]") == std::string::npos) {
|
||||
std::cerr << "multi-platform QUIC status is incomplete: "
|
||||
<< multi_status << '\n';
|
||||
return 1;
|
||||
}
|
||||
multi_platform_task.stop();
|
||||
if (multi_platform_task.state() != cmvr::task::TaskState::STOPPED) {
|
||||
std::cerr << "multi-platform QUIC task did not stop cleanly\n";
|
||||
return 1;
|
||||
}
|
||||
|
||||
auto duplicate_endpoint_config = validConfig();
|
||||
auto* duplicate_platform = duplicate_endpoint_config.add_platforms();
|
||||
duplicate_platform->set_id("duplicate-primary");
|
||||
duplicate_platform->set_server_host(
|
||||
duplicate_endpoint_config.server_host());
|
||||
duplicate_platform->set_server_port(
|
||||
duplicate_endpoint_config.server_port());
|
||||
duplicate_platform->mutable_tls()->set_allow_insecure(true);
|
||||
cmvr::task::QuicEdgeTask duplicate_endpoint_task(
|
||||
duplicate_endpoint_config);
|
||||
if (duplicate_endpoint_task.init() ||
|
||||
duplicate_endpoint_task.detailStatusString().find(
|
||||
"duplicate QUIC edge platform endpoint") == std::string::npos) {
|
||||
std::cerr << "duplicate QUIC platform endpoint was not rejected\n";
|
||||
return 1;
|
||||
}
|
||||
|
||||
std::cout << "quic_edge_task_test: PASS\n";
|
||||
return 0;
|
||||
}
|
||||
|
||||
@ -43,9 +43,27 @@ message QuicEdgeTrackConfig {
|
||||
string source_track_id = 6;
|
||||
}
|
||||
|
||||
message QuicEdgePlatformConfig {
|
||||
// Stable operator-facing identifier used in logs and task status. It must be
|
||||
// unique within the task and differ from QuicEdgeConfig.id when the
|
||||
// backward-compatible primary endpoint is also configured.
|
||||
string id = 1;
|
||||
string server_host = 2;
|
||||
uint32 server_port = 3;
|
||||
QuicEdgeTlsConfig tls = 4;
|
||||
|
||||
// False keeps registration, heartbeat, IP and device-state reporting while
|
||||
// preventing accidental duplication of configured media tracks.
|
||||
bool enable_media = 5;
|
||||
}
|
||||
|
||||
message QuicEdgeConfig {
|
||||
string id = 1;
|
||||
reserved 2;
|
||||
|
||||
// Backward-compatible primary platform endpoint. When configured, it runs
|
||||
// alongside every endpoint in platforms. New configurations may leave these
|
||||
// fields empty and declare every destination in platforms instead.
|
||||
string server_host = 3;
|
||||
uint32 server_port = 4;
|
||||
string alpn = 5;
|
||||
@ -87,6 +105,10 @@ message QuicEdgeConfig {
|
||||
// Immutable robot product identity, for example:
|
||||
// CN-CMVR-MBLRV1-CHAGAN-20260731-001.
|
||||
string robot_id = 22;
|
||||
|
||||
// Every platform owns an independent QUIC connection, registration,
|
||||
// heartbeat/ACK sequence and reconnect loop.
|
||||
repeated QuicEdgePlatformConfig platforms = 23;
|
||||
}
|
||||
|
||||
message QuicEdgeRootConfig {
|
||||
|
||||
@ -85,8 +85,8 @@ their exact scope and manual gateway options.
|
||||
|
||||
To enable node presence without media:
|
||||
|
||||
1. Configure the gateway address, ALPN `cmvr-quic-edge/1`, TLS trust, node ID,
|
||||
advertised gRPC endpoint and heartbeat values in
|
||||
1. Configure one or more gateway addresses, ALPN `cmvr-quic-edge/1`, TLS
|
||||
trust, node ID, advertised gRPC endpoint and heartbeat values in
|
||||
`cmvr-es/config/tasks/quic_edge_task/quic_edge_task.pb.txt`.
|
||||
Certificate paths are relative to the cmvr-es configuration root, for
|
||||
example `certs/quic_gateway_ca.pem`.
|
||||
@ -99,6 +99,14 @@ media must not suppress node registration, IP reporting or heartbeat.
|
||||
`maximum_frame_bytes` must fit in one atomically admitted DATAGRAM batch; reduce
|
||||
it or increase `datagram_send_queue_depth` when changing the DATAGRAM size.
|
||||
|
||||
The backward-compatible top-level `server_host`, `server_port` and `tls`
|
||||
describe the primary platform. Every `platforms` entry runs concurrently with
|
||||
that primary endpoint, or the list may be used by itself. Platform IDs and
|
||||
`host:port` pairs must be unique. Each platform owns its own connection,
|
||||
registration session, heartbeat/ACK sequence and reconnect loop. Additional
|
||||
platforms default to presence-only operation; set their `enable_media=true` to
|
||||
forward the shared `tracks` as a separate media session.
|
||||
|
||||
## Reliable control stream
|
||||
|
||||
The edge opens one bidirectional reliable stream after the QUIC handshake. Both
|
||||
|
||||
Loading…
Reference in New Issue
Block a user