cmvr-es/cmvr-es/service/quic_edge/tests/quic_edge_protocol_test.cpp

1211 lines
48 KiB
C++

#include <atomic>
#include <chrono>
#include <condition_variable>
#include <cstdint>
#include <deque>
#include <functional>
#include <future>
#include <iostream>
#include <mutex>
#include <optional>
#include <string>
#include <thread>
#include <utility>
#include <vector>
#include "cmvr/quic_edge/v1/quic_edge.pb.h"
#include "manager/media_source_hub/include/media_source_hub.h"
#include "service/quic_edge/include/control_framing.h"
#include "service/quic_edge/include/datagram_packetizer.h"
#include "service/quic_edge/include/quic_edge_service.h"
namespace {
#define CHECK_TRUE(expression) \
do { \
if (!(expression)) { \
std::cerr << "CHECK failed at line " << __LINE__ << ": " \
<< #expression << '\n'; \
return false; \
} \
} while (false)
using namespace cmvr;
class FakeTransport final : public quic_edge::QuicTransport {
public:
explicit FakeTransport(const bool accept_registration = true,
const bool acknowledge_heartbeats = true,
const bool valid_heartbeat_session = true,
const std::size_t media_session_would_block_count = 0U,
const std::uint32_t registration_heartbeat_interval_ms =
250U,
const std::uint32_t heartbeat_ack_delay_ms = 0U)
: accept_registration_(accept_registration),
acknowledge_heartbeats_(acknowledge_heartbeats),
valid_heartbeat_session_(valid_heartbeat_session),
media_session_would_block_count_(media_session_would_block_count),
registration_heartbeat_interval_ms_(
registration_heartbeat_interval_ms),
heartbeat_ack_delay_ms_(heartbeat_ack_delay_ms)
{
}
bool connect(const config::QuicEdgeConfig&,
std::chrono::milliseconds,
std::string*) override
{
connected_.store(true);
++connect_count_;
{
std::lock_guard lock(mutex_);
connect_times_.push_back(std::chrono::steady_clock::now());
}
condition_.notify_all();
return true;
}
void disconnect() override
{
connected_.store(false);
{
std::lock_guard lock(mutex_);
delayed_heartbeat_ack_.reset();
}
condition_.notify_all();
}
bool isConnected() const override { return connected_.load(); }
std::size_t maximumDatagramBytes() const override { return 1200U; }
std::size_t maximumDatagramBatchPackets() const override { return 32U; }
quic_edge::TransportSendResult sendControl(
std::vector<std::uint8_t> message, std::string*) override
{
if (!connected_.load()) return quic_edge::TransportSendResult::DISCONNECTED;
std::vector<std::vector<std::uint8_t>> frames;
std::string decode_error;
quic_edge::ControlFrameDecoder decoder(1024U * 1024U);
if (!decoder.push(message, &frames, &decode_error) || frames.size() != 1U) {
return quic_edge::TransportSendResult::ERROR;
}
cmvr::quic_edge::v1::EdgeControlEnvelope envelope;
if (!envelope.ParseFromArray(frames.front().data(),
static_cast<int>(frames.front().size()))) {
return quic_edge::TransportSendResult::ERROR;
}
{
std::lock_guard lock(mutex_);
if (envelope.has_media_session_open() &&
media_session_would_block_count_ != 0U) {
--media_session_would_block_count_;
return quic_edge::TransportSendResult::WOULD_BLOCK;
}
edge_message_sequences_.push_back(envelope.message_sequence());
controls_.push_back(std::move(message));
if (envelope.has_node_register_request()) {
const auto& request = envelope.node_register_request();
last_registered_node_id_ = request.node().node_id();
last_registered_robot_id_ = request.node().robot_id();
last_grpc_endpoint_port_ = request.node().grpc_endpoint().port();
last_interface_count_ = request.node().local_interfaces_size();
cmvr::quic_edge::v1::EdgeControlEnvelope response;
response.set_protocol_version(quic_edge::kProtocolVersion);
response.set_message_sequence(server_message_sequence_++);
auto* registration = response.mutable_node_register_response();
registration->set_accepted(accept_registration_);
registration->set_session_id(
accept_registration_ ? "test-session" : "");
registration->set_message(
accept_registration_ ? "accepted" : "rejected for test");
registration->set_heartbeat_interval_ms(
registration_heartbeat_interval_ms_);
registration->set_observed_source_ip("203.0.113.10");
enqueueEnvelopeLocked(response);
} else if (envelope.has_node_heartbeat()) {
const auto& heartbeat = envelope.node_heartbeat();
last_heartbeat_ = heartbeat;
has_last_heartbeat_ = true;
heartbeat_times_.push_back(
std::chrono::steady_clock::now());
if (!acknowledge_heartbeats_) {
condition_.notify_all();
return quic_edge::TransportSendResult::QUEUED;
}
cmvr::quic_edge::v1::EdgeControlEnvelope response;
response.set_protocol_version(quic_edge::kProtocolVersion);
response.set_message_sequence(server_message_sequence_++);
auto* ack = response.mutable_node_heartbeat_ack();
ack->set_accepted(true);
ack->set_acknowledged_sequence(heartbeat.sequence());
ack->set_session_id(valid_heartbeat_session_
? heartbeat.session_id() : "");
ack->set_observed_source_ip("203.0.113.11");
if (heartbeat_ack_delay_ms_ == 0U) {
enqueueEnvelopeLocked(response);
} else {
delayed_heartbeat_ack_ = std::move(response);
delayed_heartbeat_ack_ready_at_ =
std::chrono::steady_clock::now() +
std::chrono::milliseconds(
heartbeat_ack_delay_ms_);
}
}
}
condition_.notify_all();
return quic_edge::TransportSendResult::QUEUED;
}
quic_edge::TransportReceiveResult receiveControl(
std::vector<std::uint8_t>* chunk,
const std::chrono::milliseconds timeout,
std::string*) override
{
if (!chunk) return quic_edge::TransportReceiveResult::ERROR;
std::unique_lock lock(mutex_);
releaseDelayedHeartbeatAckLocked();
if (control_receive_queue_.empty() && connected_.load()) {
auto wait_duration = timeout;
if (delayed_heartbeat_ack_) {
const auto now = std::chrono::steady_clock::now();
if (now < delayed_heartbeat_ack_ready_at_) {
wait_duration = std::min(
wait_duration,
std::chrono::duration_cast<std::chrono::milliseconds>(
delayed_heartbeat_ack_ready_at_ - now) +
std::chrono::milliseconds(1));
}
}
condition_.wait_for(lock, wait_duration, [this]() {
return !control_receive_queue_.empty() ||
!connected_.load();
});
releaseDelayedHeartbeatAckLocked();
}
if (!control_receive_queue_.empty()) {
*chunk = std::move(control_receive_queue_.front());
control_receive_queue_.pop_front();
return quic_edge::TransportReceiveResult::DATA;
}
return connected_.load()
? quic_edge::TransportReceiveResult::TIMEOUT
: quic_edge::TransportReceiveResult::DISCONNECTED;
}
quic_edge::TransportSendResult sendDatagramBatch(
std::vector<quic_edge::DatagramPacket> packets, std::string*) override
{
if (!connected_.load()) return quic_edge::TransportSendResult::DISCONNECTED;
std::lock_guard lock(mutex_);
datagrams_.push_back(std::move(packets));
return quic_edge::TransportSendResult::QUEUED;
}
std::size_t datagramBatchCount() const
{
std::lock_guard lock(mutex_);
return datagrams_.size();
}
std::size_t controlCount() const
{
std::lock_guard lock(mutex_);
return controls_.size();
}
bool controlSequencesStrictlyIncreasing() const
{
std::lock_guard lock(mutex_);
for (std::size_t index = 1U;
index < edge_message_sequences_.size(); ++index) {
if (edge_message_sequences_[index] <=
edge_message_sequences_[index - 1U]) {
return false;
}
}
return !edge_message_sequences_.empty();
}
std::vector<std::uint64_t> datagramFrameSequences() const
{
std::lock_guard lock(mutex_);
std::vector<std::uint64_t> sequences;
for (const auto& batch : datagrams_) {
if (!batch.empty()) sequences.push_back(batch.front().header.frame_sequence);
}
return sequences;
}
std::uint64_t connectCount() const { return connect_count_.load(); }
std::vector<std::chrono::milliseconds> connectIntervals() const
{
std::lock_guard lock(mutex_);
std::vector<std::chrono::milliseconds> intervals;
for (std::size_t index = 1U; index < connect_times_.size(); ++index) {
intervals.push_back(std::chrono::duration_cast<std::chrono::milliseconds>(
connect_times_[index] - connect_times_[index - 1U]));
}
return intervals;
}
std::string lastRegisteredNodeId() const
{
std::lock_guard lock(mutex_);
return last_registered_node_id_;
}
std::string lastRegisteredRobotId() const
{
std::lock_guard lock(mutex_);
return last_registered_robot_id_;
}
std::uint32_t lastGrpcEndpointPort() const
{
std::lock_guard lock(mutex_);
return last_grpc_endpoint_port_;
}
int lastInterfaceCount() const
{
std::lock_guard lock(mutex_);
return last_interface_count_;
}
std::size_t heartbeatCount() const
{
std::lock_guard lock(mutex_);
return heartbeat_times_.size();
}
std::vector<std::chrono::milliseconds> heartbeatIntervals() const
{
std::lock_guard lock(mutex_);
std::vector<std::chrono::milliseconds> intervals;
for (std::size_t index = 1U;
index < heartbeat_times_.size(); ++index) {
intervals.push_back(
std::chrono::duration_cast<std::chrono::milliseconds>(
heartbeat_times_[index] -
heartbeat_times_[index - 1U]));
}
return intervals;
}
cmvr::quic_edge::v1::NodeHeartbeat lastHeartbeat() const
{
std::lock_guard lock(mutex_);
return has_last_heartbeat_
? last_heartbeat_
: cmvr::quic_edge::v1::NodeHeartbeat{};
}
private:
void releaseDelayedHeartbeatAckLocked()
{
if (!delayed_heartbeat_ack_ ||
std::chrono::steady_clock::now() <
delayed_heartbeat_ack_ready_at_) {
return;
}
enqueueEnvelopeLocked(*delayed_heartbeat_ack_);
delayed_heartbeat_ack_.reset();
}
void enqueueEnvelopeLocked(
const cmvr::quic_edge::v1::EdgeControlEnvelope& envelope)
{
std::string serialized;
if (!envelope.SerializeToString(&serialized)) return;
std::vector<std::uint8_t> framed;
std::string error;
if (quic_edge::ControlFrameEncoder::encode(
reinterpret_cast<const std::uint8_t*>(serialized.data()),
serialized.size(), 1024U * 1024U, &framed, &error)) {
const auto split = framed.size() / 2U;
control_receive_queue_.emplace_back(
framed.begin(), framed.begin() + static_cast<std::ptrdiff_t>(split));
control_receive_queue_.emplace_back(
framed.begin() + static_cast<std::ptrdiff_t>(split), framed.end());
}
}
mutable std::mutex mutex_;
std::condition_variable condition_;
std::atomic<bool> connected_{false};
std::atomic<std::uint64_t> connect_count_{0};
bool accept_registration_{true};
bool acknowledge_heartbeats_{true};
bool valid_heartbeat_session_{true};
std::size_t media_session_would_block_count_{0U};
std::uint32_t registration_heartbeat_interval_ms_{250U};
std::uint32_t heartbeat_ack_delay_ms_{0U};
std::uint64_t server_message_sequence_{0};
std::string last_registered_node_id_;
std::string last_registered_robot_id_;
std::uint32_t last_grpc_endpoint_port_{0};
int last_interface_count_{0};
std::vector<std::vector<std::uint8_t>> controls_;
std::deque<std::vector<std::uint8_t>> control_receive_queue_;
std::vector<std::vector<quic_edge::DatagramPacket>> datagrams_;
std::vector<std::uint64_t> edge_message_sequences_;
std::vector<std::chrono::steady_clock::time_point> connect_times_;
std::vector<std::chrono::steady_clock::time_point> heartbeat_times_;
bool has_last_heartbeat_{false};
cmvr::quic_edge::v1::NodeHeartbeat last_heartbeat_;
std::optional<cmvr::quic_edge::v1::EdgeControlEnvelope>
delayed_heartbeat_ack_;
std::chrono::steady_clock::time_point
delayed_heartbeat_ack_ready_at_{};
};
config::QuicEdgeConfig validConfig(const std::string& source_track_id)
{
config::QuicEdgeConfig config;
config.set_id("quic-test");
config.set_server_host("127.0.0.1");
config.set_server_port(4433);
config.set_alpn("cmvr-quic-edge/1");
config.set_node_id("test-node");
config.set_robot_id("CN-CMVR-MBLRV1-CHAGAN-20260731-001");
config.set_software_version("test-version");
config.set_grpc_endpoint_host("auto");
config.set_grpc_endpoint_port(50052U);
config.set_grpc_endpoint_tls(false);
config.set_heartbeat_interval_ms(250U);
config.set_control_response_timeout_ms(100U);
config.set_include_loopback_interfaces(true);
config.mutable_tls()->set_allow_insecure(true);
config.mutable_reconnect()->set_initial_delay_ms(5);
config.mutable_reconnect()->set_maximum_delay_ms(20);
config.mutable_reconnect()->set_multiplier(2.0);
config.mutable_reconnect()->set_connect_timeout_ms(50);
config.set_maximum_datagram_bytes(1200);
config.set_maximum_control_frame_bytes(4096);
config.set_maximum_frame_bytes(32 * 1024);
config.set_datagram_send_queue_depth(32);
config.set_media_poll_interval_ms(1);
auto* track = config.add_tracks();
track->set_track_id(7);
track->set_source_kind(config::QuicEdgeTrackConfig::SOURCE_KIND_CAMERA);
track->set_device_id("camera-test");
track->set_source_track_id(source_track_id);
track->set_enable(true);
return config;
}
config::QuicEdgeConfig validPresenceOnlyConfig()
{
auto config = validConfig("unused/video/color");
config.clear_tracks();
return config;
}
media::TrackDescriptorPtr videoDescriptor(const std::string& id)
{
media::TrackDescriptor::Config config;
config.id = id;
config.source_id = "camera-test";
config.kind = media::MediaKind::VIDEO;
config.codec = media::Codec::H264;
config.payload_format = media::PayloadFormat::ANNEX_B;
config.time_base = {1, 90000};
config.width = 640;
config.height = 480;
config.nominal_rate = 30;
config.generation = 0x0000000200000003ULL;
return media::makeTrackDescriptor(std::move(config));
}
bool waitUntil(const std::function<bool()>& predicate)
{
const auto deadline = std::chrono::steady_clock::now() +
std::chrono::seconds(2);
while (std::chrono::steady_clock::now() < deadline) {
if (predicate()) return true;
std::this_thread::sleep_for(std::chrono::milliseconds(2));
}
return predicate();
}
bool testControlFraming()
{
const std::vector<std::uint8_t> payload{1, 2, 3, 4, 5};
std::vector<std::uint8_t> encoded;
std::string error;
CHECK_TRUE(quic_edge::ControlFrameEncoder::encode(
payload, 64, &encoded, &error));
quic_edge::ControlFrameDecoder decoder(64);
std::vector<std::vector<std::uint8_t>> decoded;
CHECK_TRUE(decoder.push(encoded.data(), 2, &decoded, &error));
CHECK_TRUE(decoded.empty());
CHECK_TRUE(decoder.push(encoded.data() + 2, encoded.size() - 2,
&decoded, &error));
CHECK_TRUE(decoded.size() == 1 && decoded.front() == payload);
return true;
}
bool testPacketizer()
{
auto descriptor = videoDescriptor("camera-test/video/color");
media::MediaFrame::Config frame_config;
frame_config.descriptor = descriptor;
frame_config.payload.resize(2500, 0x5a);
frame_config.sequence = 11;
frame_config.capture_time_ns = 1234567000ULL;
frame_config.key_frame = true;
auto frame = media::makeMediaFrame(std::move(frame_config));
quic_edge::DatagramPacketizer packetizer(100);
std::vector<quic_edge::DatagramPacket> packets;
std::string error;
const std::uint32_t token =
quic_edge::descriptorGenerationToken(descriptor->generation);
CHECK_TRUE(packetizer.packetize(*frame, 7, token, 99, 11, true, 1200,
&packets, &error));
CHECK_TRUE(packets.size() == 3);
std::size_t payload_bytes = 0;
for (std::size_t index = 0; index < packets.size(); ++index) {
quic_edge::DatagramHeader header;
CHECK_TRUE(quic_edge::DatagramPacketizer::decodeHeader(
packets[index].bytes, &header, &error));
CHECK_TRUE(header.track_id == 7 && header.session_epoch == 99);
CHECK_TRUE(header.fragment_index == index && header.fragment_count == 3);
CHECK_TRUE((header.flags & quic_edge::DATAGRAM_FLAG_KEY_FRAME) != 0);
CHECK_TRUE((header.flags & quic_edge::DATAGRAM_FLAG_DISCONTINUITY) != 0);
payload_bytes += header.payload_size;
}
CHECK_TRUE(payload_bytes == frame->size());
return true;
}
bool testServiceWithSharedHub()
{
const std::string track_id = "camera-test/video/color";
media::MediaSourceManager hub;
media::MediaSourceManager::FrameSink sink;
std::mutex sink_mutex;
std::atomic<bool> source_started{false};
std::atomic<std::uint64_t> keyframe_requests{0};
media::MediaSourceManager::SourceCallbacks callbacks;
callbacks.start = [&](const media::MediaSourceManager::FrameSink& value,
const media::MediaSourceManager::CancelPredicate&) {
std::lock_guard lock(sink_mutex);
sink = value;
source_started.store(true);
return true;
};
callbacks.stop = [&]() { source_started.store(false); };
callbacks.request_key_frame = [&]() {
++keyframe_requests;
return true;
};
auto descriptor = videoDescriptor(track_id);
CHECK_TRUE(hub.registerSource(descriptor, std::move(callbacks), 8));
auto transport = std::make_unique<FakeTransport>(true, true, true, 1U);
FakeTransport* transport_view = transport.get();
quic_edge::QuicEdgeService service(
validConfig(track_id), std::move(transport), hub);
std::string error;
CHECK_TRUE(service.initialize(&error));
CHECK_TRUE(service.start(&error));
CHECK_TRUE(waitUntil([&]() { return source_started.load(); }));
CHECK_TRUE(waitUntil([&]() {
return service.stats().registrations_accepted == 1U &&
service.stats().heartbeats_acknowledged >= 1U;
}));
CHECK_TRUE(transport_view->lastRegisteredNodeId() == "test-node");
CHECK_TRUE(transport_view->lastRegisteredRobotId() ==
"CN-CMVR-MBLRV1-CHAGAN-20260731-001");
CHECK_TRUE(transport_view->lastGrpcEndpointPort() == 50052U);
CHECK_TRUE(transport_view->connectCount() == 1U);
media::MediaFrame::Config frame_config;
frame_config.descriptor = descriptor;
frame_config.payload.resize(1800, 0x11);
frame_config.sequence = 1;
frame_config.capture_time_ns = 1000000;
frame_config.key_frame = true;
media::MediaSourceManager::FrameSink publisher;
{
std::lock_guard lock(sink_mutex);
publisher = sink;
}
CHECK_TRUE(static_cast<bool>(publisher));
publisher(media::makeMediaFrame(std::move(frame_config)));
CHECK_TRUE(waitUntil([&]() { return transport_view->datagramBatchCount() == 1; }));
CHECK_TRUE(service.stats().frames_queued == 1);
CHECK_TRUE(keyframe_requests.load() != 0);
// Codec initialization bytes may be learned at the first encoder IDR
// without changing the source generation. The reliable descriptor must be
// refreshed rather than treating this legal enrichment as a source error.
auto enriched_descriptor_config = media::TrackDescriptor::Config{};
enriched_descriptor_config.id = track_id;
enriched_descriptor_config.source_id = "camera-test";
enriched_descriptor_config.kind = media::MediaKind::VIDEO;
enriched_descriptor_config.codec = media::Codec::H264;
enriched_descriptor_config.payload_format = media::PayloadFormat::ANNEX_B;
enriched_descriptor_config.time_base = {1, 90000};
enriched_descriptor_config.width = 640;
enriched_descriptor_config.height = 480;
enriched_descriptor_config.nominal_rate = 30;
enriched_descriptor_config.generation = descriptor->generation;
enriched_descriptor_config.codec_config = {0, 0, 0, 1, 0x67};
auto enriched_descriptor =
media::makeTrackDescriptor(std::move(enriched_descriptor_config));
media::MediaFrame::Config enriched_frame;
enriched_frame.descriptor = std::move(enriched_descriptor);
enriched_frame.payload.resize(256, 0x22);
// Simulate a source restart that resets its device-local sequence. The
// QUIC wire sequence must remain monotonic inside the media epoch.
enriched_frame.sequence = 1;
enriched_frame.capture_time_ns = 2000000;
enriched_frame.key_frame = true;
enriched_frame.discontinuity = true;
publisher(media::makeMediaFrame(std::move(enriched_frame)));
CHECK_TRUE(waitUntil([&]() { return transport_view->datagramBatchCount() == 2; }));
const auto wire_sequences = transport_view->datagramFrameSequences();
CHECK_TRUE(wire_sequences.size() == 2U &&
wire_sequences[0] == 1U && wire_sequences[1] == 2U);
CHECK_TRUE(transport_view->controlCount() >= 3);
CHECK_TRUE(transport_view->controlSequencesStrictlyIncreasing());
CHECK_TRUE(service.state() != quic_edge::QuicEdgeServiceState::FAILED);
service.stop();
return true;
}
bool testMediaActivityInterruptPreservesPresenceAndResumes()
{
const std::string track_id = "camera-stop-all/video/color";
media::MediaSourceHub hub;
media::MediaSourceManager::FrameSink sink;
std::mutex sink_mutex;
std::atomic<bool> source_started{false};
std::atomic<std::uint64_t> source_starts{0U};
std::atomic<std::uint64_t> source_stops{0U};
media::MediaSourceManager::SourceCallbacks callbacks;
callbacks.start = [&](const media::MediaSourceManager::FrameSink& value,
const media::MediaSourceManager::CancelPredicate&) {
{
std::lock_guard lock(sink_mutex);
sink = value;
}
source_started.store(true);
++source_starts;
return true;
};
callbacks.stop = [&]() {
source_started.store(false);
++source_stops;
};
callbacks.request_key_frame = [] { return true; };
const auto descriptor = videoDescriptor(track_id);
CHECK_TRUE(hub.registerSource(descriptor, std::move(callbacks), 8U));
auto transport = std::make_unique<FakeTransport>();
FakeTransport* transport_view = transport.get();
quic_edge::QuicEdgeService edge_service(
validConfig(track_id), std::move(transport), hub);
std::string error;
CHECK_TRUE(edge_service.initialize(&error));
CHECK_TRUE(edge_service.start(&error));
CHECK_TRUE(waitUntil([&] {
return edge_service.state() == quic_edge::QuicEdgeServiceState::ONLINE &&
edge_service.status().registered && source_started.load() &&
hub.subscriberCount(track_id) == 1U;
}));
auto publish = [&](const std::uint64_t sequence) {
media::MediaSourceManager::FrameSink publisher;
{
std::lock_guard lock(sink_mutex);
publisher = sink;
}
if (!publisher) return false;
media::MediaFrame::Config frame;
frame.descriptor = descriptor;
frame.payload.resize(256U, 0x5aU);
frame.sequence = sequence;
frame.capture_time_ns = sequence * 1000000U;
frame.key_frame = true;
publisher(media::makeMediaFrame(std::move(frame)));
return true;
};
CHECK_TRUE(publish(1U));
CHECK_TRUE(waitUntil([&] {
return transport_view->datagramBatchCount() == 1U;
}));
const auto heartbeat_before_stop = transport_view->heartbeatCount();
CHECK_TRUE(edge_service.interruptMediaActivities());
CHECK_TRUE(hub.subscriberCount(track_id) == 0U);
CHECK_TRUE(!source_started.load());
CHECK_TRUE(source_stops.load() == 1U);
CHECK_TRUE(edge_service.state() == quic_edge::QuicEdgeServiceState::ONLINE);
CHECK_TRUE(edge_service.status().registered);
CHECK_TRUE(transport_view->connectCount() == 1U);
// A publisher retained by the stopped generation must no longer enqueue
// frames while StopAll admission is closed.
CHECK_TRUE(publish(2U));
std::this_thread::sleep_for(std::chrono::milliseconds(30));
CHECK_TRUE(transport_view->datagramBatchCount() == 1U);
CHECK_TRUE(waitUntil([&] {
return source_starts.load() == 2U && source_started.load() &&
hub.subscriberCount(track_id) == 1U;
}));
CHECK_TRUE(publish(3U));
CHECK_TRUE(waitUntil([&] {
return transport_view->datagramBatchCount() == 2U;
}));
CHECK_TRUE(waitUntil([&] {
return transport_view->heartbeatCount() > heartbeat_before_stop;
}));
CHECK_TRUE(edge_service.stats().registrations_accepted == 1U);
CHECK_TRUE(transport_view->connectCount() == 1U);
CHECK_TRUE(edge_service.state() == quic_edge::QuicEdgeServiceState::ONLINE);
edge_service.stop();
return true;
}
bool testMediaActivityInterruptCancelsStartingSubscription()
{
const std::string track_id = "slow-stop-all/video/color";
media::MediaSourceHub hub;
std::atomic<bool> start_entered{false};
std::atomic<bool> start_cancelled{false};
media::MediaSourceManager::SourceCallbacks callbacks;
callbacks.start = [&](const media::MediaSourceManager::FrameSink&,
const media::MediaSourceManager::CancelPredicate& cancelled) {
start_entered.store(true);
while (!cancelled()) {
std::this_thread::sleep_for(std::chrono::milliseconds(2));
}
start_cancelled.store(true);
return false;
};
callbacks.stop = [] {};
CHECK_TRUE(hub.registerSource(
videoDescriptor(track_id), std::move(callbacks), 8U));
auto transport = std::make_unique<FakeTransport>();
FakeTransport* transport_view = transport.get();
quic_edge::QuicEdgeService edge_service(
validConfig(track_id), std::move(transport), hub);
std::string error;
CHECK_TRUE(edge_service.initialize(&error));
CHECK_TRUE(edge_service.start(&error));
CHECK_TRUE(waitUntil([&] {
return start_entered.load() && edge_service.status().registered;
}));
auto interrupt = std::async(std::launch::async, [&edge_service] {
return edge_service.interruptMediaActivities();
});
CHECK_TRUE(interrupt.wait_for(std::chrono::milliseconds(500)) ==
std::future_status::ready);
CHECK_TRUE(interrupt.get());
CHECK_TRUE(start_cancelled.load());
CHECK_TRUE(hub.subscriberCount(track_id) == 0U);
CHECK_TRUE(edge_service.state() == quic_edge::QuicEdgeServiceState::ONLINE);
CHECK_TRUE(edge_service.status().registered);
CHECK_TRUE(transport_view->connectCount() == 1U);
edge_service.stop();
return true;
}
bool testMissingInjectedSourceRetriesSafely()
{
media::MediaSourceManager hub;
auto transport = std::make_unique<FakeTransport>();
FakeTransport* transport_view = transport.get();
quic_edge::QuicEdgeService service(
validConfig("missing/video/color"), std::move(transport), hub);
std::string error;
CHECK_TRUE(service.initialize(&error));
CHECK_TRUE(service.start(&error));
CHECK_TRUE(waitUntil([&]() { return service.stats().source_errors >= 2; }));
CHECK_TRUE(service.state() == quic_edge::QuicEdgeServiceState::ONLINE);
CHECK_TRUE(service.status().registered);
CHECK_TRUE(service.status().active_media_tracks == 0U);
CHECK_TRUE(transport_view->connectCount() == 1U);
service.stop();
return true;
}
bool testPresenceOnlyWithoutMedia()
{
media::MediaSourceManager hub;
auto transport = std::make_unique<FakeTransport>();
FakeTransport* transport_view = transport.get();
quic_edge::QuicEdgeService service(
validPresenceOnlyConfig(), std::move(transport), hub);
std::string error;
CHECK_TRUE(service.initialize(&error));
CHECK_TRUE(service.start(&error));
CHECK_TRUE(waitUntil([&]() {
return service.stats().registrations_accepted == 1U &&
service.stats().heartbeats_acknowledged >= 1U;
}));
const auto status = service.status();
CHECK_TRUE(status.registered);
CHECK_TRUE(status.session_id == "test-session");
CHECK_TRUE(status.observed_source_ip == "203.0.113.11");
CHECK_TRUE(status.active_media_tracks == 0U);
CHECK_TRUE(transport_view->datagramBatchCount() == 0U);
CHECK_TRUE(transport_view->connectCount() == 1U);
CHECK_TRUE(service.state() == quic_edge::QuicEdgeServiceState::ONLINE);
service.stop();
return true;
}
bool testDeviceManagerSnapshotInHeartbeat()
{
device::DeviceManagerSnapshot snapshot;
snapshot.name = "edge-device-manager";
snapshot.version = "2.3.4";
snapshot.description = "heartbeat snapshot test";
device::ManagedDeviceSnapshot disabled;
disabled.id = "camera-disabled";
disabled.kind = device::DeviceKind::Camera;
disabled.type_name = "DEVICE_TYPE_CAMERA";
disabled.enabled = false;
disabled.state = device::ManagedDeviceState::Disabled;
disabled.health.state = device::DeviceHealthState::Unknown;
disabled.status_updated_at_unix_ms = 101U;
snapshot.devices.push_back(disabled);
device::ManagedDeviceSnapshot running;
running.id = "src1100";
running.kind = device::DeviceKind::AGV;
running.type_name = "SeerRobokitAgv";
running.enabled = true;
running.state = device::ManagedDeviceState::Running;
running.health.state = device::DeviceHealthState::Healthy;
running.status_updated_at_unix_ms = 202U;
snapshot.devices.push_back(running);
device::ManagedDeviceSnapshot failed;
failed.id = "microphone-failed";
failed.kind = device::DeviceKind::Microphone;
failed.type_name = "FfmpegMicrophone";
failed.enabled = true;
failed.state = device::ManagedDeviceState::Error;
failed.health.state = device::DeviceHealthState::Fault;
failed.abnormal = true;
failed.error_message = "device start returned false";
failed.status_updated_at_unix_ms = 303U;
snapshot.devices.push_back(failed);
media::MediaSourceManager hub;
auto transport = std::make_unique<FakeTransport>();
FakeTransport* transport_view = transport.get();
quic_edge::QuicEdgeService service(
validPresenceOnlyConfig(), std::move(transport), hub,
[snapshot]() { return snapshot; });
std::string error;
CHECK_TRUE(service.initialize(&error));
CHECK_TRUE(service.start(&error));
CHECK_TRUE(waitUntil([&]() {
return service.stats().heartbeats_acknowledged >= 1U;
}));
const auto heartbeat = transport_view->lastHeartbeat();
service.stop();
CHECK_TRUE(heartbeat.robot_id() ==
"CN-CMVR-MBLRV1-CHAGAN-20260731-001");
CHECK_TRUE(heartbeat.has_device_manager());
CHECK_TRUE(heartbeat.device_manager().manager_name() ==
"edge-device-manager");
CHECK_TRUE(heartbeat.device_manager().manager_version() == "2.3.4");
CHECK_TRUE(heartbeat.device_manager().manager_description() ==
"heartbeat snapshot test");
CHECK_TRUE(heartbeat.device_manager().sampled_at_unix_ms() ==
heartbeat.sent_at_unix_ms());
CHECK_TRUE(heartbeat.device_manager().devices_size() == 2);
const auto& wire_failed = heartbeat.device_manager().devices(0);
CHECK_TRUE(wire_failed.device_id() == "microphone-failed");
CHECK_TRUE(wire_failed.enabled());
CHECK_TRUE(wire_failed.kind() ==
cmvr::quic_edge::v1::DEVICE_KIND_MICROPHONE);
CHECK_TRUE(wire_failed.manager_state() ==
cmvr::quic_edge::v1::MANAGED_DEVICE_STATE_ERROR);
CHECK_TRUE(wire_failed.health() ==
cmvr::quic_edge::v1::DEVICE_HEALTH_STATUS_FAULT);
CHECK_TRUE(wire_failed.has_error());
CHECK_TRUE(wire_failed.error_message() ==
"device start returned false");
const auto& wire_running = heartbeat.device_manager().devices(1);
CHECK_TRUE(wire_running.device_id() == "src1100");
CHECK_TRUE(wire_running.enabled());
CHECK_TRUE(wire_running.kind() ==
cmvr::quic_edge::v1::DEVICE_KIND_AGV);
CHECK_TRUE(wire_running.manager_state() ==
cmvr::quic_edge::v1::MANAGED_DEVICE_STATE_RUNNING);
CHECK_TRUE(wire_running.health() ==
cmvr::quic_edge::v1::DEVICE_HEALTH_STATUS_HEALTHY);
CHECK_TRUE(!wire_running.has_error());
for (const auto& wire_device : heartbeat.device_manager().devices()) {
CHECK_TRUE(wire_device.enabled());
CHECK_TRUE(wire_device.device_id() != "camera-disabled");
}
return true;
}
bool testAllDeviceKindAndStateMappings()
{
using ProtoKind = cmvr::quic_edge::v1::DeviceKind;
using ProtoState = cmvr::quic_edge::v1::ManagedDeviceState;
using ProtoHealth = cmvr::quic_edge::v1::DeviceHealthStatus;
const std::vector<std::pair<device::DeviceKind, ProtoKind>> kinds{
{device::DeviceKind::Unknown,
cmvr::quic_edge::v1::DEVICE_KIND_UNSPECIFIED},
{device::DeviceKind::AGV, cmvr::quic_edge::v1::DEVICE_KIND_AGV},
{device::DeviceKind::Arm, cmvr::quic_edge::v1::DEVICE_KIND_ARM},
{device::DeviceKind::Battery,
cmvr::quic_edge::v1::DEVICE_KIND_BATTERY},
{device::DeviceKind::BioHead,
cmvr::quic_edge::v1::DEVICE_KIND_BIO_HEAD},
{device::DeviceKind::Camera,
cmvr::quic_edge::v1::DEVICE_KIND_CAMERA},
{device::DeviceKind::CanBus,
cmvr::quic_edge::v1::DEVICE_KIND_CAN_BUS},
{device::DeviceKind::DexHand,
cmvr::quic_edge::v1::DEVICE_KIND_DEX_HAND},
{device::DeviceKind::Gripper,
cmvr::quic_edge::v1::DEVICE_KIND_GRIPPER},
{device::DeviceKind::Microphone,
cmvr::quic_edge::v1::DEVICE_KIND_MICROPHONE},
{device::DeviceKind::Motor,
cmvr::quic_edge::v1::DEVICE_KIND_MOTOR},
{device::DeviceKind::MotorSystem,
cmvr::quic_edge::v1::DEVICE_KIND_MOTOR_SYSTEM},
{device::DeviceKind::Robot,
cmvr::quic_edge::v1::DEVICE_KIND_ROBOT},
{device::DeviceKind::Speaker,
cmvr::quic_edge::v1::DEVICE_KIND_SPEAKER},
};
const std::vector<std::pair<device::ManagedDeviceState, ProtoState>> states{
{device::ManagedDeviceState::Unknown,
cmvr::quic_edge::v1::MANAGED_DEVICE_STATE_UNSPECIFIED},
{device::ManagedDeviceState::Disabled,
cmvr::quic_edge::v1::MANAGED_DEVICE_STATE_DISABLED},
{device::ManagedDeviceState::Initializing,
cmvr::quic_edge::v1::MANAGED_DEVICE_STATE_INITIALIZING},
{device::ManagedDeviceState::Registered,
cmvr::quic_edge::v1::MANAGED_DEVICE_STATE_REGISTERED},
{device::ManagedDeviceState::Ready,
cmvr::quic_edge::v1::MANAGED_DEVICE_STATE_READY},
{device::ManagedDeviceState::Running,
cmvr::quic_edge::v1::MANAGED_DEVICE_STATE_RUNNING},
{device::ManagedDeviceState::Stopped,
cmvr::quic_edge::v1::MANAGED_DEVICE_STATE_STOPPED},
{device::ManagedDeviceState::Error,
cmvr::quic_edge::v1::MANAGED_DEVICE_STATE_ERROR},
};
const std::vector<std::pair<device::DeviceHealthState, ProtoHealth>> health{
{device::DeviceHealthState::Unknown,
cmvr::quic_edge::v1::DEVICE_HEALTH_STATUS_UNSPECIFIED},
{device::DeviceHealthState::Healthy,
cmvr::quic_edge::v1::DEVICE_HEALTH_STATUS_HEALTHY},
{device::DeviceHealthState::Degraded,
cmvr::quic_edge::v1::DEVICE_HEALTH_STATUS_DEGRADED},
{device::DeviceHealthState::Fault,
cmvr::quic_edge::v1::DEVICE_HEALTH_STATUS_FAULT},
};
device::DeviceManagerSnapshot snapshot;
snapshot.name = "mapping-test";
for (std::size_t index = 0U; index < kinds.size(); ++index) {
device::ManagedDeviceSnapshot row;
row.id = std::string("kind-") + (index < 10U ? "0" : "") +
std::to_string(index);
row.kind = kinds[index].first;
row.type_name = "mapping";
row.enabled = true;
row.state = states[index % states.size()].first;
row.health.state = health[index % health.size()].first;
row.abnormal = index % 2U != 0U;
snapshot.devices.push_back(std::move(row));
}
media::MediaSourceManager hub;
auto transport = std::make_unique<FakeTransport>();
FakeTransport* transport_view = transport.get();
quic_edge::QuicEdgeService service(
validPresenceOnlyConfig(), std::move(transport), hub,
[snapshot]() { return snapshot; });
std::string error;
CHECK_TRUE(service.initialize(&error));
CHECK_TRUE(service.start(&error));
CHECK_TRUE(waitUntil([&]() {
return service.stats().heartbeats_acknowledged >= 1U;
}));
const auto heartbeat = transport_view->lastHeartbeat();
service.stop();
CHECK_TRUE(heartbeat.device_manager().devices_size() ==
static_cast<int>(kinds.size()));
for (std::size_t index = 0U; index < kinds.size(); ++index) {
const auto& row =
heartbeat.device_manager().devices(static_cast<int>(index));
CHECK_TRUE(row.kind() == kinds[index].second);
CHECK_TRUE(row.manager_state() ==
states[index % states.size()].second);
CHECK_TRUE(row.health() == health[index % health.size()].second);
CHECK_TRUE(row.has_error() == (index % 2U != 0U));
}
return true;
}
bool testConfiguredHeartbeatIntervalWithoutGatewayOverride()
{
media::MediaSourceManager hub;
auto transport = std::make_unique<FakeTransport>(
true, true, true, 0U, 0U);
FakeTransport* transport_view = transport.get();
auto config = validPresenceOnlyConfig();
config.set_heartbeat_interval_ms(400U);
quic_edge::QuicEdgeService service(
config, std::move(transport), hub);
std::string error;
CHECK_TRUE(service.initialize(&error));
CHECK_TRUE(service.start(&error));
CHECK_TRUE(waitUntil([&]() {
return transport_view->heartbeatCount() >= 3U;
}));
const auto intervals = transport_view->heartbeatIntervals();
service.stop();
CHECK_TRUE(intervals.size() >= 2U);
CHECK_TRUE(intervals[0].count() >= 350);
CHECK_TRUE(intervals[1].count() >= 350);
CHECK_TRUE(intervals[0].count() <= 900);
CHECK_TRUE(intervals[1].count() <= 900);
return true;
}
bool testSnapshotLatencyDoesNotConsumeAckDeadline()
{
media::MediaSourceManager hub;
auto transport = std::make_unique<FakeTransport>(
true, true, true, 0U, 250U, 70U);
FakeTransport* transport_view = transport.get();
auto config = validPresenceOnlyConfig();
config.set_control_response_timeout_ms(100U);
quic_edge::QuicEdgeService service(
config, std::move(transport), hub, [] {
std::this_thread::sleep_for(std::chrono::milliseconds(120));
return device::DeviceManagerSnapshot{};
});
std::string error;
CHECK_TRUE(service.initialize(&error));
CHECK_TRUE(service.start(&error));
CHECK_TRUE(waitUntil([&]() {
return service.stats().heartbeats_acknowledged >= 2U;
}));
CHECK_TRUE(transport_view->connectCount() == 1U);
CHECK_TRUE(service.stats().heartbeat_timeouts == 0U);
service.stop();
return true;
}
bool testHeartbeatTimeoutReconnectsWithoutTaskFailure()
{
media::MediaSourceManager hub;
auto transport = std::make_unique<FakeTransport>(true, false);
FakeTransport* transport_view = transport.get();
auto config = validPresenceOnlyConfig();
config.set_control_response_timeout_ms(20U);
quic_edge::QuicEdgeService service(config, std::move(transport), hub);
std::string error;
CHECK_TRUE(service.initialize(&error));
CHECK_TRUE(service.start(&error));
CHECK_TRUE(waitUntil([&]() {
return service.stats().heartbeat_timeouts >= 1U &&
transport_view->connectCount() >= 2U;
}));
CHECK_TRUE(service.state() != quic_edge::QuicEdgeServiceState::FAILED);
service.stop();
return true;
}
bool testRegistrationRejectionBacksOff()
{
media::MediaSourceManager hub;
auto transport = std::make_unique<FakeTransport>(false, true);
FakeTransport* transport_view = transport.get();
auto config = validPresenceOnlyConfig();
config.mutable_reconnect()->set_initial_delay_ms(20U);
config.mutable_reconnect()->set_maximum_delay_ms(80U);
quic_edge::QuicEdgeService service(
config, std::move(transport), hub);
std::string error;
CHECK_TRUE(service.initialize(&error));
CHECK_TRUE(service.start(&error));
CHECK_TRUE(waitUntil([&]() {
return service.stats().registrations_rejected >= 3U &&
transport_view->connectCount() >= 4U;
}));
CHECK_TRUE(service.stats().heartbeats_sent == 0U);
const auto intervals = transport_view->connectIntervals();
CHECK_TRUE(intervals.size() >= 3U);
CHECK_TRUE(intervals[0].count() >= 15);
CHECK_TRUE(intervals[1].count() >= 35);
CHECK_TRUE(intervals[2].count() >= 70);
CHECK_TRUE(service.state() != quic_edge::QuicEdgeServiceState::FAILED);
service.stop();
return true;
}
bool testHeartbeatAckRequiresSessionId()
{
media::MediaSourceManager hub;
auto transport = std::make_unique<FakeTransport>(true, true, false);
FakeTransport* transport_view = transport.get();
quic_edge::QuicEdgeService service(
validPresenceOnlyConfig(), std::move(transport), hub);
std::string error;
CHECK_TRUE(service.initialize(&error));
CHECK_TRUE(service.start(&error));
CHECK_TRUE(waitUntil([&]() {
return transport_view->connectCount() >= 2U &&
service.stats().heartbeats_sent >= 2U;
}));
CHECK_TRUE(service.stats().heartbeats_acknowledged == 0U);
CHECK_TRUE(service.state() != quic_edge::QuicEdgeServiceState::FAILED);
service.stop();
return true;
}
bool testSlowMediaStartDoesNotBlockHeartbeat()
{
const std::string track_id = "slow-camera/video/color";
media::MediaSourceManager hub;
std::atomic<bool> start_entered{false};
std::atomic<bool> start_exited{false};
std::atomic<bool> release_start{false};
media::MediaSourceManager::SourceCallbacks callbacks;
callbacks.start = [&](const media::MediaSourceManager::FrameSink&,
const media::MediaSourceManager::CancelPredicate& cancelled) {
start_entered.store(true);
while (!release_start.load() && !cancelled()) {
std::this_thread::sleep_for(std::chrono::milliseconds(2));
}
start_exited.store(true);
return !cancelled();
};
callbacks.stop = []() {};
CHECK_TRUE(hub.registerSource(
videoDescriptor(track_id), std::move(callbacks), 8U));
auto transport = std::make_unique<FakeTransport>();
FakeTransport* transport_view = transport.get();
quic_edge::QuicEdgeService service(
validConfig(track_id), std::move(transport), hub);
std::string error;
CHECK_TRUE(service.initialize(&error));
CHECK_TRUE(service.start(&error));
const bool entered = waitUntil([&]() { return start_entered.load(); });
const bool heartbeat_while_blocked = entered && waitUntil([&]() {
return service.stats().heartbeats_acknowledged >= 1U;
});
const bool registered_while_blocked = service.status().registered;
auto stop_future = std::async(std::launch::async, [&service] {
service.stop();
});
const bool stop_while_start_blocked =
stop_future.wait_for(std::chrono::milliseconds(500)) ==
std::future_status::ready;
release_start.store(true);
stop_future.get();
const bool start_cancelled = waitUntil([&]() { return start_exited.load(); });
CHECK_TRUE(entered);
CHECK_TRUE(heartbeat_while_blocked);
CHECK_TRUE(registered_while_blocked);
CHECK_TRUE(stop_while_start_blocked);
CHECK_TRUE(start_cancelled);
CHECK_TRUE(transport_view->connectCount() == 1U);
CHECK_TRUE(transport_view->controlSequencesStrictlyIncreasing());
return true;
}
bool testLifecycleStateGuards()
{
media::MediaSourceManager hub;
{
auto transport = std::make_unique<FakeTransport>();
quic_edge::QuicEdgeService service(
validPresenceOnlyConfig(), std::move(transport), hub);
service.stop();
CHECK_TRUE(service.state() ==
quic_edge::QuicEdgeServiceState::UNINITIALIZED);
std::string error;
CHECK_TRUE(!service.start(&error));
}
{
auto transport = std::make_unique<FakeTransport>();
quic_edge::QuicEdgeService service(
validPresenceOnlyConfig(), std::move(transport), hub);
std::string error;
CHECK_TRUE(service.initialize(&error));
CHECK_TRUE(service.start(&error));
CHECK_TRUE(waitUntil([&]() {
return service.stats().registrations_accepted == 1U;
}));
error.clear();
CHECK_TRUE(!service.initialize(&error));
CHECK_TRUE(!error.empty());
service.stop();
}
return true;
}
bool testRobotIdIsRequired()
{
media::MediaSourceManager hub;
auto config = validPresenceOnlyConfig();
config.clear_robot_id();
auto transport = std::make_unique<FakeTransport>();
quic_edge::QuicEdgeService service(
std::move(config), std::move(transport), hub);
std::string error;
CHECK_TRUE(!service.initialize(&error));
CHECK_TRUE(error.find("robot_id") != std::string::npos);
return true;
}
} // namespace
int main()
{
if (!testControlFraming() || !testPacketizer() ||
!testServiceWithSharedHub() ||
!testMediaActivityInterruptPreservesPresenceAndResumes() ||
!testMediaActivityInterruptCancelsStartingSubscription() ||
!testMissingInjectedSourceRetriesSafely() ||
!testPresenceOnlyWithoutMedia() ||
!testDeviceManagerSnapshotInHeartbeat() ||
!testAllDeviceKindAndStateMappings() ||
!testConfiguredHeartbeatIntervalWithoutGatewayOverride() ||
!testSnapshotLatencyDoesNotConsumeAckDeadline() ||
!testHeartbeatTimeoutReconnectsWithoutTaskFailure() ||
!testRegistrationRejectionBacksOff() ||
!testHeartbeatAckRequiresSessionId() ||
!testSlowMediaStartDoesNotBlockHeartbeat() ||
!testLifecycleStateGuards() || !testRobotIdIsRequired()) {
return 1;
}
std::cout << "quic_edge_protocol_test: PASS\n";
return 0;
}