fix(quic): restore device inventory heartbeat
This commit is contained in:
parent
7ecf0324f1
commit
d7c2d439a4
@ -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
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@ -1,7 +1,9 @@
|
||||
#ifndef CMVR_ES_DEVICE_TYPES_H
|
||||
#define CMVR_ES_DEVICE_TYPES_H
|
||||
|
||||
#include <cstdint>
|
||||
#include <string>
|
||||
#include <vector>
|
||||
|
||||
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<ManagedDeviceSnapshot> devices;
|
||||
};
|
||||
|
||||
} // namespace cmvr::device
|
||||
|
||||
#endif // CMVR_ES_DEVICE_TYPES_H
|
||||
|
||||
@ -30,6 +30,7 @@ namespace cmvr::device {
|
||||
void getDeviceList(std::list<std::pair<std::string, std::string>> &device_list);
|
||||
void registerDevice(const std::shared_ptr<AbstractDevice>& device);
|
||||
void registerDevice(const std::string& device_id, const std::shared_ptr<AbstractDevice>& device);
|
||||
DeviceManagerSnapshot snapshot() const;
|
||||
|
||||
std::string version() const;
|
||||
std::string name() const;
|
||||
|
||||
@ -5,6 +5,8 @@
|
||||
|
||||
#include "../include/device_manager.h"
|
||||
|
||||
#include <algorithm>
|
||||
|
||||
#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();
|
||||
}
|
||||
|
||||
@ -70,10 +70,14 @@ struct QuicEdgeStatus {
|
||||
|
||||
class QuicEdgeService {
|
||||
public:
|
||||
using DeviceSnapshotProvider =
|
||||
std::function<device::DeviceManagerSnapshot()>;
|
||||
|
||||
explicit QuicEdgeService(config::QuicEdgeConfig config);
|
||||
QuicEdgeService(config::QuicEdgeConfig config,
|
||||
std::unique_ptr<QuicTransport> 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_;
|
||||
|
||||
@ -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](
|
||||
|
||||
@ -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<unsigned char>(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<const device::ManagedDeviceSnapshot*> 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<QuicTransport> 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::DeviceManagerSnapshot> 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");
|
||||
|
||||
Loading…
Reference in New Issue
Block a user