cmvr-es/cmvr-es/service/quic_edge/src/quic_edge_service.cpp
2026-07-27 16:07:41 +08:00

1892 lines
72 KiB
C++

#include "service/quic_edge/include/quic_edge_service.h"
#include <algorithm>
#include <array>
#include <cerrno>
#include <chrono>
#include <cmath>
#include <cstring>
#include <exception>
#include <fstream>
#include <iomanip>
#include <limits>
#include <random>
#include <set>
#include <sstream>
#include <unordered_set>
#include <utility>
#include <arpa/inet.h>
#include <ifaddrs.h>
#include <net/if.h>
#include <sys/socket.h>
#include <unistd.h>
#include "cmvr/quic_edge/v1/quic_edge.pb.h"
#include "common/base/logging/logger.h"
namespace cmvr::quic_edge {
namespace {
constexpr std::uint32_t kMinimumDatagramBytes =
static_cast<std::uint32_t>(kDatagramHeaderBytes + 1U);
constexpr std::uint32_t kMaximumDatagramBytes = 65527U;
constexpr std::uint32_t kMaximumControlBytes = 16U * 1024U * 1024U;
constexpr std::uint32_t kMaximumConfiguredFrameBytes = 256U * 1024U * 1024U;
constexpr std::size_t kMaximumDrainPerPoll = 4096U;
constexpr std::size_t kMaximumDeviceErrorBytes = 512U;
constexpr std::uint32_t kMinimumHeartbeatIntervalMs = 250U;
constexpr std::uint32_t kMaximumHeartbeatIntervalMs = 60U * 60U * 1000U;
constexpr auto kControlPollInterval = std::chrono::milliseconds(50);
constexpr auto kMediaSourceRetryInterval = std::chrono::seconds(1);
constexpr std::size_t kMaximumReservedControlSends = 8U;
void setError(std::string* error, const std::string& message)
{
if (error) {
*error = message;
}
}
std::string boundedDeviceError(std::string message)
{
if (message.size() <= kMaximumDeviceErrorBytes) return message;
std::size_t end = kMaximumDeviceErrorBytes;
while (end > 0U &&
(static_cast<unsigned char>(message[end]) & 0xc0U) == 0x80U) {
--end;
}
message.resize(end);
return message;
}
std::string trim(std::string value)
{
const auto first = value.find_first_not_of(" \t\r\n");
if (first == std::string::npos) return {};
const auto last = value.find_last_not_of(" \t\r\n");
return value.substr(first, last - first + 1U);
}
std::string readTextFile(const std::string& path)
{
std::ifstream input(path, std::ios::in | std::ios::binary);
if (!input) return {};
std::ostringstream contents;
contents << input.rdbuf();
return input.good() || input.eof() ? contents.str() : std::string{};
}
std::string randomBootId()
{
std::random_device random_device;
std::mt19937_64 generator(random_device());
std::uniform_int_distribution<std::uint64_t> distribution;
const std::uint64_t high = distribution(generator);
const std::uint64_t low = distribution(generator);
std::ostringstream stream;
stream << std::hex << std::setfill('0')
<< std::setw(8) << static_cast<std::uint32_t>(high >> 32U) << '-'
<< std::setw(4) << static_cast<std::uint16_t>(high >> 16U) << '-'
<< std::setw(4) << static_cast<std::uint16_t>(high) << '-'
<< std::setw(4) << static_cast<std::uint16_t>(low >> 48U) << '-'
<< std::setw(12) << (low & 0x0000ffffffffffffULL);
return stream.str();
}
std::string resolveNodeId(const std::string& configured_node_id)
{
if (!configured_node_id.empty() && configured_node_id != "auto") {
return configured_node_id;
}
std::array<char, 256> hostname{};
if (::gethostname(hostname.data(), hostname.size() - 1U) == 0) {
hostname.back() = '\0';
const std::string resolved = trim(hostname.data());
if (!resolved.empty()) return resolved;
}
return "cmvr-es";
}
std::string resolveBootId()
{
const std::string kernel_boot_id =
trim(readTextFile("/proc/sys/kernel/random/boot_id"));
return kernel_boot_id.empty() ? randomBootId() : kernel_boot_id;
}
std::uint64_t unixTimeMs()
{
const auto now = std::chrono::system_clock::now().time_since_epoch();
return static_cast<std::uint64_t>(
std::chrono::duration_cast<std::chrono::milliseconds>(now).count());
}
std::vector<v1::NetworkInterfaceAddress> collectInterfaceAddresses(
const bool include_loopback)
{
std::vector<v1::NetworkInterfaceAddress> addresses;
ifaddrs* interface_list = nullptr;
if (::getifaddrs(&interface_list) != 0 || interface_list == nullptr) {
return addresses;
}
std::set<std::string> seen;
for (const ifaddrs* current = interface_list;
current != nullptr; current = current->ifa_next) {
if (!current->ifa_addr || !current->ifa_name ||
(current->ifa_flags & IFF_UP) == 0) {
continue;
}
const bool loopback = (current->ifa_flags & IFF_LOOPBACK) != 0;
if (loopback && !include_loopback) continue;
const int family = current->ifa_addr->sa_family;
std::array<char, INET6_ADDRSTRLEN> buffer{};
const void* source = nullptr;
v1::NetworkInterfaceAddress::AddressFamily proto_family =
v1::NetworkInterfaceAddress::ADDRESS_FAMILY_UNSPECIFIED;
if (family == AF_INET) {
source = &reinterpret_cast<const sockaddr_in*>(
current->ifa_addr)->sin_addr;
proto_family = v1::NetworkInterfaceAddress::ADDRESS_FAMILY_IPV4;
} else if (family == AF_INET6) {
source = &reinterpret_cast<const sockaddr_in6*>(
current->ifa_addr)->sin6_addr;
proto_family = v1::NetworkInterfaceAddress::ADDRESS_FAMILY_IPV6;
} else {
continue;
}
if (::inet_ntop(family, source, buffer.data(), buffer.size()) == nullptr) {
continue;
}
const std::string ip_address(buffer.data());
if (ip_address.empty() || ip_address == "0.0.0.0" || ip_address == "::") {
continue;
}
const std::string key = std::string(current->ifa_name) + '\n' + ip_address;
if (!seen.insert(key).second) continue;
v1::NetworkInterfaceAddress address;
address.set_interface_name(current->ifa_name);
address.set_ip_address(ip_address);
address.set_family(proto_family);
address.set_loopback(loopback);
addresses.push_back(std::move(address));
}
::freeifaddrs(interface_list);
std::sort(addresses.begin(), addresses.end(),
[](const v1::NetworkInterfaceAddress& lhs,
const v1::NetworkInterfaceAddress& rhs) {
if (lhs.loopback() != rhs.loopback()) return !lhs.loopback();
if (lhs.interface_name() != rhs.interface_name()) {
return lhs.interface_name() < rhs.interface_name();
}
if (lhs.family() != rhs.family()) return lhs.family() < rhs.family();
return lhs.ip_address() < rhs.ip_address();
});
return addresses;
}
std::string advertisedGrpcHost(
const config::QuicEdgeConfig& config,
const std::vector<v1::NetworkInterfaceAddress>& interfaces)
{
if (!config.grpc_endpoint_host().empty() &&
config.grpc_endpoint_host() != "auto") {
return config.grpc_endpoint_host();
}
for (const auto& address : interfaces) {
if (!address.loopback() &&
address.family() == v1::NetworkInterfaceAddress::ADDRESS_FAMILY_IPV4) {
return address.ip_address();
}
}
for (const auto& address : interfaces) {
if (!address.loopback()) return address.ip_address();
}
return interfaces.empty() ? "127.0.0.1" : interfaces.front().ip_address();
}
void populateDescriptor(const config::QuicEdgeConfig& config,
const std::string& node_id,
const std::string& boot_id,
const std::string& software_version,
v1::NodeDescriptor* descriptor)
{
if (!descriptor) return;
descriptor->set_node_id(node_id);
descriptor->set_boot_id(boot_id);
descriptor->set_software_version(software_version);
const auto interfaces =
collectInterfaceAddresses(config.include_loopback_interfaces());
for (const auto& address : interfaces) {
*descriptor->add_local_interfaces() = address;
}
auto* endpoint = descriptor->mutable_grpc_endpoint();
endpoint->set_host(advertisedGrpcHost(config, interfaces));
endpoint->set_port(config.grpc_endpoint_port());
endpoint->set_tls(config.grpc_endpoint_tls());
}
void populateHeartbeatNetwork(const config::QuicEdgeConfig& config,
v1::NodeHeartbeat* heartbeat)
{
if (!heartbeat) return;
const auto interfaces =
collectInterfaceAddresses(config.include_loopback_interfaces());
for (const auto& address : interfaces) {
*heartbeat->add_local_interfaces() = address;
}
auto* endpoint = heartbeat->mutable_grpc_endpoint();
endpoint->set_host(advertisedGrpcHost(config, interfaces));
endpoint->set_port(config.grpc_endpoint_port());
endpoint->set_tls(config.grpc_endpoint_tls());
}
v1::DeviceKind toProtoDeviceKind(const device::DeviceKind kind)
{
switch (kind) {
case device::DeviceKind::AGV: return v1::DEVICE_KIND_AGV;
case device::DeviceKind::Arm: return v1::DEVICE_KIND_ARM;
case device::DeviceKind::Battery: return v1::DEVICE_KIND_BATTERY;
case device::DeviceKind::BioHead: return v1::DEVICE_KIND_BIO_HEAD;
case device::DeviceKind::Camera: return v1::DEVICE_KIND_CAMERA;
case device::DeviceKind::CanBus: return v1::DEVICE_KIND_CAN_BUS;
case device::DeviceKind::DexHand: return v1::DEVICE_KIND_DEX_HAND;
case device::DeviceKind::Gripper: return v1::DEVICE_KIND_GRIPPER;
case device::DeviceKind::Microphone: return v1::DEVICE_KIND_MICROPHONE;
case device::DeviceKind::Motor: return v1::DEVICE_KIND_MOTOR;
case device::DeviceKind::MotorSystem:
return v1::DEVICE_KIND_MOTOR_SYSTEM;
case device::DeviceKind::Robot: return v1::DEVICE_KIND_ROBOT;
case device::DeviceKind::Speaker: return v1::DEVICE_KIND_SPEAKER;
case device::DeviceKind::Unknown: break;
}
return v1::DEVICE_KIND_UNSPECIFIED;
}
v1::ManagedDeviceState toProtoManagedDeviceState(
const device::ManagedDeviceState state)
{
switch (state) {
case device::ManagedDeviceState::Disabled:
return v1::MANAGED_DEVICE_STATE_DISABLED;
case device::ManagedDeviceState::Initializing:
return v1::MANAGED_DEVICE_STATE_INITIALIZING;
case device::ManagedDeviceState::Registered:
return v1::MANAGED_DEVICE_STATE_REGISTERED;
case device::ManagedDeviceState::Ready:
return v1::MANAGED_DEVICE_STATE_READY;
case device::ManagedDeviceState::Running:
return v1::MANAGED_DEVICE_STATE_RUNNING;
case device::ManagedDeviceState::Stopped:
return v1::MANAGED_DEVICE_STATE_STOPPED;
case device::ManagedDeviceState::Error:
return v1::MANAGED_DEVICE_STATE_ERROR;
case device::ManagedDeviceState::Unknown:
break;
}
return v1::MANAGED_DEVICE_STATE_UNSPECIFIED;
}
v1::DeviceHealthStatus toProtoDeviceHealthStatus(
const device::DeviceHealthState state)
{
switch (state) {
case device::DeviceHealthState::Healthy:
return v1::DEVICE_HEALTH_STATUS_HEALTHY;
case device::DeviceHealthState::Degraded:
return v1::DEVICE_HEALTH_STATUS_DEGRADED;
case device::DeviceHealthState::Fault:
return v1::DEVICE_HEALTH_STATUS_FAULT;
case device::DeviceHealthState::Unknown:
break;
}
return v1::DEVICE_HEALTH_STATUS_UNSPECIFIED;
}
void populateDeviceManagerSnapshot(
const device::DeviceManagerSnapshot& source,
const std::uint64_t sampled_at_unix_ms,
v1::DeviceManagerSnapshot* destination)
{
if (!destination) return;
destination->set_manager_name(source.name);
destination->set_manager_version(source.version);
destination->set_manager_description(source.description);
destination->set_sampled_at_unix_ms(sampled_at_unix_ms);
std::vector<const device::ManagedDeviceSnapshot*> ordered_devices;
ordered_devices.reserve(source.devices.size());
for (const auto& source_device : source.devices) {
// DeviceManager keeps disabled entries for local configuration and
// diagnostics, but the platform heartbeat only advertises devices
// that are enabled on this edge node. Enabled entries remain visible
// even when creation, initialization or start has failed.
if (!source_device.enabled) continue;
ordered_devices.push_back(&source_device);
}
std::sort(
ordered_devices.begin(), ordered_devices.end(),
[](const auto* lhs, const auto* rhs) {
if (lhs->id != rhs->id) return lhs->id < rhs->id;
if (lhs->kind != rhs->kind) return lhs->kind < rhs->kind;
return lhs->type_name < rhs->type_name;
});
for (const auto* source_device : ordered_devices) {
auto* destination_device = destination->add_devices();
destination_device->set_device_id(source_device->id);
destination_device->set_kind(toProtoDeviceKind(source_device->kind));
destination_device->set_type_name(source_device->type_name);
destination_device->set_enabled(source_device->enabled);
destination_device->set_manager_state(
toProtoManagedDeviceState(source_device->state));
destination_device->set_health(
toProtoDeviceHealthStatus(source_device->health.state));
destination_device->set_has_error(source_device->abnormal);
destination_device->set_error_message(boundedDeviceError(
source_device->error_message.empty()
? source_device->health.error_message
: source_device->error_message));
destination_device->set_status_updated_at_unix_ms(
source_device->status_updated_at_unix_ms);
}
}
v1::MediaKind toProtoKind(const media::MediaKind kind)
{
switch (kind) {
case media::MediaKind::VIDEO: return v1::MEDIA_KIND_VIDEO;
case media::MediaKind::AUDIO: return v1::MEDIA_KIND_AUDIO;
case media::MediaKind::UNKNOWN: break;
}
return v1::MEDIA_KIND_UNSPECIFIED;
}
media::MediaKind expectedKind(const config::QuicEdgeTrackConfig& track)
{
switch (track.source_kind()) {
case config::QuicEdgeTrackConfig::SOURCE_KIND_CAMERA:
return media::MediaKind::VIDEO;
case config::QuicEdgeTrackConfig::SOURCE_KIND_MICROPHONE:
return media::MediaKind::AUDIO;
case config::QuicEdgeTrackConfig::SOURCE_KIND_UNSPECIFIED:
break;
}
return media::MediaKind::UNKNOWN;
}
std::string sourceTrackId(const config::QuicEdgeTrackConfig& track)
{
if (!track.source_track_id().empty()) {
return track.source_track_id();
}
switch (track.source_kind()) {
case config::QuicEdgeTrackConfig::SOURCE_KIND_CAMERA:
return track.device_id() + "/video/color";
case config::QuicEdgeTrackConfig::SOURCE_KIND_MICROPHONE:
return track.device_id() + "/audio/main";
case config::QuicEdgeTrackConfig::SOURCE_KIND_UNSPECIFIED:
break;
}
return {};
}
std::size_t reservedControlSendSlots(const std::size_t queue_depth)
{
if (queue_depth <= 1U) return 0U;
return std::min<std::size_t>(
kMaximumReservedControlSends,
std::max<std::size_t>(1U, queue_depth / 16U));
}
} // namespace
const char* toString(const QuicEdgeServiceState state)
{
switch (state) {
case QuicEdgeServiceState::UNINITIALIZED: return "UNINITIALIZED";
case QuicEdgeServiceState::DISABLED: return "DISABLED";
case QuicEdgeServiceState::STOPPED: return "STOPPED";
case QuicEdgeServiceState::CONNECTING: return "CONNECTING";
case QuicEdgeServiceState::REGISTERING: return "REGISTERING";
case QuicEdgeServiceState::ONLINE: return "ONLINE";
case QuicEdgeServiceState::BACKOFF: return "BACKOFF";
case QuicEdgeServiceState::FAILED: return "FAILED";
}
return "UNKNOWN";
}
QuicEdgeService::QuicEdgeService(config::QuicEdgeConfig config,
std::unique_ptr<QuicTransport> transport,
media::MediaSourceHub& media_hub,
DeviceSnapshotProvider device_snapshot_provider)
: config_(std::move(config)),
transport_(std::move(transport)),
media_hub_(&media_hub),
device_snapshot_provider_(std::move(device_snapshot_provider))
{
initializeIdentity();
}
QuicEdgeService::~QuicEdgeService()
{
stop();
}
bool QuicEdgeService::validateConfig(const config::QuicEdgeConfig& config,
std::string* error)
{
if (config.id().empty()) {
setError(error, "QUIC edge task id is empty");
return false;
}
if (!config.enable()) {
return true;
}
if (config.server_host().empty()) {
setError(error, "QUIC edge server_host is empty");
return false;
}
if (config.server_port() == 0U || config.server_port() > 65535U) {
setError(error, "QUIC edge server_port must be in [1, 65535]");
return false;
}
if (config.alpn().empty() || config.alpn().size() > 255U) {
setError(error, "QUIC edge ALPN must contain 1 to 255 bytes");
return false;
}
if (!config.has_tls()) {
setError(error, "QUIC edge TLS configuration is missing");
return false;
}
if (!config.tls().allow_insecure() &&
(config.tls().ca_file().empty() || config.tls().server_name().empty())) {
setError(error, "secure QUIC edge TLS requires ca_file and server_name");
return false;
}
if (!config.tls().allow_insecure() &&
config.tls().server_name() != config.server_host()) {
setError(error,
"QUIC edge v1 requires tls.server_name to match server_host");
return false;
}
const bool has_certificate = !config.tls().certificate_file().empty();
const bool has_private_key = !config.tls().private_key_file().empty();
if (has_certificate != has_private_key) {
setError(error, "QUIC edge client certificate and private key must be configured together");
return false;
}
if (!config.has_reconnect() || config.reconnect().initial_delay_ms() == 0U ||
config.reconnect().maximum_delay_ms() < config.reconnect().initial_delay_ms() ||
config.reconnect().connect_timeout_ms() == 0U) {
setError(error, "QUIC edge reconnect configuration is invalid");
return false;
}
if (!std::isfinite(config.reconnect().multiplier()) ||
config.reconnect().multiplier() < 1.0 ||
config.reconnect().jitter_percent() > 100U) {
setError(error, "QUIC edge reconnect multiplier or jitter is invalid");
return false;
}
if (config.maximum_datagram_bytes() < kMinimumDatagramBytes ||
config.maximum_datagram_bytes() > kMaximumDatagramBytes) {
setError(error, "QUIC edge maximum_datagram_bytes is outside the v1 range");
return false;
}
if (config.maximum_control_frame_bytes() == 0U ||
config.maximum_control_frame_bytes() > kMaximumControlBytes) {
setError(error, "QUIC edge maximum_control_frame_bytes is invalid");
return false;
}
if (config.maximum_frame_bytes() == 0U ||
config.maximum_frame_bytes() > kMaximumConfiguredFrameBytes) {
setError(error, "QUIC edge maximum_frame_bytes is invalid");
return false;
}
const std::uint64_t fragment_payload =
config.maximum_datagram_bytes() - kDatagramHeaderBytes;
if (config.maximum_frame_bytes() > fragment_payload * 65535U) {
setError(error, "QUIC edge maximum_frame_bytes exceeds v1 fragment capacity");
return false;
}
if (config.datagram_send_queue_depth() < 2U) {
setError(error, "QUIC edge datagram_send_queue_depth must be at least 2");
return false;
}
const bool has_enabled_media = std::any_of(
config.tracks().begin(), config.tracks().end(),
[](const config::QuicEdgeTrackConfig& track) { return track.enable(); });
const std::uint64_t media_send_slots =
config.datagram_send_queue_depth() -
reservedControlSendSlots(config.datagram_send_queue_depth());
if (has_enabled_media &&
config.maximum_frame_bytes() > fragment_payload * media_send_slots) {
setError(error,
"maximum_frame_bytes exceeds the atomic DATAGRAM send-queue capacity");
return false;
}
if (config.media_poll_interval_ms() == 0U ||
config.media_poll_interval_ms() > 1000U) {
setError(error, "QUIC edge media_poll_interval_ms must be in [1, 1000]");
return false;
}
if (config.grpc_endpoint_port() == 0U ||
config.grpc_endpoint_port() > 65535U) {
setError(error, "advertised gRPC endpoint port must be in [1, 65535]");
return false;
}
if (config.heartbeat_interval_ms() < kMinimumHeartbeatIntervalMs ||
config.heartbeat_interval_ms() > kMaximumHeartbeatIntervalMs) {
setError(error, "heartbeat_interval_ms must be in [250, 3600000]");
return false;
}
if (config.control_response_timeout_ms() == 0U ||
config.control_response_timeout_ms() > kMaximumHeartbeatIntervalMs) {
setError(error, "control_response_timeout_ms must be in [1, 3600000]");
return false;
}
std::unordered_set<std::uint32_t> wire_track_ids;
std::unordered_set<std::string> source_track_ids;
for (const auto& track : config.tracks()) {
if (!track.enable()) {
continue;
}
if (track.track_id() == 0U ||
!wire_track_ids.insert(track.track_id()).second) {
setError(error, "enabled QUIC edge track IDs must be unique and non-zero");
return false;
}
if (expectedKind(track) == media::MediaKind::UNKNOWN) {
setError(error, "enabled QUIC edge track has unsupported source_kind");
return false;
}
if (track.device_id().empty()) {
setError(error, "enabled QUIC edge track has empty device_id");
return false;
}
const std::string source_id = sourceTrackId(track);
if (source_id.empty() || !source_track_ids.insert(source_id).second) {
setError(error, "enabled QUIC edge source track IDs must be unique and non-empty");
return false;
}
if (track.max_frame_bytes() > config.maximum_frame_bytes()) {
setError(error, "track max_frame_bytes exceeds the global frame limit");
return false;
}
}
return true;
}
bool QuicEdgeService::initialize(std::string* error)
{
std::lock_guard lifecycle_lock(lifecycle_mutex_);
if (worker_.joinable()) {
setError(error, "cannot initialize QUIC edge service while it is running");
return false;
}
std::string validation_error;
if (!validateConfig(config_, &validation_error)) {
setState(QuicEdgeServiceState::FAILED, validation_error);
setError(error, validation_error);
return false;
}
if (!config_.enable()) {
setState(QuicEdgeServiceState::DISABLED);
return true;
}
if (!transport_ || !media_hub_) {
const std::string message = "QUIC edge transport or MediaSourceHub is null";
setState(QuicEdgeServiceState::FAILED, message);
setError(error, message);
return false;
}
if (using_default_transport_ && !hasCompiledMsQuicSupport()) {
const std::string message =
"QUIC edge is enabled but CMVR_HAS_MSQUIC is not compiled";
setState(QuicEdgeServiceState::FAILED, message);
setError(error, message);
return false;
}
setState(QuicEdgeServiceState::STOPPED);
return true;
}
bool QuicEdgeService::start(std::string* error)
{
std::lock_guard lifecycle_lock(lifecycle_mutex_);
{
std::lock_guard lock(mutex_);
if (state_ == QuicEdgeServiceState::DISABLED) {
return true;
}
if (state_ == QuicEdgeServiceState::UNINITIALIZED) {
setError(error, "QUIC edge service is not initialized");
return false;
}
if (state_ == QuicEdgeServiceState::FAILED) {
setError(error, last_error_);
return false;
}
if (worker_.joinable()) {
return true;
}
stop_requested_ = false;
}
try {
worker_ = std::thread(&QuicEdgeService::run, this);
} catch (const std::exception& exception) {
const std::string message =
std::string("failed to start QUIC edge worker: ") + exception.what();
setState(QuicEdgeServiceState::FAILED, message);
setError(error, message);
return false;
}
return true;
}
void QuicEdgeService::stop()
{
std::lock_guard lifecycle_lock(lifecycle_mutex_);
{
std::lock_guard lock(mutex_);
stop_requested_ = true;
media_stop_requested_ = true;
}
stop_cv_.notify_all();
media_stop_cv_.notify_all();
if (transport_) {
transport_->disconnect();
}
if (worker_.joinable()) {
worker_.join();
}
resetConnectionStatus();
std::lock_guard lock(mutex_);
if (state_ != QuicEdgeServiceState::UNINITIALIZED &&
state_ != QuicEdgeServiceState::DISABLED &&
state_ != QuicEdgeServiceState::FAILED) {
state_ = QuicEdgeServiceState::STOPPED;
}
}
QuicEdgeServiceState QuicEdgeService::state() const
{
std::lock_guard lock(mutex_);
return state_;
}
std::string QuicEdgeService::lastError() const
{
std::lock_guard lock(mutex_);
return last_error_;
}
QuicEdgeStats QuicEdgeService::stats() const
{
std::lock_guard lock(mutex_);
return stats_;
}
QuicEdgeStatus QuicEdgeService::status() const
{
std::lock_guard lock(mutex_);
QuicEdgeStatus result;
result.registered = registered_;
result.node_id = node_id_;
result.boot_id = boot_id_;
result.session_id = session_id_;
result.observed_source_ip = observed_source_ip_;
result.last_media_error = last_media_error_;
result.heartbeat_sequence = heartbeat_sequence_;
result.last_heartbeat_ack_unix_ms = last_heartbeat_ack_unix_ms_;
result.active_media_tracks = active_media_tracks_;
return result;
}
void QuicEdgeService::run()
{
try {
auto backoff = std::chrono::milliseconds(config_.reconnect().initial_delay_ms());
while (true) {
{
std::lock_guard lock(mutex_);
if (stop_requested_) break;
++stats_.connection_attempts;
}
resetConnectionStatus();
setState(QuicEdgeServiceState::CONNECTING);
CMVR_LOG(INFO) << "[QuicEdgeService] connecting"
<< ", target=" << config_.server_host() << ':'
<< config_.server_port()
<< ", tls_server_name=" << config_.tls().server_name()
<< ", alpn=" << config_.alpn()
<< ", node_id=" << node_id_;
std::string error;
if (!transport_->connect(
config_,
std::chrono::milliseconds(config_.reconnect().connect_timeout_ms()),
&error)) {
CMVR_LOG(ERROR) << "[QuicEdgeService] connect failed"
<< ", target=" << config_.server_host() << ':'
<< config_.server_port()
<< ", error="
<< (error.empty() ? "QUIC connect failed" : error)
<< ", backoff_ms=" << backoff.count();
setState(QuicEdgeServiceState::BACKOFF,
error.empty() ? "QUIC connect failed" : error);
if (waitForStop(jittered(backoff))) break;
backoff = nextBackoff(backoff);
continue;
}
std::uint64_t acknowledged_heartbeats_at_connect = 0U;
{
std::lock_guard lock(mutex_);
++stats_.successful_connections;
if (stats_.successful_connections > 1U) ++stats_.reconnects;
acknowledged_heartbeats_at_connect =
stats_.heartbeats_acknowledged;
}
session_epoch_ = nextSessionEpoch();
control_message_sequence_ = 0U;
inbound_message_sequence_ = 0U;
has_inbound_message_sequence_ = false;
ControlFrameDecoder decoder(config_.maximum_control_frame_bytes());
setState(QuicEdgeServiceState::REGISTERING);
CMVR_LOG(INFO) << "[QuicEdgeService] connected, registering node"
<< ", target=" << config_.server_host() << ':'
<< config_.server_port()
<< ", node_id=" << node_id_;
if (!performRegistration(&decoder, &error)) {
transport_->disconnect();
{
std::lock_guard lock(mutex_);
if (stop_requested_) break;
}
CMVR_LOG(ERROR) << "[QuicEdgeService] registration failed"
<< ", target=" << config_.server_host() << ':'
<< config_.server_port()
<< ", error=" << error
<< ", backoff_ms=" << backoff.count();
setState(QuicEdgeServiceState::BACKOFF, error);
if (waitForStop(jittered(backoff))) break;
backoff = nextBackoff(backoff);
continue;
}
setState(QuicEdgeServiceState::ONLINE);
CMVR_LOG(INFO) << "[QuicEdgeService] online"
<< ", target=" << config_.server_host() << ':'
<< config_.server_port()
<< ", node_id=" << node_id_
<< ", session=" << session_id_;
if (!startMediaWorker(&error)) {
transport_->disconnect();
{
std::lock_guard lock(mutex_);
if (stop_requested_) break;
}
setState(QuicEdgeServiceState::BACKOFF, error);
if (waitForStop(jittered(backoff))) break;
backoff = nextBackoff(backoff);
continue;
}
bool connection_failed = false;
bool connection_became_healthy = false;
while (true) {
{
std::lock_guard lock(mutex_);
if (stop_requested_) break;
}
bool received_control = false;
if (!receiveAndDispatchControl(
&decoder, std::chrono::milliseconds(0),
&received_control, &error)) {
connection_failed = true;
break;
}
if (!transport_->isConnected()) {
// receiveControl drains bytes already delivered by MsQuic
// before it reports the disconnect on the next iteration.
continue;
}
if (!connection_became_healthy) {
std::lock_guard lock(mutex_);
if (stats_.heartbeats_acknowledged >
acknowledged_heartbeats_at_connect) {
backoff = std::chrono::milliseconds(
config_.reconnect().initial_delay_ms());
connection_became_healthy = true;
}
}
const auto now = std::chrono::steady_clock::now();
std::uint64_t outstanding_sequence = 0U;
std::chrono::steady_clock::time_point heartbeat_deadline;
{
std::lock_guard lock(mutex_);
outstanding_sequence = outstanding_heartbeat_sequence_;
heartbeat_deadline = heartbeat_deadline_;
}
if (outstanding_sequence != 0U && now >= heartbeat_deadline) {
{
std::lock_guard lock(mutex_);
++stats_.heartbeat_timeouts;
}
error = "QUIC node heartbeat acknowledgement timed out";
connection_failed = true;
break;
}
if (outstanding_sequence == 0U && now >= next_heartbeat_) {
std::uint64_t sequence = 0U;
{
std::lock_guard lock(mutex_);
if (heartbeat_sequence_ ==
std::numeric_limits<std::uint64_t>::max()) {
heartbeat_sequence_ = 0U;
}
sequence = ++heartbeat_sequence_;
}
if (!sendHeartbeat(sequence, &error)) {
connection_failed = true;
break;
}
const auto heartbeat_sent_at =
std::chrono::steady_clock::now();
{
std::lock_guard lock(mutex_);
outstanding_heartbeat_sequence_ = sequence;
heartbeat_deadline_ =
heartbeat_sent_at + std::chrono::milliseconds(
config_.control_response_timeout_ms());
++stats_.heartbeats_sent;
}
next_heartbeat_ =
heartbeat_sent_at + std::chrono::milliseconds(
effective_heartbeat_interval_ms_);
}
{
std::lock_guard lock(mutex_);
if (media_connection_failed_) {
error = media_connection_error_.empty()
? "QUIC media worker reported a connection failure"
: media_connection_error_;
connection_failed = true;
}
}
if (connection_failed) break;
if (!receiveAndDispatchControl(
&decoder, kControlPollInterval,
&received_control, &error)) {
connection_failed = true;
break;
}
}
stopMediaWorker();
transport_->disconnect();
resetConnectionStatus();
{
std::lock_guard lock(mutex_);
if (stop_requested_) break;
}
setState(QuicEdgeServiceState::BACKOFF,
error.empty() ? "QUIC edge connection closed" : error);
if (waitForStop(jittered(backoff))) break;
backoff = nextBackoff(backoff);
}
stopMediaWorker();
transport_->disconnect();
resetConnectionStatus();
setState(QuicEdgeServiceState::STOPPED);
} catch (const std::exception& exception) {
stopMediaWorker();
transport_->disconnect();
resetConnectionStatus();
setState(QuicEdgeServiceState::FAILED,
std::string("QUIC edge worker exception: ") + exception.what());
} catch (...) {
stopMediaWorker();
transport_->disconnect();
resetConnectionStatus();
setState(QuicEdgeServiceState::FAILED, "unknown QUIC edge worker exception");
}
}
bool QuicEdgeService::startMediaWorker(std::string* error)
{
if (!hasEnabledMediaTracks()) {
std::lock_guard lock(mutex_);
active_media_tracks_ = 0U;
return true;
}
if (media_worker_.joinable()) {
setError(error, "QUIC media worker is already running");
return false;
}
{
std::lock_guard lock(mutex_);
if (stop_requested_) {
setError(error, "QUIC edge service is stopping");
return false;
}
media_stop_requested_ = false;
media_connection_failed_ = false;
media_connection_error_.clear();
active_media_tracks_ = 0U;
}
try {
media_worker_ = std::thread(&QuicEdgeService::runMedia, this);
} catch (const std::exception& exception) {
const std::string message =
std::string("failed to start QUIC media worker: ") + exception.what();
recordMediaConnectionFailure(message);
setError(error, message);
return false;
}
return true;
}
void QuicEdgeService::stopMediaWorker()
{
{
std::lock_guard lock(mutex_);
media_stop_requested_ = true;
}
media_stop_cv_.notify_all();
if (media_worker_.joinable()) media_worker_.join();
std::lock_guard lock(mutex_);
active_media_tracks_ = 0U;
}
void QuicEdgeService::recordMediaConnectionFailure(const std::string& error)
{
{
std::lock_guard lock(mutex_);
if (!media_connection_failed_) {
media_connection_error_ = error.empty()
? "unknown QUIC media connection failure" : error;
}
media_connection_failed_ = true;
media_stop_requested_ = true;
last_media_error_ = media_connection_error_;
}
media_stop_cv_.notify_all();
}
void QuicEdgeService::runMedia()
{
std::vector<ActiveTrack> tracks;
try {
next_media_source_retry_ = std::chrono::steady_clock::now();
while (true) {
{
std::lock_guard lock(mutex_);
if (stop_requested_ || media_stop_requested_) break;
}
if (!transport_->isConnected()) break;
const auto now = std::chrono::steady_clock::now();
std::string error;
if (now >= next_media_source_retry_) {
tracks.erase(
std::remove_if(tracks.begin(), tracks.end(),
[](const ActiveTrack& track) {
return !track.subscription.valid();
}),
tracks.end());
if (!openMediaSession(&tracks, &error)) {
recordMediaConnectionFailure(error);
break;
}
next_media_source_retry_ = now + kMediaSourceRetryInterval;
}
bool sent_anything = false;
bool failed = false;
for (auto& track : tracks) {
if (!processTrack(&track, &sent_anything, &error)) {
failed = true;
break;
}
}
if (failed) {
recordMediaConnectionFailure(error);
break;
}
tracks.erase(
std::remove_if(tracks.begin(), tracks.end(),
[](const ActiveTrack& track) {
return !track.subscription.valid();
}),
tracks.end());
{
std::lock_guard lock(mutex_);
active_media_tracks_ = tracks.size();
}
if (!sent_anything) {
const auto poll_wait = tracks.empty()
? kControlPollInterval
: std::chrono::milliseconds(config_.media_poll_interval_ms());
std::unique_lock lock(mutex_);
media_stop_cv_.wait_for(lock, poll_wait, [this]() {
return stop_requested_ || media_stop_requested_;
});
}
}
} catch (const std::exception& exception) {
recordMediaConnectionFailure(
std::string("QUIC media worker exception: ") + exception.what());
} catch (...) {
recordMediaConnectionFailure("unknown QUIC media worker exception");
}
tracks.clear();
std::lock_guard lock(mutex_);
active_media_tracks_ = 0U;
}
bool QuicEdgeService::performRegistration(ControlFrameDecoder* decoder,
std::string* error)
{
if (!decoder) {
setError(error, "control frame decoder is null");
return false;
}
if (!sendNodeRegistration(error)) return false;
const auto deadline = std::chrono::steady_clock::now() +
std::chrono::milliseconds(config_.control_response_timeout_ms());
while (transport_->isConnected()) {
{
std::lock_guard lock(mutex_);
if (stop_requested_) {
setError(error, "QUIC edge service is stopping");
return false;
}
if (registered_) return true;
}
const auto now = std::chrono::steady_clock::now();
if (now >= deadline) {
setError(error, "QUIC node registration response timed out");
return false;
}
const auto remaining = std::chrono::duration_cast<std::chrono::milliseconds>(
deadline - now);
bool received = false;
if (!receiveAndDispatchControl(
decoder, std::min(remaining, kControlPollInterval),
&received, error)) {
return false;
}
}
setError(error, "QUIC connection closed during node registration");
return false;
}
bool QuicEdgeService::sendNodeRegistration(std::string* error)
{
std::lock_guard control_lock(control_send_mutex_);
v1::EdgeControlEnvelope envelope;
envelope.set_protocol_version(kProtocolVersion);
envelope.set_message_sequence(control_message_sequence_++);
auto* request = envelope.mutable_node_register_request();
populateDescriptor(config_, node_id_, boot_id_, software_version_,
request->mutable_node());
request->set_sent_at_unix_ms(unixTimeMs());
std::string serialized;
if (!envelope.SerializeToString(&serialized)) {
setError(error, "failed to serialize NodeRegisterRequest");
return false;
}
if (!sendControlEnvelope(serialized, error)) return false;
const auto& node = request->node();
CMVR_LOG(INFO) << "[QuicEdgeService] sent NodeRegisterRequest"
<< ", message_sequence=" << envelope.message_sequence()
<< ", payload_bytes=" << serialized.size()
<< ", node_id=" << node.node_id()
<< ", boot_id=" << node.boot_id()
<< ", grpc_endpoint=" << node.grpc_endpoint().host()
<< ':' << node.grpc_endpoint().port()
<< ", grpc_tls=" << node.grpc_endpoint().tls()
<< ", interfaces=" << node.local_interfaces_size()
<< ", sent_at_unix_ms=" << request->sent_at_unix_ms();
std::lock_guard lock(mutex_);
++stats_.registrations_sent;
return true;
}
bool QuicEdgeService::sendHeartbeat(const std::uint64_t sequence,
std::string* error)
{
std::optional<device::DeviceManagerSnapshot> device_manager_snapshot;
if (device_snapshot_provider_) {
try {
device_manager_snapshot = device_snapshot_provider_();
} catch (const std::exception& exception) {
setError(
error,
std::string("failed to snapshot DeviceManager for heartbeat: ") +
exception.what());
return false;
} catch (...) {
setError(error,
"failed to snapshot DeviceManager for heartbeat: "
"unknown exception");
return false;
}
}
std::lock_guard control_lock(control_send_mutex_);
std::string session_id;
{
std::lock_guard lock(mutex_);
session_id = session_id_;
}
if (session_id.empty()) {
setError(error, "cannot send heartbeat before node registration");
return false;
}
v1::EdgeControlEnvelope envelope;
envelope.set_protocol_version(kProtocolVersion);
envelope.set_message_sequence(control_message_sequence_++);
auto* heartbeat = envelope.mutable_node_heartbeat();
heartbeat->set_node_id(node_id_);
heartbeat->set_boot_id(boot_id_);
heartbeat->set_session_id(session_id);
heartbeat->set_sequence(sequence);
const std::uint64_t sent_at_unix_ms = unixTimeMs();
heartbeat->set_sent_at_unix_ms(sent_at_unix_ms);
heartbeat->set_software_version(software_version_);
populateHeartbeatNetwork(config_, heartbeat);
if (device_manager_snapshot.has_value()) {
populateDeviceManagerSnapshot(
*device_manager_snapshot, sent_at_unix_ms,
heartbeat->mutable_device_manager());
}
std::string serialized;
if (!envelope.SerializeToString(&serialized)) {
setError(error, "failed to serialize NodeHeartbeat");
return false;
}
if (!sendControlEnvelope(serialized, error)) return false;
// CMVR_LOG(INFO) << "[QuicEdgeService] sent NodeHeartbeat"
// << ", message_sequence=" << envelope.message_sequence()
// << ", payload_bytes=" << serialized.size()
// << ", session=" << session_id
// << ", heartbeat_sequence=" << sequence
// << ", grpc_endpoint=" << heartbeat->grpc_endpoint().host()
// << ':' << heartbeat->grpc_endpoint().port()
// << ", interfaces=" << heartbeat->local_interfaces_size()
// << ", devices=" << heartbeat->device_manager().devices_size()
// << ", sent_at_unix_ms=" << sent_at_unix_ms;
return true;
}
bool QuicEdgeService::receiveAndDispatchControl(
ControlFrameDecoder* decoder,
const std::chrono::milliseconds timeout,
bool* received,
std::string* error)
{
if (!decoder || !received) {
setError(error, "control receive arguments are null");
return false;
}
*received = false;
std::vector<std::uint8_t> chunk;
const auto result = transport_->receiveControl(&chunk, timeout, error);
if (result == TransportReceiveResult::TIMEOUT) return true;
if (result == TransportReceiveResult::DISCONNECTED) {
if (error && error->empty()) *error = "QUIC control stream disconnected";
return false;
}
if (result == TransportReceiveResult::ERROR) {
if (error && error->empty()) *error = "QUIC control stream receive failed";
return false;
}
if (chunk.empty()) {
setError(error, "QUIC control stream returned an empty data chunk");
return false;
}
std::vector<std::vector<std::uint8_t>> frames;
if (!decoder->push(chunk, &frames, error)) return false;
*received = true;
for (const auto& frame : frames) {
if (!dispatchControlFrame(frame, error)) return false;
}
return true;
}
bool QuicEdgeService::dispatchControlFrame(
const std::vector<std::uint8_t>& frame,
std::string* error)
{
v1::EdgeControlEnvelope envelope;
if (!envelope.ParseFromArray(frame.data(), static_cast<int>(frame.size()))) {
setError(error, "failed to parse EdgeControlEnvelope");
return false;
}
if (envelope.protocol_version() != kProtocolVersion) {
setError(error, "unsupported QUIC edge protocol version");
return false;
}
if (has_inbound_message_sequence_ &&
envelope.message_sequence() <= inbound_message_sequence_) {
setError(error, "QUIC control message sequence did not increase");
return false;
}
inbound_message_sequence_ = envelope.message_sequence();
has_inbound_message_sequence_ = true;
if (envelope.has_node_register_response()) {
const auto& response = envelope.node_register_response();
if (!response.accepted()) {
{
std::lock_guard lock(mutex_);
++stats_.registrations_rejected;
}
setError(error, response.message().empty()
? "gateway rejected node registration"
: response.message());
return false;
}
if (response.session_id().empty()) {
setError(error, "gateway accepted registration without a session ID");
return false;
}
const std::uint32_t requested_interval =
response.heartbeat_interval_ms() == 0U
? config_.heartbeat_interval_ms()
: response.heartbeat_interval_ms();
{
std::lock_guard lock(mutex_);
if (registered_) {
setError(error, "duplicate node registration response");
return false;
}
registered_ = true;
session_id_ = response.session_id();
observed_source_ip_ = response.observed_source_ip();
effective_heartbeat_interval_ms_ = std::clamp(
requested_interval, kMinimumHeartbeatIntervalMs,
kMaximumHeartbeatIntervalMs);
next_heartbeat_ = std::chrono::steady_clock::now();
++stats_.registrations_accepted;
}
return true;
}
if (envelope.has_node_heartbeat_ack()) {
const auto& ack = envelope.node_heartbeat_ack();
std::lock_guard lock(mutex_);
if (!registered_ || outstanding_heartbeat_sequence_ == 0U) {
setError(error, "unexpected heartbeat acknowledgement");
return false;
}
if (!ack.accepted()) {
setError(error, ack.message().empty()
? "gateway rejected node heartbeat"
: ack.message());
return false;
}
if (ack.session_id() != session_id_ ||
ack.acknowledged_sequence() != outstanding_heartbeat_sequence_) {
setError(error, "heartbeat acknowledgement session or sequence mismatch");
return false;
}
outstanding_heartbeat_sequence_ = 0U;
if (!ack.observed_source_ip().empty()) {
observed_source_ip_ = ack.observed_source_ip();
}
last_heartbeat_ack_unix_ms_ = unixTimeMs();
++stats_.heartbeats_acknowledged;
return true;
}
if (envelope.has_protocol_error()) {
const auto& protocol_error = envelope.protocol_error();
{
std::lock_guard lock(mutex_);
++stats_.protocol_errors;
last_error_ = protocol_error.message();
}
if (protocol_error.fatal()) {
setError(error, protocol_error.message().empty()
? "gateway reported a fatal protocol error"
: protocol_error.message());
return false;
}
return true;
}
setError(error, "gateway sent an unexpected QUIC control message");
return false;
}
bool QuicEdgeService::openMediaSession(std::vector<ActiveTrack>* tracks,
std::string* error)
{
if (!tracks) {
setError(error, "active track output is null");
return false;
}
if (!hasEnabledMediaTracks()) {
std::lock_guard lock(mutex_);
active_media_tracks_ = 0U;
return true;
}
if (!media_session_announced_) {
if (!sendSessionOpen(error)) {
if (!transport_->isConnected()) return false;
recordMediaError(error ? *error : "failed to open media session");
if (error) error->clear();
return true;
}
media_session_announced_ = true;
{
std::lock_guard lock(mutex_);
++stats_.media_sessions_opened;
}
}
refreshMediaTracks(tracks);
return true;
}
void QuicEdgeService::refreshMediaTracks(std::vector<ActiveTrack>* tracks)
{
if (!tracks) return;
for (const auto& track_config : config_.tracks()) {
if (!track_config.enable()) continue;
const bool already_active = std::any_of(
tracks->begin(), tracks->end(),
[&track_config](const ActiveTrack& active) {
return active.config.track_id() == track_config.track_id();
});
if (already_active) continue;
ActiveTrack track;
track.config = track_config;
track.source_track_id = sourceTrackId(track_config);
std::string source_error;
if (!ensureSourceRegistered(track_config, track.source_track_id,
&source_error)) {
recordMediaError(source_error);
continue;
}
track.subscription = media_hub_->subscribe(
track.source_track_id,
media::MediaSourceHub::StartPosition::LATEST_AVAILABLE,
[this] {
std::lock_guard lock(mutex_);
return stop_requested_ || media_stop_requested_;
});
if (!track.subscription.valid()) {
recordMediaError(
"MediaSourceHub source unavailable: " + track.source_track_id);
continue;
}
track.waiting_for_keyframe =
expectedKind(track_config) == media::MediaKind::VIDEO;
track.next_frame_discontinuous = true;
if (track.waiting_for_keyframe) requestKeyFrame(&track);
tracks->push_back(std::move(track));
}
std::lock_guard lock(mutex_);
active_media_tracks_ = tracks->size();
}
bool QuicEdgeService::hasEnabledMediaTracks() const
{
return std::any_of(config_.tracks().begin(), config_.tracks().end(),
[](const config::QuicEdgeTrackConfig& track) {
return track.enable();
});
}
bool QuicEdgeService::ensureSourceRegistered(
const config::QuicEdgeTrackConfig& track,
const std::string& source_track_id,
std::string* error)
{
if (media_hub_->hasSource(source_track_id)) return true;
if (!using_global_media_hub_ || !source_registrar_) {
setError(error, "MediaSourceHub source unavailable: " + source_track_id);
return false;
}
const bool registered = source_registrar_(track, source_track_id, error);
if (!registered && !media_hub_->hasSource(source_track_id)) {
if (!error || error->empty()) {
setError(error,
"failed to register MediaSourceHub source: " + source_track_id);
}
return false;
}
if (!media_hub_->hasSource(source_track_id)) {
setError(error,
"configured source_track_id does not match the device adapter track: " +
source_track_id);
return false;
}
return true;
}
bool QuicEdgeService::sendSessionOpen(std::string* error)
{
std::lock_guard control_lock(control_send_mutex_);
std::string session_id;
{
std::lock_guard lock(mutex_);
session_id = session_id_;
}
if (session_id.empty()) {
setError(error, "cannot open media before node registration");
return false;
}
v1::EdgeControlEnvelope envelope;
envelope.set_protocol_version(kProtocolVersion);
envelope.set_message_sequence(control_message_sequence_++);
auto* open = envelope.mutable_media_session_open();
open->set_node_id(node_id_);
open->set_session_epoch(session_epoch_);
open->set_session_id(session_id);
std::string serialized;
if (!envelope.SerializeToString(&serialized)) {
setError(error, "failed to serialize MediaSessionOpen");
return false;
}
if (!sendControlEnvelope(serialized, error)) return false;
CMVR_LOG(INFO) << "[QuicEdgeService] sent MediaSessionOpen"
<< ", message_sequence=" << envelope.message_sequence()
<< ", payload_bytes=" << serialized.size()
<< ", session=" << session_id
<< ", session_epoch=" << session_epoch_
<< ", node_id=" << node_id_;
return true;
}
bool QuicEdgeService::sendTrackDescription(
const MediaTrackDescription& description,
std::string* error)
{
std::lock_guard control_lock(control_send_mutex_);
v1::EdgeControlEnvelope envelope;
envelope.set_protocol_version(kProtocolVersion);
envelope.set_message_sequence(control_message_sequence_++);
auto* descriptor = envelope.mutable_media_track_descriptor();
descriptor->set_track_id(description.track_id);
descriptor->set_kind(toProtoKind(description.kind));
descriptor->set_device_id(description.source_id);
descriptor->set_codec(codecName(description.codec));
descriptor->set_codec_generation(description.codec_generation);
descriptor->set_source_track_id(description.source_track_id);
descriptor->set_codec_generation_token(description.codec_generation_token);
descriptor->set_payload_format(payloadFormatName(description.payload_format));
descriptor->set_width(description.width);
descriptor->set_height(description.height);
descriptor->set_frames_per_second(description.nominal_rate);
descriptor->set_sample_rate(description.sample_rate);
descriptor->set_channels(description.channels);
if (!description.codec_config.empty()) {
descriptor->set_codec_config(description.codec_config.data(),
description.codec_config.size());
}
std::string serialized;
if (!envelope.SerializeToString(&serialized)) {
setError(error, "failed to serialize MediaTrackDescriptor");
return false;
}
if (!sendControlEnvelope(serialized, error)) return false;
CMVR_LOG(INFO) << "[QuicEdgeService] sent MediaTrackDescriptor"
<< ", message_sequence=" << envelope.message_sequence()
<< ", payload_bytes=" << serialized.size()
<< ", track_id=" << description.track_id
<< ", source_track_id=" << description.source_track_id
<< ", source_id=" << description.source_id
<< ", kind=" << static_cast<int>(description.kind)
<< ", codec=" << codecName(description.codec)
<< ", payload_format=" << payloadFormatName(description.payload_format)
<< ", codec_generation=" << description.codec_generation
<< ", codec_generation_token="
<< description.codec_generation_token
<< ", size=" << description.width << 'x' << description.height
<< ", fps=" << description.nominal_rate
<< ", sample_rate=" << description.sample_rate
<< ", channels=" << description.channels
<< ", codec_config_bytes=" << description.codec_config.size();
return true;
}
bool QuicEdgeService::sendControlEnvelope(const std::string& serialized,
std::string* error)
{
std::vector<std::uint8_t> framed;
if (!ControlFrameEncoder::encode(
reinterpret_cast<const std::uint8_t*>(serialized.data()),
serialized.size(), config_.maximum_control_frame_bytes(),
&framed, error)) {
return false;
}
const auto result = transport_->sendControl(std::move(framed), error);
if (result == TransportSendResult::QUEUED) return true;
if (error && error->empty()) {
*error = result == TransportSendResult::WOULD_BLOCK
? "QUIC reliable control stream is congested"
: "QUIC reliable control stream is disconnected";
}
return false;
}
bool QuicEdgeService::processTrack(ActiveTrack* track,
bool* sent_anything,
std::string* error)
{
if (!track || !sent_anything) {
setError(error, "active track or sent flag is null");
return false;
}
auto read = track->subscription.tryRead();
if (!read) {
if (!track->subscription.valid()) {
recordMediaError(
"MediaSourceHub subscription stopped: " + track->source_track_id);
}
return true;
}
*sent_anything = true;
const media::MediaFramePtr& frame = read->value;
if (!frame || !frame->descriptor ||
frame->descriptor->id != track->source_track_id ||
frame->descriptor->kind != expectedKind(track->config)) {
recordMediaError("MediaSourceHub returned an invalid or mismatched frame");
track->next_frame_discontinuous = true;
return true;
}
bool discontinuity = track->next_frame_discontinuous ||
frame->discontinuity ||
read->generation_changed ||
read->dropped_since_last_read != 0U;
if (read->dropped_since_last_read != 0U) {
std::lock_guard lock(mutex_);
stats_.frames_dropped_source += read->dropped_since_last_read;
}
const auto description = describeTrack(track->config.track_id(),
*frame->descriptor);
const bool descriptor_changed =
!track->last_description || *track->last_description != description;
if (track->last_description && descriptor_changed &&
description.codec_generation < track->last_description->codec_generation) {
recordMediaError("MediaSourceHub descriptor generation regressed");
track->next_frame_discontinuous = true;
return true;
}
if (descriptor_changed) {
if (!sendTrackDescription(description, error)) {
if (!transport_->isConnected()) return false;
recordMediaError(error ? *error : "failed to send media track descriptor");
if (error) error->clear();
track->next_frame_discontinuous = true;
return true;
}
track->last_description = description;
discontinuity = true;
if (description.kind == media::MediaKind::VIDEO) {
track->waiting_for_keyframe = true;
track->keyframe_requested = false;
requestKeyFrame(track);
}
}
if (description.kind == media::MediaKind::VIDEO && discontinuity) {
if (!track->waiting_for_keyframe) {
track->keyframe_requested = false;
}
track->waiting_for_keyframe = true;
requestKeyFrame(track);
}
if (frame->descriptor->kind == media::MediaKind::VIDEO &&
track->waiting_for_keyframe && !frame->key_frame) {
track->next_frame_discontinuous = discontinuity;
requestKeyFrame(track);
std::lock_guard lock(mutex_);
++stats_.frames_skipped_waiting_keyframe;
return true;
}
if (frame->empty() || frame->size() > maximumFrameBytes(*track)) {
const auto drained = shedBufferedFrames(track);
track->next_frame_discontinuous = true;
if (frame->descriptor->kind == media::MediaKind::VIDEO) {
track->waiting_for_keyframe = true;
track->keyframe_requested = false;
requestKeyFrame(track);
}
std::lock_guard lock(mutex_);
stats_.frames_dropped_oversize += 1U + drained;
return true;
}
const std::size_t peer_maximum = transport_->maximumDatagramBytes();
if (peer_maximum <= kDatagramHeaderBytes) {
const auto drained = shedBufferedFrames(track);
track->next_frame_discontinuous = true;
if (frame->descriptor->kind == media::MediaKind::VIDEO) {
track->waiting_for_keyframe = true;
track->keyframe_requested = false;
requestKeyFrame(track);
}
std::lock_guard lock(mutex_);
stats_.frames_dropped_no_datagram += 1U + drained;
last_media_error_ = "peer/path has no usable QUIC DATAGRAM support";
return true;
}
const std::size_t effective_maximum =
std::min<std::size_t>(config_.maximum_datagram_bytes(), peer_maximum);
const std::size_t fragment_payload_capacity =
effective_maximum - kDatagramHeaderBytes;
const std::size_t required_fragments =
(frame->size() + fragment_payload_capacity - 1U) /
fragment_payload_capacity;
const std::size_t maximum_batch_packets =
transport_->maximumDatagramBatchPackets();
if (maximum_batch_packets == 0U ||
required_fragments > maximum_batch_packets) {
const auto drained = shedBufferedFrames(track);
track->next_frame_discontinuous = true;
if (frame->descriptor->kind == media::MediaKind::VIDEO) {
track->waiting_for_keyframe = true;
track->keyframe_requested = false;
requestKeyFrame(track);
}
std::lock_guard lock(mutex_);
stats_.frames_dropped_oversize += 1U + drained;
last_media_error_ =
"media frame exceeds the atomic DATAGRAM batch capacity";
return true;
}
auto& next_wire_sequence =
next_wire_frame_sequence_[track->config.track_id()];
if (next_wire_sequence == 0U) next_wire_sequence = 1U;
if (next_wire_sequence == std::numeric_limits<std::uint64_t>::max()) {
setError(error, "QUIC media wire frame sequence exhausted");
return false;
}
const std::uint64_t wire_frame_sequence = next_wire_sequence++;
std::vector<DatagramPacket> packets;
if (!packetizer_.packetize(
*frame, track->config.track_id(),
description.codec_generation_token, session_epoch_,
wire_frame_sequence, discontinuity, effective_maximum,
&packets, error)) {
const std::string packet_error =
error && !error->empty() ? *error : "media frame packetization failed";
track->next_frame_discontinuous = true;
if (frame->descriptor->kind == media::MediaKind::VIDEO) {
track->waiting_for_keyframe = true;
track->keyframe_requested = false;
requestKeyFrame(track);
}
std::lock_guard lock(mutex_);
++stats_.frames_dropped_oversize;
last_media_error_ = packet_error;
if (error) error->clear();
return true;
}
const std::size_t datagram_count = packets.size();
const auto result = transport_->sendDatagramBatch(std::move(packets), error);
if (result == TransportSendResult::QUEUED) {
track->next_frame_discontinuous = false;
if (frame->descriptor->kind == media::MediaKind::VIDEO && frame->key_frame) {
track->waiting_for_keyframe = false;
track->keyframe_requested = false;
}
std::lock_guard lock(mutex_);
++stats_.frames_queued;
stats_.datagrams_queued += datagram_count;
CMVR_LOG(INFO) << "[QuicEdgeService] queued media frame"
<< ", track_id=" << track->config.track_id()
<< ", source_track_id=" << track->source_track_id
<< ", frame_sequence=" << frame->sequence
<< ", wire_frame_sequence=" << wire_frame_sequence
<< ", payload_bytes=" << frame->size()
<< ", datagrams=" << datagram_count
<< ", key_frame=" << frame->key_frame
<< ", discontinuity=" << discontinuity
<< ", codec=" << codecName(description.codec)
<< ", payload_format="
<< payloadFormatName(description.payload_format)
<< ", capture_timestamp_us="
<< frame->capture_time_ns / 1000U
<< ", source_timestamp=" << frame->source_timestamp
<< ", source_frame_number=" << frame->source_frame_number;
return true;
}
if (result == TransportSendResult::WOULD_BLOCK) {
const auto drained = shedBufferedFrames(track);
track->next_frame_discontinuous = true;
if (frame->descriptor->kind == media::MediaKind::VIDEO) {
track->waiting_for_keyframe = true;
track->keyframe_requested = false;
requestKeyFrame(track);
}
std::lock_guard lock(mutex_);
stats_.frames_dropped_backpressure += 1U + drained;
return true;
}
if (error && error->empty()) {
*error = result == TransportSendResult::DISCONNECTED
? "QUIC transport disconnected while sending media"
: "QUIC transport failed while sending media";
}
if (result == TransportSendResult::DISCONNECTED ||
!transport_->isConnected()) {
return false;
}
recordMediaError(error ? *error : "QUIC DATAGRAM send failed");
if (error) error->clear();
track->next_frame_discontinuous = true;
return true;
}
std::uint64_t QuicEdgeService::shedBufferedFrames(ActiveTrack* track)
{
if (!track) return 0U;
std::uint64_t drained = 0U;
while (drained < kMaximumDrainPerPoll && track->subscription.tryRead()) {
++drained;
}
return drained;
}
void QuicEdgeService::requestKeyFrame(ActiveTrack* track)
{
if (!track || track->keyframe_requested || !media_hub_ ||
expectedKind(track->config) != media::MediaKind::VIDEO) {
return;
}
(void)media_hub_->requestKeyFrame(track->source_track_id);
track->keyframe_requested = true;
}
std::uint32_t QuicEdgeService::maximumFrameBytes(const ActiveTrack& track) const
{
return track.config.max_frame_bytes() == 0U
? config_.maximum_frame_bytes()
: track.config.max_frame_bytes();
}
void QuicEdgeService::initializeIdentity()
{
node_id_ = resolveNodeId(config_.node_id());
boot_id_ = resolveBootId();
software_version_ = config_.software_version().empty()
? "unknown" : config_.software_version();
effective_heartbeat_interval_ms_ = config_.heartbeat_interval_ms();
}
void QuicEdgeService::resetConnectionStatus()
{
std::lock_guard lock(mutex_);
registered_ = false;
session_id_.clear();
observed_source_ip_.clear();
outstanding_heartbeat_sequence_ = 0U;
active_media_tracks_ = 0U;
last_media_error_.clear();
media_stop_requested_ = false;
media_connection_failed_ = false;
media_connection_error_.clear();
effective_heartbeat_interval_ms_ = config_.heartbeat_interval_ms();
next_wire_frame_sequence_.clear();
media_session_announced_ = false;
heartbeat_deadline_ = {};
next_heartbeat_ = {};
}
void QuicEdgeService::recordMediaError(const std::string& error)
{
std::lock_guard lock(mutex_);
++stats_.source_errors;
last_media_error_ = error.empty() ? "unknown media source error" : error;
}
void QuicEdgeService::setState(const QuicEdgeServiceState state,
const std::string& error)
{
std::lock_guard lock(mutex_);
state_ = state;
if (!error.empty()) {
last_error_ = error;
} else if (state == QuicEdgeServiceState::ONLINE ||
state == QuicEdgeServiceState::STOPPED ||
state == QuicEdgeServiceState::DISABLED) {
last_error_.clear();
}
}
bool QuicEdgeService::waitForStop(const std::chrono::milliseconds duration)
{
std::unique_lock lock(mutex_);
return stop_cv_.wait_for(lock, duration, [this]() { return stop_requested_; });
}
std::chrono::milliseconds QuicEdgeService::nextBackoff(
const std::chrono::milliseconds current)
{
const double multiplied = static_cast<double>(current.count()) *
config_.reconnect().multiplier();
const auto bounded = std::min<double>(
multiplied, static_cast<double>(config_.reconnect().maximum_delay_ms()));
return std::chrono::milliseconds(static_cast<std::int64_t>(bounded));
}
std::chrono::milliseconds QuicEdgeService::jittered(
const std::chrono::milliseconds base)
{
const auto jitter_percent = config_.reconnect().jitter_percent();
if (jitter_percent == 0U) return base;
thread_local std::mt19937 generator(std::random_device{}());
const double fraction = static_cast<double>(jitter_percent) / 100.0;
std::uniform_real_distribution<double> distribution(1.0 - fraction,
1.0 + fraction);
const double value = static_cast<double>(base.count()) * distribution(generator);
return std::chrono::milliseconds(
std::max<std::int64_t>(1, static_cast<std::int64_t>(value)));
}
std::uint64_t QuicEdgeService::nextSessionEpoch()
{
if (session_epoch_ == 0U) {
const auto wall_clock = static_cast<std::uint64_t>(
std::chrono::duration_cast<std::chrono::microseconds>(
std::chrono::system_clock::now().time_since_epoch()).count());
std::random_device random;
session_epoch_ = wall_clock ^
(static_cast<std::uint64_t>(random()) << 32U) ^ random();
if (session_epoch_ == 0U) session_epoch_ = 1U;
return session_epoch_;
}
session_epoch_ = session_epoch_ == std::numeric_limits<std::uint64_t>::max()
? 1U : session_epoch_ + 1U;
return session_epoch_;
}
} // namespace cmvr::quic_edge