1879 lines
72 KiB
C++
1879 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_robot_id(config.robot_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::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.robot_id().empty()) {
|
|
setError(error, "QUIC edge robot_id is empty");
|
|
return false;
|
|
}
|
|
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 (!transport_ || !media_hub_) {
|
|
const std::string message = "QUIC edge transport or MediaSourceHub is null";
|
|
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::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::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()
|
|
<< ", robot_id=" << node.robot_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_robot_id(config_.robot_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) {
|
|
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
|