From c6617b89cd8f4e50824655e28dd56f0e43012720 Mon Sep 17 00:00:00 2001 From: linbo <1034003879@qq.com> Date: Mon, 24 Aug 2026 16:00:51 +0800 Subject: [PATCH] fix(quic): align media hub targets with linbo branch --- .../include/media_source_hub.h | 4 + cmvr-es/service/quic_edge/CMakeLists.txt | 5 +- .../quic_edge/include/quic_edge_service.h | 8 +- .../quic_edge/include/quic_edge_types.h | 2 +- .../src/quic_edge_device_adapter.cpp | 6 +- .../quic_edge/src/quic_edge_service.cpp | 32 +++----- .../tests/quic_edge_protocol_test.cpp | 14 +--- cmvr-es/task/quic_edge_task/CMakeLists.txt | 2 - .../quic_edge_task/src/quic_edge_task.cpp | 75 ++++++------------- .../tests/quic_edge_task_test.cpp | 51 +------------ 10 files changed, 53 insertions(+), 146 deletions(-) diff --git a/cmvr-es/manager/media_source_hub/include/media_source_hub.h b/cmvr-es/manager/media_source_hub/include/media_source_hub.h index c7dcb553..39daa7bd 100644 --- a/cmvr-es/manager/media_source_hub/include/media_source_hub.h +++ b/cmvr-es/manager/media_source_hub/include/media_source_hub.h @@ -123,6 +123,10 @@ private: std::shared_ptr impl_; }; +// Compatibility alias for older protocol tests and integrations. New code should +// use MediaSourceHub directly. +using MediaSourceManager = MediaSourceHub; + } // namespace cmvr::media #endif // CMVR_ES_MANAGER_MEDIA_SOURCE_HUB_H diff --git a/cmvr-es/service/quic_edge/CMakeLists.txt b/cmvr-es/service/quic_edge/CMakeLists.txt index 3efed1be..6d349e99 100644 --- a/cmvr-es/service/quic_edge/CMakeLists.txt +++ b/cmvr-es/service/quic_edge/CMakeLists.txt @@ -14,7 +14,7 @@ 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_manager + cmvr_es::media_source_hub PRIVATE cmvr_es::device_media_source_adapter cmvr_es::device_manager @@ -29,8 +29,7 @@ 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_manager - cmvr_es::stop_all_admission_gate + cmvr_es::media_source_hub Threads::Threads ) add_test(NAME quic_edge_protocol_test COMMAND quic_edge_protocol_test) diff --git a/cmvr-es/service/quic_edge/include/quic_edge_service.h b/cmvr-es/service/quic_edge/include/quic_edge_service.h index 9b86af96..23eae61b 100644 --- a/cmvr-es/service/quic_edge/include/quic_edge_service.h +++ b/cmvr-es/service/quic_edge/include/quic_edge_service.h @@ -15,7 +15,7 @@ #include "cmvr/config/quic_edge_config/quic_edge_config.pb.h" #include "devices/device_types.h" -#include "manager/media_source_manager/include/media_source_manager.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_transport.h" @@ -76,7 +76,7 @@ public: explicit QuicEdgeService(config::QuicEdgeConfig config); QuicEdgeService(config::QuicEdgeConfig config, std::unique_ptr transport, - media::MediaSourceManager& media_hub, + media::MediaSourceHub& media_hub, DeviceSnapshotProvider device_snapshot_provider = {}); ~QuicEdgeService(); @@ -106,7 +106,7 @@ private: struct ActiveTrack { config::QuicEdgeTrackConfig config; std::string source_track_id; - media::MediaSourceManager::Subscription subscription; + media::MediaSourceHub::Subscription subscription; std::optional last_description; bool waiting_for_keyframe{false}; bool keyframe_requested{false}; @@ -157,7 +157,7 @@ private: config::QuicEdgeConfig config_; std::unique_ptr transport_; - media::MediaSourceManager* media_hub_{nullptr}; + media::MediaSourceHub* media_hub_{nullptr}; bool using_global_media_hub_{false}; SourceRegistrar source_registrar_; DeviceSnapshotProvider device_snapshot_provider_; diff --git a/cmvr-es/service/quic_edge/include/quic_edge_types.h b/cmvr-es/service/quic_edge/include/quic_edge_types.h index cabfdd36..36fda822 100644 --- a/cmvr-es/service/quic_edge/include/quic_edge_types.h +++ b/cmvr-es/service/quic_edge/include/quic_edge_types.h @@ -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 MediaSourceManager descriptor generation +// is a compact token for the full 64-bit MediaSourceHub descriptor generation // announced on the reliable control stream. struct DatagramHeader { std::uint8_t protocol_version{kProtocolVersion}; diff --git a/cmvr-es/service/quic_edge/src/quic_edge_device_adapter.cpp b/cmvr-es/service/quic_edge/src/quic_edge_device_adapter.cpp index 0e1ba0fc..856535dc 100644 --- a/cmvr-es/service/quic_edge/src/quic_edge_device_adapter.cpp +++ b/cmvr-es/service/quic_edge/src/quic_edge_device_adapter.cpp @@ -3,14 +3,14 @@ #include #include "manager/device_manager/include/device_manager.h" -#include "manager/media_source_manager/include/device_media_source_adapter.h" +#include "manager/media_source_hub/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::globalMediaSourceManager()), + media_hub_(&media::globalMediaSourceHub()), 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 MediaSourceManager source: " + + *error = "failed to register MediaSourceHub source: " + source_track_id; } return false; diff --git a/cmvr-es/service/quic_edge/src/quic_edge_service.cpp b/cmvr-es/service/quic_edge/src/quic_edge_service.cpp index 2ea62c79..30a2f14d 100644 --- a/cmvr-es/service/quic_edge/src/quic_edge_service.cpp +++ b/cmvr-es/service/quic_edge/src/quic_edge_service.cpp @@ -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_manager/include/device_media_source_adapter.h" +#include "manager/media_source_hub/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 transport, - media::MediaSourceManager& media_hub, + media::MediaSourceHub& 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 MediaSourceManager is null"; + const std::string message = "QUIC edge transport or MediaSourceHub is null"; setState(QuicEdgeServiceState::FAILED, message); setError(error, message); return false; @@ -1450,21 +1450,9 @@ 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::MediaSourceManager::StartPosition::LATEST_AVAILABLE, + media::MediaSourceHub::StartPosition::LATEST_AVAILABLE, [this, activity_generation] { std::lock_guard lock(mutex_); return stop_requested_ || media_stop_requested_ || @@ -1472,7 +1460,7 @@ void QuicEdgeService::refreshMediaTracks( }); if (!track.subscription.valid()) { recordMediaError( - "MediaSourceManager source unavailable: " + track.source_track_id); + "MediaSourceHub source unavailable: " + track.source_track_id); continue; } track.waiting_for_keyframe = @@ -1500,14 +1488,14 @@ bool QuicEdgeService::ensureSourceRegistered( { if (media_hub_->hasSource(source_track_id)) return true; if (!using_global_media_hub_ || !source_registrar_) { - setError(error, "MediaSourceManager source unavailable: " + source_track_id); + 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 MediaSourceManager source: " + source_track_id); + "failed to register MediaSourceHub source: " + source_track_id); } return false; } @@ -1638,7 +1626,7 @@ bool QuicEdgeService::processTrack(ActiveTrack* track, if (!read) { if (!track->subscription.valid()) { recordMediaError( - "MediaSourceManager subscription stopped: " + track->source_track_id); + "MediaSourceHub subscription stopped: " + track->source_track_id); } return true; } @@ -1647,7 +1635,7 @@ bool QuicEdgeService::processTrack(ActiveTrack* track, if (!frame || !frame->descriptor || frame->descriptor->id != track->source_track_id || frame->descriptor->kind != expectedKind(track->config)) { - recordMediaError("MediaSourceManager returned an invalid or mismatched frame"); + recordMediaError("MediaSourceHub returned an invalid or mismatched frame"); track->next_frame_discontinuous = true; return true; } @@ -1667,7 +1655,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("MediaSourceManager descriptor generation regressed"); + recordMediaError("MediaSourceHub descriptor generation regressed"); track->next_frame_discontinuous = true; return true; } diff --git a/cmvr-es/service/quic_edge/tests/quic_edge_protocol_test.cpp b/cmvr-es/service/quic_edge/tests/quic_edge_protocol_test.cpp index b652a904..f97ac7f1 100644 --- a/cmvr-es/service/quic_edge/tests/quic_edge_protocol_test.cpp +++ b/cmvr-es/service/quic_edge/tests/quic_edge_protocol_test.cpp @@ -14,11 +14,10 @@ #include #include "cmvr/quic_edge/v1/quic_edge.pb.h" -#include "manager/media_source_manager/include/media_source_manager.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" -#include "service/grpc/stop_all/include/stop_all_admission_gate.h" namespace { @@ -580,8 +579,7 @@ bool testServiceWithSharedHub() bool testMediaActivityInterruptPreservesPresenceAndResumes() { const std::string track_id = "camera-stop-all/video/color"; - service::StopAllAdmissionGate admission; - media::MediaSourceManager hub(&admission); + media::MediaSourceHub hub; media::MediaSourceManager::FrameSink sink; std::mutex sink_mutex; std::atomic source_started{false}; @@ -642,7 +640,6 @@ 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()); @@ -656,8 +653,6 @@ 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; @@ -679,8 +674,7 @@ bool testMediaActivityInterruptPreservesPresenceAndResumes() bool testMediaActivityInterruptCancelsStartingSubscription() { const std::string track_id = "slow-stop-all/video/color"; - service::StopAllAdmissionGate admission; - media::MediaSourceManager hub(&admission); + media::MediaSourceHub hub; std::atomic start_entered{false}; std::atomic start_cancelled{false}; media::MediaSourceManager::SourceCallbacks callbacks; @@ -708,7 +702,6 @@ 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(); }); @@ -720,7 +713,6 @@ 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; } diff --git a/cmvr-es/task/quic_edge_task/CMakeLists.txt b/cmvr-es/task/quic_edge_task/CMakeLists.txt index 8b2db539..2a8e70d9 100644 --- a/cmvr-es/task/quic_edge_task/CMakeLists.txt +++ b/cmvr-es/task/quic_edge_task/CMakeLists.txt @@ -9,7 +9,6 @@ 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) @@ -18,7 +17,6 @@ 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) diff --git a/cmvr-es/task/quic_edge_task/src/quic_edge_task.cpp b/cmvr-es/task/quic_edge_task/src/quic_edge_task.cpp index fe5d2293..d7ddbd21 100644 --- a/cmvr-es/task/quic_edge_task/src/quic_edge_task.cpp +++ b/cmvr-es/task/quic_edge_task/src/quic_edge_task.cpp @@ -8,7 +8,6 @@ #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" @@ -218,62 +217,36 @@ bool QuicEdgeTask::init() bool QuicEdgeTask::start() { - 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::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; } - { - 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"; + 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; } - - 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(); + ++started_services; } - - 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; + last_error_.clear(); + state_ = TaskState::RUNNING; + CMVR_LOG(INFO) << "[QuicEdgeTask] Started, id=" << id_ + << ", platforms=" << services_.size(); + return true; } bool QuicEdgeTask::step(const double dt) diff --git a/cmvr-es/task/quic_edge_task/tests/quic_edge_task_test.cpp b/cmvr-es/task/quic_edge_task/tests/quic_edge_task_test.cpp index 06e6124c..a2f13f5c 100644 --- a/cmvr-es/task/quic_edge_task/tests/quic_edge_task_test.cpp +++ b/cmvr-es/task/quic_edge_task/tests/quic_edge_task_test.cpp @@ -12,9 +12,8 @@ #include #include "cmvr/quic_edge/v1/quic_edge.pb.h" -#include "manager/media_source_manager/include/media_source_manager.h" +#include "manager/media_source_hub/include/media_source_hub.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 { @@ -235,9 +234,6 @@ cmvr::config::QuicEdgeConfig multiPlatformConfig() int main() { - auto& admission = cmvr::service::globalStopAllAdmissionGate(); - admission.clearForTesting(); - cmvr::config::QuicEdgeConfig config; config.set_id("quic-invalid-config-test"); @@ -256,50 +252,7 @@ int main() return 1; } - 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; + cmvr::media::MediaSourceHub media_hub; std::vector transports; std::vector services; std::vector expanded_configs;