From d7c2d439a46884bdc81eab8c5ec8c04232e39831 Mon Sep 17 00:00:00 2001 From: linbo <1034003879@qq.com> Date: Sun, 20 Sep 2026 15:03:35 +0800 Subject: [PATCH] fix(quic): restore device inventory heartbeat --- cmvr-es/config/manager/device_manager.pb.txt | 45 +++++- cmvr-es/devices/device_types.h | 83 ++++++++++ .../device_manager/include/device_manager.h | 1 + .../device_manager/src/device_manager.cpp | 31 ++++ .../quic_edge/include/quic_edge_service.h | 7 +- .../src/quic_edge_device_adapter.cpp | 5 +- .../quic_edge/src/quic_edge_service.cpp | 152 +++++++++++++++++- 7 files changed, 319 insertions(+), 5 deletions(-) diff --git a/cmvr-es/config/manager/device_manager.pb.txt b/cmvr-es/config/manager/device_manager.pb.txt index c387cc85..1b8c3f6d 100644 --- a/cmvr-es/config/manager/device_manager.pb.txt +++ b/cmvr-es/config/manager/device_manager.pb.txt @@ -144,9 +144,52 @@ device_manager { } devices { - id: "agv_1" + id: "src1100" type: DEVICE_TYPE_AGV config_file: "devices/agv/agv.pb.txt" enable: false } + + devices { + id: "hikvision_cam" + type: DEVICE_TYPE_CAMERA + config_file: "devices/camera/camera.pb.txt" + enable: false + } + + devices { + id: "hikvision_thermal_cam" + type: DEVICE_TYPE_CAMERA + config_file: "devices/camera/camera.pb.txt" + enable: false + } + + devices { + id: "mic1" + type: DEVICE_TYPE_MICROPHONE + config_file: "devices/microphone/microphone.pb.txt" + enable: false + } + + devices { + id: "spk1" + type: DEVICE_TYPE_SPEAKER + config_file: "devices/speaker/speaker.pb.txt" + enable: false + } + + devices { + id: "real_cam1" + type: DEVICE_TYPE_CAMERA + config_file: "devices/camera/camera.pb.txt" + enable: false + } + + devices { + id: "usb_cam1" + type: DEVICE_TYPE_CAMERA + config_file: "devices/camera/camera.pb.txt" + enable: false + } + } diff --git a/cmvr-es/devices/device_types.h b/cmvr-es/devices/device_types.h index 38b8d5b3..ddf55a62 100644 --- a/cmvr-es/devices/device_types.h +++ b/cmvr-es/devices/device_types.h @@ -1,7 +1,9 @@ #ifndef CMVR_ES_DEVICE_TYPES_H #define CMVR_ES_DEVICE_TYPES_H +#include #include +#include namespace cmvr::device { @@ -68,6 +70,87 @@ namespace cmvr::device { std::string type_name; }; + // DeviceManager lifecycle and device-reported health are deliberately + // separate. A device can, for example, be READY from the manager's point + // of view while its backend has not implemented health reporting yet. + enum class ManagedDeviceState { + Unknown, + Disabled, + Initializing, + Registered, + Ready, + Running, + Stopped, + Error, + }; + + enum class DeviceHealthState { + Unknown, + Healthy, + Degraded, + Fault, + }; + + inline std::string toString(ManagedDeviceState state) { + switch (state) { + case ManagedDeviceState::Disabled: + return "Disabled"; + case ManagedDeviceState::Initializing: + return "Initializing"; + case ManagedDeviceState::Registered: + return "Registered"; + case ManagedDeviceState::Ready: + return "Ready"; + case ManagedDeviceState::Running: + return "Running"; + case ManagedDeviceState::Stopped: + return "Stopped"; + case ManagedDeviceState::Error: + return "Error"; + case ManagedDeviceState::Unknown: + default: + return "Unknown"; + } + } + + inline std::string toString(DeviceHealthState state) { + switch (state) { + case DeviceHealthState::Healthy: + return "Healthy"; + case DeviceHealthState::Degraded: + return "Degraded"; + case DeviceHealthState::Fault: + return "Fault"; + case DeviceHealthState::Unknown: + default: + return "Unknown"; + } + } + + struct DeviceHealthSnapshot { + DeviceHealthState state = DeviceHealthState::Unknown; + std::string error_message; + }; + + struct ManagedDeviceSnapshot { + std::string id; + DeviceKind kind = DeviceKind::Unknown; + std::string type_name; + bool enabled = false; + ManagedDeviceState state = ManagedDeviceState::Unknown; + DeviceHealthSnapshot health; + bool abnormal = false; + std::string error_message; + std::uint64_t status_updated_at_unix_ms = 0; + }; + + struct DeviceManagerSnapshot { + std::string name; + std::string version; + std::string description; + std::vector devices; + }; + } // namespace cmvr::device #endif // CMVR_ES_DEVICE_TYPES_H diff --git a/cmvr-es/manager/device_manager/include/device_manager.h b/cmvr-es/manager/device_manager/include/device_manager.h index 58bedc74..afe408a6 100644 --- a/cmvr-es/manager/device_manager/include/device_manager.h +++ b/cmvr-es/manager/device_manager/include/device_manager.h @@ -30,6 +30,7 @@ namespace cmvr::device { void getDeviceList(std::list> &device_list); void registerDevice(const std::shared_ptr& device); void registerDevice(const std::string& device_id, const std::shared_ptr& device); + DeviceManagerSnapshot snapshot() const; std::string version() const; std::string name() const; diff --git a/cmvr-es/manager/device_manager/src/device_manager.cpp b/cmvr-es/manager/device_manager/src/device_manager.cpp index d8d8ac0c..a2fb4952 100644 --- a/cmvr-es/manager/device_manager/src/device_manager.cpp +++ b/cmvr-es/manager/device_manager/src/device_manager.cpp @@ -5,6 +5,8 @@ #include "../include/device_manager.h" +#include + #include "devices/agv/abstract_agv.h" #include "devices/arm/robot_arm.h" #include "devices/battery/abstract_battery.h" @@ -202,6 +204,35 @@ void DeviceManager::registerDevice(const std::string& device_id, << ", kind=" << toString(device->kind()); } +DeviceManagerSnapshot DeviceManager::snapshot() const +{ + DeviceManagerSnapshot result; + result.name = name(); + result.version = version(); + result.description = description(); + result.devices.reserve(devices_.size()); + + for (const auto& [id, record] : devices_) { + ManagedDeviceSnapshot status; + status.id = id; + status.kind = record.kind; + status.type_name = record.type_name; + status.enabled = true; + status.state = ManagedDeviceState::Ready; + // Current dev has no AbstractDevice::healthSnapshot(). Keep the + // linbo_dev wire contract without inventing health information. + status.health.state = DeviceHealthState::Unknown; + result.devices.push_back(std::move(status)); + } + + std::sort(result.devices.begin(), result.devices.end(), + [](const ManagedDeviceSnapshot& lhs, + const ManagedDeviceSnapshot& rhs) { + return lhs.id < rhs.id; + }); + return result; +} + std::string DeviceManager::version() const { return cfg_.version().empty() ? "1.0" : cfg_.version(); } 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 68375a4c..23eae61b 100644 --- a/cmvr-es/service/quic_edge/include/quic_edge_service.h +++ b/cmvr-es/service/quic_edge/include/quic_edge_service.h @@ -70,10 +70,14 @@ struct QuicEdgeStatus { class QuicEdgeService { public: + using DeviceSnapshotProvider = + std::function; + explicit QuicEdgeService(config::QuicEdgeConfig config); QuicEdgeService(config::QuicEdgeConfig config, std::unique_ptr transport, - media::MediaSourceHub& media_hub); + media::MediaSourceHub& media_hub, + DeviceSnapshotProvider device_snapshot_provider = {}); ~QuicEdgeService(); QuicEdgeService(const QuicEdgeService&) = delete; @@ -156,6 +160,7 @@ private: media::MediaSourceHub* media_hub_{nullptr}; bool using_global_media_hub_{false}; SourceRegistrar source_registrar_; + DeviceSnapshotProvider device_snapshot_provider_; std::mutex lifecycle_mutex_; std::mutex control_send_mutex_; 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 b9819542..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 @@ -11,7 +11,10 @@ QuicEdgeService::QuicEdgeService(config::QuicEdgeConfig config) : config_(std::move(config)), transport_(createDefaultQuicTransport(config_.datagram_send_queue_depth())), media_hub_(&media::globalMediaSourceHub()), - using_global_media_hub_(true) + using_global_media_hub_(true), + device_snapshot_provider_([] { + return device::DeviceManager::getInstance().snapshot(); + }) { initializeIdentity(); source_registrar_ = [this]( 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 47d1df83..30a2f14d 100644 --- a/cmvr-es/service/quic_edge/src/quic_edge_service.cpp +++ b/cmvr-es/service/quic_edge/src/quic_edge_service.cpp @@ -36,6 +36,7 @@ constexpr std::uint32_t kMaximumDatagramBytes = 65527U; constexpr std::uint32_t kMaximumControlBytes = 16U * 1024U * 1024U; constexpr std::uint32_t kMaximumConfiguredFrameBytes = 256U * 1024U * 1024U; constexpr std::size_t kMaximumDrainPerPoll = 4096U; +constexpr std::size_t kMaximumDeviceErrorBytes = 512U; constexpr std::uint32_t kMinimumHeartbeatIntervalMs = 250U; constexpr std::uint32_t kMaximumHeartbeatIntervalMs = 60U * 60U * 1000U; constexpr auto kControlPollInterval = std::chrono::milliseconds(50); @@ -49,6 +50,18 @@ void setError(std::string* error, const std::string& message) } } +std::string boundedDeviceError(std::string message) +{ + if (message.size() <= kMaximumDeviceErrorBytes) return message; + std::size_t end = kMaximumDeviceErrorBytes; + while (end > 0U && + (static_cast(message[end]) & 0xc0U) == 0x80U) { + --end; + } + message.resize(end); + return message; +} + std::string trim(std::string value) { const auto first = value.find_first_not_of(" \t\r\n"); @@ -235,6 +248,116 @@ void populateHeartbeatNetwork(const config::QuicEdgeConfig& config, endpoint->set_tls(config.grpc_endpoint_tls()); } +v1::DeviceKind toProtoDeviceKind(const device::DeviceKind kind) +{ + switch (kind) { + case device::DeviceKind::AGV: return v1::DEVICE_KIND_AGV; + case device::DeviceKind::Arm: return v1::DEVICE_KIND_ARM; + case device::DeviceKind::Battery: return v1::DEVICE_KIND_BATTERY; + case device::DeviceKind::BioHead: return v1::DEVICE_KIND_BIO_HEAD; + case device::DeviceKind::Camera: return v1::DEVICE_KIND_CAMERA; + case device::DeviceKind::CanBus: return v1::DEVICE_KIND_CAN_BUS; + case device::DeviceKind::DexHand: return v1::DEVICE_KIND_DEX_HAND; + case device::DeviceKind::Gripper: return v1::DEVICE_KIND_GRIPPER; + case device::DeviceKind::Microphone: return v1::DEVICE_KIND_MICROPHONE; + case device::DeviceKind::Motor: return v1::DEVICE_KIND_MOTOR; + case device::DeviceKind::MotorSystem: + return v1::DEVICE_KIND_MOTOR_SYSTEM; + case device::DeviceKind::Robot: return v1::DEVICE_KIND_ROBOT; + case device::DeviceKind::Speaker: return v1::DEVICE_KIND_SPEAKER; + case device::DeviceKind::Unknown: break; + } + return v1::DEVICE_KIND_UNSPECIFIED; +} + +v1::ManagedDeviceState toProtoManagedDeviceState( + const device::ManagedDeviceState state) +{ + switch (state) { + case device::ManagedDeviceState::Disabled: + return v1::MANAGED_DEVICE_STATE_DISABLED; + case device::ManagedDeviceState::Initializing: + return v1::MANAGED_DEVICE_STATE_INITIALIZING; + case device::ManagedDeviceState::Registered: + return v1::MANAGED_DEVICE_STATE_REGISTERED; + case device::ManagedDeviceState::Ready: + return v1::MANAGED_DEVICE_STATE_READY; + case device::ManagedDeviceState::Running: + return v1::MANAGED_DEVICE_STATE_RUNNING; + case device::ManagedDeviceState::Stopped: + return v1::MANAGED_DEVICE_STATE_STOPPED; + case device::ManagedDeviceState::Error: + return v1::MANAGED_DEVICE_STATE_ERROR; + case device::ManagedDeviceState::Unknown: + break; + } + return v1::MANAGED_DEVICE_STATE_UNSPECIFIED; +} + +v1::DeviceHealthStatus toProtoDeviceHealthStatus( + const device::DeviceHealthState state) +{ + switch (state) { + case device::DeviceHealthState::Healthy: + return v1::DEVICE_HEALTH_STATUS_HEALTHY; + case device::DeviceHealthState::Degraded: + return v1::DEVICE_HEALTH_STATUS_DEGRADED; + case device::DeviceHealthState::Fault: + return v1::DEVICE_HEALTH_STATUS_FAULT; + case device::DeviceHealthState::Unknown: + break; + } + return v1::DEVICE_HEALTH_STATUS_UNSPECIFIED; +} + +void populateDeviceManagerSnapshot( + const device::DeviceManagerSnapshot& source, + const std::uint64_t sampled_at_unix_ms, + v1::DeviceManagerSnapshot* destination) +{ + if (!destination) return; + destination->set_manager_name(source.name); + destination->set_manager_version(source.version); + destination->set_manager_description(source.description); + destination->set_sampled_at_unix_ms(sampled_at_unix_ms); + std::vector ordered_devices; + ordered_devices.reserve(source.devices.size()); + for (const auto& source_device : source.devices) { + // DeviceManager keeps disabled entries for local configuration and + // diagnostics, but the platform heartbeat only advertises devices + // that are enabled on this edge node. Enabled entries remain visible + // even when creation, initialization or start has failed. + if (!source_device.enabled) continue; + ordered_devices.push_back(&source_device); + } + std::sort( + ordered_devices.begin(), ordered_devices.end(), + [](const auto* lhs, const auto* rhs) { + if (lhs->id != rhs->id) return lhs->id < rhs->id; + if (lhs->kind != rhs->kind) return lhs->kind < rhs->kind; + return lhs->type_name < rhs->type_name; + }); + + for (const auto* source_device : ordered_devices) { + auto* destination_device = destination->add_devices(); + destination_device->set_device_id(source_device->id); + destination_device->set_kind(toProtoDeviceKind(source_device->kind)); + destination_device->set_type_name(source_device->type_name); + destination_device->set_enabled(source_device->enabled); + destination_device->set_manager_state( + toProtoManagedDeviceState(source_device->state)); + destination_device->set_health( + toProtoDeviceHealthStatus(source_device->health.state)); + destination_device->set_has_error(source_device->abnormal); + destination_device->set_error_message(boundedDeviceError( + source_device->error_message.empty() + ? source_device->health.error_message + : source_device->error_message)); + destination_device->set_status_updated_at_unix_ms( + source_device->status_updated_at_unix_ms); + } +} + v1::MediaKind toProtoKind(const media::MediaKind kind) { switch (kind) { @@ -300,10 +423,12 @@ const char* toString(const QuicEdgeServiceState state) QuicEdgeService::QuicEdgeService(config::QuicEdgeConfig config, std::unique_ptr transport, - media::MediaSourceHub& media_hub) + media::MediaSourceHub& media_hub, + DeviceSnapshotProvider device_snapshot_provider) : config_(std::move(config)), transport_(std::move(transport)), - media_hub_(&media_hub) + media_hub_(&media_hub), + device_snapshot_provider_(std::move(device_snapshot_provider)) { initializeIdentity(); } @@ -1062,6 +1187,24 @@ bool QuicEdgeService::sendNodeRegistration(std::string* error) bool QuicEdgeService::sendHeartbeat(const std::uint64_t sequence, std::string* error) { + std::optional device_manager_snapshot; + if (device_snapshot_provider_) { + try { + device_manager_snapshot = device_snapshot_provider_(); + } catch (const std::exception& exception) { + setError( + error, + std::string("failed to snapshot DeviceManager for heartbeat: ") + + exception.what()); + return false; + } catch (...) { + setError(error, + "failed to snapshot DeviceManager for heartbeat: " + "unknown exception"); + return false; + } + } + std::lock_guard control_lock(control_send_mutex_); std::string session_id; { @@ -1086,6 +1229,11 @@ bool QuicEdgeService::sendHeartbeat(const std::uint64_t sequence, heartbeat->set_sent_at_unix_ms(sent_at_unix_ms); heartbeat->set_software_version(software_version_); populateHeartbeatNetwork(config_, heartbeat); + if (device_manager_snapshot.has_value()) { + populateDeviceManagerSnapshot( + *device_manager_snapshot, sent_at_unix_ms, + heartbeat->mutable_device_manager()); + } std::string serialized; if (!envelope.SerializeToString(&serialized)) { setError(error, "failed to serialize NodeHeartbeat");