fix(quic): align media hub targets with linbo branch

This commit is contained in:
linbo 2026-08-24 16:00:51 +08:00
parent ab27c97e87
commit c6617b89cd
10 changed files with 53 additions and 146 deletions

View File

@ -123,6 +123,10 @@ private:
std::shared_ptr<Impl> impl_; std::shared_ptr<Impl> impl_;
}; };
// Compatibility alias for older protocol tests and integrations. New code should
// use MediaSourceHub directly.
using MediaSourceManager = MediaSourceHub;
} // namespace cmvr::media } // namespace cmvr::media
#endif // CMVR_ES_MANAGER_MEDIA_SOURCE_HUB_H #endif // CMVR_ES_MANAGER_MEDIA_SOURCE_HUB_H

View File

@ -14,7 +14,7 @@ target_include_directories(quic_edge_service PUBLIC ${PROJECT_SOURCE_DIR}/cmvr-e
target_link_libraries(quic_edge_service target_link_libraries(quic_edge_service
PUBLIC PUBLIC
cmvr_es::proto cmvr_es::proto
cmvr_es::media_source_manager cmvr_es::media_source_hub
PRIVATE PRIVATE
cmvr_es::device_media_source_adapter cmvr_es::device_media_source_adapter
cmvr_es::device_manager cmvr_es::device_manager
@ -29,8 +29,7 @@ if(BUILD_TESTING)
target_compile_features(quic_edge_protocol_test PRIVATE cxx_std_17) target_compile_features(quic_edge_protocol_test PRIVATE cxx_std_17)
target_link_libraries(quic_edge_protocol_test PRIVATE target_link_libraries(quic_edge_protocol_test PRIVATE
cmvr_es::quic_edge_service cmvr_es::quic_edge_service
cmvr_es::media_source_manager cmvr_es::media_source_hub
cmvr_es::stop_all_admission_gate
Threads::Threads Threads::Threads
) )
add_test(NAME quic_edge_protocol_test COMMAND quic_edge_protocol_test) 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 "cmvr/config/quic_edge_config/quic_edge_config.pb.h"
#include "devices/device_types.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/control_framing.h"
#include "service/quic_edge/include/datagram_packetizer.h" #include "service/quic_edge/include/datagram_packetizer.h"
#include "service/quic_edge/include/quic_transport.h" #include "service/quic_edge/include/quic_transport.h"
@ -76,7 +76,7 @@ public:
explicit QuicEdgeService(config::QuicEdgeConfig config); explicit QuicEdgeService(config::QuicEdgeConfig config);
QuicEdgeService(config::QuicEdgeConfig config, QuicEdgeService(config::QuicEdgeConfig config,
std::unique_ptr<QuicTransport> transport, std::unique_ptr<QuicTransport> transport,
media::MediaSourceManager& media_hub, media::MediaSourceHub& media_hub,
DeviceSnapshotProvider device_snapshot_provider = {}); DeviceSnapshotProvider device_snapshot_provider = {});
~QuicEdgeService(); ~QuicEdgeService();
@ -106,7 +106,7 @@ private:
struct ActiveTrack { struct ActiveTrack {
config::QuicEdgeTrackConfig config; config::QuicEdgeTrackConfig config;
std::string source_track_id; std::string source_track_id;
media::MediaSourceManager::Subscription subscription; media::MediaSourceHub::Subscription subscription;
std::optional<MediaTrackDescription> last_description; std::optional<MediaTrackDescription> last_description;
bool waiting_for_keyframe{false}; bool waiting_for_keyframe{false};
bool keyframe_requested{false}; bool keyframe_requested{false};
@ -157,7 +157,7 @@ private:
config::QuicEdgeConfig config_; config::QuicEdgeConfig config_;
std::unique_ptr<QuicTransport> transport_; std::unique_ptr<QuicTransport> transport_;
media::MediaSourceManager* media_hub_{nullptr}; media::MediaSourceHub* media_hub_{nullptr};
bool using_global_media_hub_{false}; bool using_global_media_hub_{false};
SourceRegistrar source_registrar_; SourceRegistrar source_registrar_;
DeviceSnapshotProvider device_snapshot_provider_; 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 // 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. // announced on the reliable control stream.
struct DatagramHeader { struct DatagramHeader {
std::uint8_t protocol_version{kProtocolVersion}; std::uint8_t protocol_version{kProtocolVersion};

View File

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

View File

@ -25,7 +25,7 @@
#include "cmvr/quic_edge/v1/quic_edge.pb.h" #include "cmvr/quic_edge/v1/quic_edge.pb.h"
#include "common/base/logging/logger.h" #include "common/base/logging/logger.h"
#include "manager/device_manager/include/device_manager.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 cmvr::quic_edge {
namespace { namespace {
@ -423,7 +423,7 @@ const char* toString(const QuicEdgeServiceState state)
QuicEdgeService::QuicEdgeService(config::QuicEdgeConfig config, QuicEdgeService::QuicEdgeService(config::QuicEdgeConfig config,
std::unique_ptr<QuicTransport> transport, std::unique_ptr<QuicTransport> transport,
media::MediaSourceManager& media_hub, media::MediaSourceHub& media_hub,
DeviceSnapshotProvider device_snapshot_provider) DeviceSnapshotProvider device_snapshot_provider)
: config_(std::move(config)), : config_(std::move(config)),
transport_(std::move(transport)), transport_(std::move(transport)),
@ -598,7 +598,7 @@ bool QuicEdgeService::initialize(std::string* error)
return false; return false;
} }
if (!transport_ || !media_hub_) { 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); setState(QuicEdgeServiceState::FAILED, message);
setError(error, message); setError(error, message);
return false; return false;
@ -1450,21 +1450,9 @@ void QuicEdgeService::refreshMediaTracks(
recordMediaError(source_error); recordMediaError(source_error);
continue; 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.subscription = media_hub_->subscribe(
track.source_track_id, track.source_track_id,
media::MediaSourceManager::StartPosition::LATEST_AVAILABLE, media::MediaSourceHub::StartPosition::LATEST_AVAILABLE,
[this, activity_generation] { [this, activity_generation] {
std::lock_guard lock(mutex_); std::lock_guard lock(mutex_);
return stop_requested_ || media_stop_requested_ || return stop_requested_ || media_stop_requested_ ||
@ -1472,7 +1460,7 @@ void QuicEdgeService::refreshMediaTracks(
}); });
if (!track.subscription.valid()) { if (!track.subscription.valid()) {
recordMediaError( recordMediaError(
"MediaSourceManager source unavailable: " + track.source_track_id); "MediaSourceHub source unavailable: " + track.source_track_id);
continue; continue;
} }
track.waiting_for_keyframe = track.waiting_for_keyframe =
@ -1500,14 +1488,14 @@ bool QuicEdgeService::ensureSourceRegistered(
{ {
if (media_hub_->hasSource(source_track_id)) return true; if (media_hub_->hasSource(source_track_id)) return true;
if (!using_global_media_hub_ || !source_registrar_) { 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; return false;
} }
const bool registered = source_registrar_(track, source_track_id, error); const bool registered = source_registrar_(track, source_track_id, error);
if (!registered && !media_hub_->hasSource(source_track_id)) { if (!registered && !media_hub_->hasSource(source_track_id)) {
if (!error || error->empty()) { if (!error || error->empty()) {
setError(error, setError(error,
"failed to register MediaSourceManager source: " + source_track_id); "failed to register MediaSourceHub source: " + source_track_id);
} }
return false; return false;
} }
@ -1638,7 +1626,7 @@ bool QuicEdgeService::processTrack(ActiveTrack* track,
if (!read) { if (!read) {
if (!track->subscription.valid()) { if (!track->subscription.valid()) {
recordMediaError( recordMediaError(
"MediaSourceManager subscription stopped: " + track->source_track_id); "MediaSourceHub subscription stopped: " + track->source_track_id);
} }
return true; return true;
} }
@ -1647,7 +1635,7 @@ bool QuicEdgeService::processTrack(ActiveTrack* track,
if (!frame || !frame->descriptor || if (!frame || !frame->descriptor ||
frame->descriptor->id != track->source_track_id || frame->descriptor->id != track->source_track_id ||
frame->descriptor->kind != expectedKind(track->config)) { 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; track->next_frame_discontinuous = true;
return true; return true;
} }
@ -1667,7 +1655,7 @@ bool QuicEdgeService::processTrack(ActiveTrack* track,
!track->last_description || *track->last_description != description; !track->last_description || *track->last_description != description;
if (track->last_description && descriptor_changed && if (track->last_description && descriptor_changed &&
description.codec_generation < track->last_description->codec_generation) { description.codec_generation < track->last_description->codec_generation) {
recordMediaError("MediaSourceManager descriptor generation regressed"); recordMediaError("MediaSourceHub descriptor generation regressed");
track->next_frame_discontinuous = true; track->next_frame_discontinuous = true;
return true; return true;
} }

View File

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

View File

@ -9,7 +9,6 @@ target_link_libraries(quic_edge_task
cmvr_es::proto cmvr_es::proto
cmvr_es::logging cmvr_es::logging
cmvr_es::device_manager cmvr_es::device_manager
cmvr_es::stop_all_admission_gate
) )
add_library(cmvr_es::quic_edge_task ALIAS quic_edge_task) 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_compile_features(quic_edge_task_test PRIVATE cxx_std_17)
target_link_libraries(quic_edge_task_test PRIVATE target_link_libraries(quic_edge_task_test PRIVATE
cmvr_es::quic_edge_task cmvr_es::quic_edge_task
cmvr_es::stop_all_admission_gate
) )
add_test(NAME quic_edge_task_test COMMAND quic_edge_task_test) add_test(NAME quic_edge_task_test COMMAND quic_edge_task_test)
if(UNIX AND NOT APPLE) if(UNIX AND NOT APPLE)

View File

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

View File

@ -12,9 +12,8 @@
#include <vector> #include <vector>
#include "cmvr/quic_edge/v1/quic_edge.pb.h" #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/control_framing.h"
#include "service/grpc/stop_all/include/stop_all_admission_gate.h"
#include "task/quic_edge_task/include/quic_edge_task.h" #include "task/quic_edge_task/include/quic_edge_task.h"
namespace { namespace {
@ -235,9 +234,6 @@ cmvr::config::QuicEdgeConfig multiPlatformConfig()
int main() int main()
{ {
auto& admission = cmvr::service::globalStopAllAdmissionGate();
admission.clearForTesting();
cmvr::config::QuicEdgeConfig config; cmvr::config::QuicEdgeConfig config;
config.set_id("quic-invalid-config-test"); config.set_id("quic-invalid-config-test");
@ -256,50 +252,7 @@ int main()
return 1; return 1;
} }
cmvr::task::QuicEdgeTask admission_task(validConfig()); cmvr::media::MediaSourceHub media_hub;
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<HeartbeatTransport*> transports;
std::vector<QuicEdgeService*> services; std::vector<QuicEdgeService*> services;
std::vector<cmvr::config::QuicEdgeConfig> expanded_configs; std::vector<cmvr::config::QuicEdgeConfig> expanded_configs;