fix(merge): align QUIC media and SocketCAN interfaces

This commit is contained in:
linbo 2026-08-27 11:29:07 +08:00
parent 19ac37b784
commit 0ac5e79751
11 changed files with 167 additions and 51 deletions

View File

@ -14,9 +14,9 @@ target_include_directories(quic_edge_service PUBLIC ${PROJECT_SOURCE_DIR}/cmvr-e
target_link_libraries(quic_edge_service
PUBLIC
cmvr_es::proto
cmvr_es::media_source_hub
cmvr_es::media_source_manager
PRIVATE
cmvr_es::media_source_hub_device_adapter
cmvr_es::device_media_source_adapter
cmvr_es::device_manager
Threads::Threads
msquic
@ -29,7 +29,8 @@ if(BUILD_TESTING)
target_compile_features(quic_edge_protocol_test PRIVATE cxx_std_17)
target_link_libraries(quic_edge_protocol_test PRIVATE
cmvr_es::quic_edge_service
cmvr_es::media_source_hub
cmvr_es::media_source_manager
cmvr_es::stop_all_admission_gate
Threads::Threads
)
add_test(NAME quic_edge_protocol_test COMMAND quic_edge_protocol_test)

View File

@ -15,7 +15,7 @@
#include "cmvr/config/quic_edge_config/quic_edge_config.pb.h"
#include "devices/device_types.h"
#include "manager/media_source_hub/include/media_source_hub.h"
#include "manager/media_source_manager/include/media_source_manager.h"
#include "service/quic_edge/include/control_framing.h"
#include "service/quic_edge/include/datagram_packetizer.h"
#include "service/quic_edge/include/quic_transport.h"
@ -76,7 +76,7 @@ public:
explicit QuicEdgeService(config::QuicEdgeConfig config);
QuicEdgeService(config::QuicEdgeConfig config,
std::unique_ptr<QuicTransport> transport,
media::MediaSourceHub& media_hub,
media::MediaSourceManager& media_hub,
DeviceSnapshotProvider device_snapshot_provider = {});
~QuicEdgeService();
@ -106,7 +106,7 @@ private:
struct ActiveTrack {
config::QuicEdgeTrackConfig config;
std::string source_track_id;
media::MediaSourceHub::Subscription subscription;
media::MediaSourceManager::Subscription subscription;
std::optional<MediaTrackDescription> last_description;
bool waiting_for_keyframe{false};
bool keyframe_requested{false};
@ -157,7 +157,7 @@ private:
config::QuicEdgeConfig config_;
std::unique_ptr<QuicTransport> transport_;
media::MediaSourceHub* media_hub_{nullptr};
media::MediaSourceManager* media_hub_{nullptr};
bool using_global_media_hub_{false};
SourceRegistrar source_registrar_;
DeviceSnapshotProvider device_snapshot_provider_;

View File

@ -22,7 +22,7 @@ enum DatagramFlag : std::uint16_t {
};
// All integer fields are serialized in network byte order. codec_generation
// is a compact token for the full 64-bit MediaSourceHub descriptor generation
// is a compact token for the full 64-bit MediaSourceManager descriptor generation
// announced on the reliable control stream.
struct DatagramHeader {
std::uint8_t protocol_version{kProtocolVersion};

View File

@ -3,14 +3,14 @@
#include <utility>
#include "manager/device_manager/include/device_manager.h"
#include "manager/media_source_hub/include/device_media_source_adapter.h"
#include "manager/media_source_manager/include/device_media_source_adapter.h"
namespace cmvr::quic_edge {
QuicEdgeService::QuicEdgeService(config::QuicEdgeConfig config)
: config_(std::move(config)),
transport_(createDefaultQuicTransport(config_.datagram_send_queue_depth())),
media_hub_(&media::globalMediaSourceHub()),
media_hub_(&media::globalMediaSourceManager()),
using_global_media_hub_(true),
device_snapshot_provider_([] {
return device::DeviceManager::getInstance().snapshot();
@ -57,7 +57,7 @@ QuicEdgeService::QuicEdgeService(config::QuicEdgeConfig config)
}
if (!registered && !media_hub_->hasSource(source_track_id)) {
if (error) {
*error = "failed to register MediaSourceHub source: " +
*error = "failed to register MediaSourceManager source: " +
source_track_id;
}
return false;

View File

@ -25,7 +25,7 @@
#include "cmvr/quic_edge/v1/quic_edge.pb.h"
#include "common/base/logging/logger.h"
#include "manager/device_manager/include/device_manager.h"
#include "manager/media_source_hub/include/device_media_source_adapter.h"
#include "manager/media_source_manager/include/device_media_source_adapter.h"
namespace cmvr::quic_edge {
namespace {
@ -423,7 +423,7 @@ const char* toString(const QuicEdgeServiceState state)
QuicEdgeService::QuicEdgeService(config::QuicEdgeConfig config,
std::unique_ptr<QuicTransport> transport,
media::MediaSourceHub& media_hub,
media::MediaSourceManager& media_hub,
DeviceSnapshotProvider device_snapshot_provider)
: config_(std::move(config)),
transport_(std::move(transport)),
@ -598,7 +598,7 @@ bool QuicEdgeService::initialize(std::string* error)
return false;
}
if (!transport_ || !media_hub_) {
const std::string message = "QUIC edge transport or MediaSourceHub is null";
const std::string message = "QUIC edge transport or MediaSourceManager is null";
setState(QuicEdgeServiceState::FAILED, message);
setError(error, message);
return false;
@ -1240,6 +1240,17 @@ bool QuicEdgeService::sendHeartbeat(const std::uint64_t sequence,
return false;
}
if (!sendControlEnvelope(serialized, error)) return false;
CMVR_LOG(INFO) << "[QuicEdgeService] sent NodeHeartbeat"
<< ", message_sequence=" << envelope.message_sequence()
<< ", payload_bytes=" << serialized.size()
<< ", robot_id=" << heartbeat->robot_id()
<< ", 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;
}
@ -1439,9 +1450,21 @@ void QuicEdgeService::refreshMediaTracks(
recordMediaError(source_error);
continue;
}
safety::DispatchGuard source_dispatch;
if (using_global_media_hub_) {
source_dispatch = media::beginMediaSourceStartDispatch(
device::DeviceManager::getInstance().safetyManager(),
track_config.device_id());
if (!source_dispatch.acquired()) {
recordMediaError(
"MediaSourceManager safety admission rejected: " +
track.source_track_id);
continue;
}
}
track.subscription = media_hub_->subscribe(
track.source_track_id,
media::MediaSourceHub::StartPosition::LATEST_AVAILABLE,
media::MediaSourceManager::StartPosition::LATEST_AVAILABLE,
[this, activity_generation] {
std::lock_guard lock(mutex_);
return stop_requested_ || media_stop_requested_ ||
@ -1449,7 +1472,7 @@ void QuicEdgeService::refreshMediaTracks(
});
if (!track.subscription.valid()) {
recordMediaError(
"MediaSourceHub source unavailable: " + track.source_track_id);
"MediaSourceManager source unavailable: " + track.source_track_id);
continue;
}
track.waiting_for_keyframe =
@ -1477,14 +1500,14 @@ bool QuicEdgeService::ensureSourceRegistered(
{
if (media_hub_->hasSource(source_track_id)) return true;
if (!using_global_media_hub_ || !source_registrar_) {
setError(error, "MediaSourceHub source unavailable: " + source_track_id);
setError(error, "MediaSourceManager 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);
"failed to register MediaSourceManager source: " + source_track_id);
}
return false;
}
@ -1615,7 +1638,7 @@ bool QuicEdgeService::processTrack(ActiveTrack* track,
if (!read) {
if (!track->subscription.valid()) {
recordMediaError(
"MediaSourceHub subscription stopped: " + track->source_track_id);
"MediaSourceManager subscription stopped: " + track->source_track_id);
}
return true;
}
@ -1624,7 +1647,7 @@ bool QuicEdgeService::processTrack(ActiveTrack* track,
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");
recordMediaError("MediaSourceManager returned an invalid or mismatched frame");
track->next_frame_discontinuous = true;
return true;
}
@ -1644,7 +1667,7 @@ bool QuicEdgeService::processTrack(ActiveTrack* track,
!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");
recordMediaError("MediaSourceManager descriptor generation regressed");
track->next_frame_discontinuous = true;
return true;
}

View File

@ -14,10 +14,11 @@
#include <vector>
#include "cmvr/quic_edge/v1/quic_edge.pb.h"
#include "manager/media_source_hub/include/media_source_hub.h"
#include "manager/media_source_manager/include/media_source_manager.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"
#include "service/grpc/stop_all/include/stop_all_admission_gate.h"
namespace {
@ -579,7 +580,8 @@ bool testServiceWithSharedHub()
bool testMediaActivityInterruptPreservesPresenceAndResumes()
{
const std::string track_id = "camera-stop-all/video/color";
media::MediaSourceHub hub;
service::StopAllAdmissionGate admission;
media::MediaSourceManager hub(&admission);
media::MediaSourceManager::FrameSink sink;
std::mutex sink_mutex;
std::atomic<bool> source_started{false};
@ -640,6 +642,7 @@ bool testMediaActivityInterruptPreservesPresenceAndResumes()
}));
const auto heartbeat_before_stop = transport_view->heartbeatCount();
const auto ticket = admission.beginStopAll();
CHECK_TRUE(edge_service.interruptMediaActivities());
CHECK_TRUE(hub.subscriberCount(track_id) == 0U);
CHECK_TRUE(!source_started.load());
@ -653,6 +656,8 @@ bool testMediaActivityInterruptPreservesPresenceAndResumes()
CHECK_TRUE(publish(2U));
std::this_thread::sleep_for(std::chrono::milliseconds(30));
CHECK_TRUE(transport_view->datagramBatchCount() == 1U);
CHECK_TRUE(admission.finishStopAll(ticket, true));
CHECK_TRUE(waitUntil([&] {
return source_starts.load() == 2U && source_started.load() &&
hub.subscriberCount(track_id) == 1U;
@ -674,7 +679,8 @@ bool testMediaActivityInterruptPreservesPresenceAndResumes()
bool testMediaActivityInterruptCancelsStartingSubscription()
{
const std::string track_id = "slow-stop-all/video/color";
media::MediaSourceHub hub;
service::StopAllAdmissionGate admission;
media::MediaSourceManager hub(&admission);
std::atomic<bool> start_entered{false};
std::atomic<bool> start_cancelled{false};
media::MediaSourceManager::SourceCallbacks callbacks;
@ -702,6 +708,7 @@ bool testMediaActivityInterruptCancelsStartingSubscription()
return start_entered.load() && edge_service.status().registered;
}));
const auto ticket = admission.beginStopAll();
auto interrupt = std::async(std::launch::async, [&edge_service] {
return edge_service.interruptMediaActivities();
});
@ -713,6 +720,7 @@ bool testMediaActivityInterruptCancelsStartingSubscription()
CHECK_TRUE(edge_service.state() == quic_edge::QuicEdgeServiceState::ONLINE);
CHECK_TRUE(edge_service.status().registered);
CHECK_TRUE(transport_view->connectCount() == 1U);
CHECK_TRUE(admission.finishStopAll(ticket, true));
edge_service.stop();
return true;
}

View File

@ -9,6 +9,7 @@ target_link_libraries(quic_edge_task
cmvr_es::proto
cmvr_es::logging
cmvr_es::device_manager
cmvr_es::stop_all_admission_gate
)
add_library(cmvr_es::quic_edge_task ALIAS quic_edge_task)
@ -17,6 +18,7 @@ if(BUILD_TESTING)
target_compile_features(quic_edge_task_test PRIVATE cxx_std_17)
target_link_libraries(quic_edge_task_test PRIVATE
cmvr_es::quic_edge_task
cmvr_es::stop_all_admission_gate
)
add_test(NAME quic_edge_task_test COMMAND quic_edge_task_test)
if(UNIX AND NOT APPLE)

View File

@ -29,7 +29,7 @@ public:
bool start() override;
bool step(double dt) override;
void stop() override;
bool stopActivity();
bool stopActivity() override;
TaskState state() const override;
bool isBusy() const override;

View File

@ -8,6 +8,7 @@
#include "cmvr/config/task_manager_config/task_manager_config.pb.h"
#include "common/base/logging/logger.h"
#include "service/grpc/stop_all/include/stop_all_admission_gate.h"
#include "common/config/config_files.h"
#include "manager/device_manager/include/device_manager.h"
#include "task/task_factory.h"
@ -217,36 +218,62 @@ bool QuicEdgeTask::init()
bool QuicEdgeTask::start()
{
std::lock_guard lock(mutex_);
if (state_ == TaskState::RUNNING) {
return true;
}
if ((state_ != TaskState::IDLE && state_ != TaskState::STOPPED) ||
services_.empty()) {
last_error_ = "QUIC edge task is not initialized";
state_ = TaskState::FAILED;
return false;
auto& admission_gate = service::globalStopAllAdmissionGate();
std::uint64_t admission_generation = 0U;
{
auto admission = admission_gate.lockAdmission();
if (!admission.accepting()) {
return false;
}
admission_generation = admission.generation();
}
std::size_t started_services = 0U;
for (auto& platform : services_) {
std::string error;
if (!platform.service->start(&error)) {
for (std::size_t index = 0U;
index < started_services; ++index) {
services_[index].service->stop();
}
last_error_ = "platform " + platform.id + ": " + error;
{
std::lock_guard lock(mutex_);
if (state_ == TaskState::RUNNING) {
auto admission = admission_gate.lockAdmission();
return admission.accepting() &&
admission.generation() == admission_generation;
}
if ((state_ != TaskState::IDLE && state_ != TaskState::STOPPED) ||
services_.empty()) {
last_error_ = "QUIC edge task is not initialized";
state_ = TaskState::FAILED;
return false;
}
++started_services;
std::size_t started_services = 0U;
for (auto& platform : services_) {
std::string error;
if (!platform.service->start(&error)) {
for (std::size_t index = 0U;
index < started_services; ++index) {
services_[index].service->stop();
}
last_error_ = "platform " + platform.id + ": " + error;
state_ = TaskState::FAILED;
return false;
}
++started_services;
}
last_error_.clear();
state_ = TaskState::RUNNING;
CMVR_LOG(INFO) << "[QuicEdgeTask] Started, id=" << id_
<< ", platforms=" << services_.size();
}
last_error_.clear();
state_ = TaskState::RUNNING;
CMVR_LOG(INFO) << "[QuicEdgeTask] Started, id=" << id_
<< ", platforms=" << services_.size();
return true;
bool admission_current = false;
{
auto admission = admission_gate.lockAdmission();
admission_current = admission.accepting() &&
admission.generation() == admission_generation;
}
if (admission_current) {
return true;
}
stop();
return false;
}
bool QuicEdgeTask::step(const double dt)

View File

@ -12,8 +12,9 @@
#include <vector>
#include "cmvr/quic_edge/v1/quic_edge.pb.h"
#include "manager/media_source_hub/include/media_source_hub.h"
#include "manager/media_source_manager/include/media_source_manager.h"
#include "service/quic_edge/include/control_framing.h"
#include "service/grpc/stop_all/include/stop_all_admission_gate.h"
#include "task/quic_edge_task/include/quic_edge_task.h"
namespace {
@ -234,6 +235,9 @@ cmvr::config::QuicEdgeConfig multiPlatformConfig()
int main()
{
auto& admission = cmvr::service::globalStopAllAdmissionGate();
admission.clearForTesting();
cmvr::config::QuicEdgeConfig config;
config.set_id("quic-invalid-config-test");
@ -252,7 +256,50 @@ int main()
return 1;
}
cmvr::media::MediaSourceHub media_hub;
cmvr::task::QuicEdgeTask admission_task(validConfig());
if (!admission_task.init()) {
std::cerr << "valid QUIC admission task did not initialize\n";
return 1;
}
auto ticket = admission.beginStopAll();
if (admission_task.start() || admission_task.isBusy()) {
std::cerr << "QUIC public start bypassed closed StopAll admission\n";
return 1;
}
if (!admission.finishStopAll(ticket, true) ||
!admission_task.start() || !admission_task.isBusy()) {
std::cerr << "QUIC public start was not restored after StopAll\n";
return 1;
}
if (!admission_task.stopActivity()) {
std::cerr << "restarted QUIC activity did not stop cleanly\n";
return 1;
}
if (!admission_task.isBusy() ||
admission_task.state() != cmvr::task::TaskState::RUNNING) {
std::cerr << "QUIC activity stop terminated the service lifecycle\n";
return 1;
}
ticket = admission.beginStopAll();
if (admission.finishStopAll(ticket, false) || admission_task.start()) {
std::cerr << "failed StopAll did not keep QUIC start fail-closed\n";
return 1;
}
ticket = admission.beginStopAll();
if (!admission.finishStopAll(ticket, true) ||
!admission_task.start() || !admission_task.stopActivity()) {
std::cerr << "successful StopAll did not restore QUIC restart\n";
return 1;
}
if (!admission_task.isBusy()) {
std::cerr << "QUIC service did not remain available after StopAll\n";
return 1;
}
admission_task.stop();
admission.clearForTesting();
cmvr::media::MediaSourceManager media_hub;
std::vector<HeartbeatTransport*> transports;
std::vector<QuicEdgeService*> services;
std::vector<cmvr::config::QuicEdgeConfig> expanded_configs;

View File

@ -52,6 +52,14 @@ message EtherCATDcConfig {
message SocketCanConfig {
string dev_id = 1;
int32 channel_id = 2;
// Optional raw SocketCAN settings used by CAN-FD backends.
optional string interface_name = 3;
optional bool enable_fd = 4;
optional bool bitrate_switch = 5;
optional uint32 receive_timeout_us = 6;
optional bool receive_own_messages = 7;
optional bool enable_error_frames = 8;
optional uint32 send_timeout_us = 9;
}
message EtherCATConfig {