diff --git a/.gitignore b/.gitignore index a547bf3..5f34c60 100644 --- a/.gitignore +++ b/.gitignore @@ -22,3 +22,10 @@ dist-ssr *.njsproj *.sln *.sw? + +# Python runtime artifacts and local TLS private material +__pycache__/ +*.py[cod] +server/certs/*.key +server/certs/*.csr +server/certs/*.srl diff --git a/components.d.ts b/components.d.ts index 650b5d0..c1980dc 100644 --- a/components.d.ts +++ b/components.d.ts @@ -24,6 +24,8 @@ declare module 'vue' { ElRadioGroup: typeof import('element-plus/es')['ElRadioGroup'] ElSelect: typeof import('element-plus/es')['ElSelect'] ElSlider: typeof import('element-plus/es')['ElSlider'] + ElTable: typeof import('element-plus/es')['ElTable'] + ElTableColumn: typeof import('element-plus/es')['ElTableColumn'] ElTag: typeof import('element-plus/es')['ElTag'] IPlayer: typeof import('./src/components/IPlayer/index.vue')['default'] LandscapeGuard: typeof import('./src/components/LandscapeGuard.vue')['default'] diff --git a/package.json b/package.json index 5002f4d..d731774 100644 --- a/package.json +++ b/package.json @@ -7,6 +7,7 @@ "dev": "concurrently \"npm run dev:server\" \"npm run dev:client\"", "dev:client": "vite", "dev:server": "node server/index.js", + "quic:server": "py -3 server/quic/server.py", "start": "cross-env NODE_ENV=production node server/index.js", "build": "vite build", "preview": "npm run build && cross-env NODE_ENV=production node server/index.js" diff --git a/server/certs/cmvr-quic-ca.crt b/server/certs/cmvr-quic-ca.crt new file mode 100644 index 0000000..91e2619 --- /dev/null +++ b/server/certs/cmvr-quic-ca.crt @@ -0,0 +1,25 @@ +-----BEGIN CERTIFICATE----- +MIIEKzCCApOgAwIBAgIUTtZCyKM8INYxZKkfwbKdpHKLjn4wDQYJKoZIhvcNAQEL +BQAwHTEbMBkGA1UEAwwSQ01WUiBRVUlDIExvY2FsIENBMB4XDTI2MDcyNDA2NDQy +MFoXDTM2MDcyMTA2NDQyMFowHTEbMBkGA1UEAwwSQ01WUiBRVUlDIExvY2FsIENB +MIIBojANBgkqhkiG9w0BAQEFAAOCAY8AMIIBigKCAYEAmq2rHldOobaemqNfWggS +OVj3inKy6AYjfgtcXUfKs48DDbpZ9gyEd/YPJXA8C2jGPXpxzgnc7a4UCUVQZ8ah +ddoJtFcC+Q6BgjeMVqUdUubu5Y9HpkfU3lvnp4KhzvOeFnkKtrCzYIPa2nK3zLc7 +uCiuLlB+91KQSRXPFbc6N7H/EAfGmUHIwlZGysAkRN7b2TAoR4C7E96JLVtuUQsS +VtlEGpunSfuefFzeeZCMS6avLbB+a8Q6yUzLt6pqnheNsDB+jCCXodlJs5XS1AOB +W2GOpGFMj7dLoTD+eBIlAlrhWFcwKjzFmtp6LGl/Jy0O+E99X4TL72oNSOdTrO2O +iavXABz9IvWR2BrAyo5AKlTJqO6tmZw77iVti8jYi+HsXIVGQKMYwvv5k0jdHeIZ +FaDToUbFPP/zj0m8ZraMy+8eNAhScnx6Zs56fcncBuDti6pT+zKisjV1rH/sFvZY +wO1UvJlOEbvXrfYLPp58Aqe/toG2nV45a2Q4+x6AETV1AgMBAAGjYzBhMB0GA1Ud +DgQWBBRdHz+g2vVmqghDZi2kEHtZfoyHZzAfBgNVHSMEGDAWgBRdHz+g2vVmqghD +Zi2kEHtZfoyHZzAPBgNVHRMBAf8EBTADAQH/MA4GA1UdDwEB/wQEAwIBBjANBgkq +hkiG9w0BAQsFAAOCAYEASuEoMFiVcg4iTMxO2kshFTJ6LIrqqGXBn+1j+yQ1ennG +mqPMo5fBOe/Kp3YWCnREQWu0+EEPEC9qWgIDOIm3v7ch4mZiW31GUOqae6bjprBe +er7ySElKGZ5GefKAq++we19A6WHnNxtNAT9BE1VSKUmkxEsnIkuwd+QMmQ9eaIRm +8RvzshWdUcyiBJg07sI3rPzPD/YpnfcYlAa2a0+oXJ3o3zlhdbs2S+9fIl+Y/zFN +wBaUV6ZNk0RFzOCA+jUq2jU5Y1pODcot3Mp2jCHzp4uTFyVqsoqc5USRB2TUILCP +acSbEjD9GnKaNU4miTcftnC85fnBAiZ5btZkXIs72ij+J3VW7QnbcFQG8bElngoR +mWizh6ByPQrsvyT0evJjdjpuOzILsVTAn64m0cYpnlR4g5orcEkIJHv4DzMhpXc2 +CvpcWwk2vJrUuws/LGxDATTAcqtIUIxAkxffqUsdB/6j8fdykNCvXFwNh/uPcawd +4kzh6dr+RFzHSMQ6oN16 +-----END CERTIFICATE----- diff --git a/server/certs/quic-gateway.crt b/server/certs/quic-gateway.crt new file mode 100644 index 0000000..dbdfded --- /dev/null +++ b/server/certs/quic-gateway.crt @@ -0,0 +1,23 @@ +-----BEGIN CERTIFICATE----- +MIIDyzCCAjOgAwIBAgIUGtonFlB5GKMtTb45iJltgKRUrWkwDQYJKoZIhvcNAQEL +BQAwHTEbMBkGA1UEAwwSQ01WUiBRVUlDIExvY2FsIENBMB4XDTI2MDgxNDA3NTMy +NFoXDTI4MTExNjA3NTgyNFowGDEWMBQGA1UEAwwNMTkyLjE2OC4wLjE5OTCCASIw +DQYJKoZIhvcNAQEBBQADggEPADCCAQoCggEBAKdM3i1FYFKqNWJzhfhsD9nRUAuK +pzilz5uqCKAt8lKYYC9WnLHOYdiEjcHGnGr02yd6sWFH/LBbxNhzx8M7h4S4izuO +bhSlG1EIhkMiojzVD1e3P7YzXdEoVxTCfmMgBZQJG63GNOfzRawFYtEeGv7ndFVw +kitCYlyTza5KlBNlWpiNOPmmx4dLTdGLUk8a5TUm0zJ+b/LyzVsWUtr9sxKmKeG8 +0/77AeiL0hQE3xUt5QROTZjRTVhNHowv410dFMJyIfY4sab8ndc4SIwE9PCKg068 +4705vbFInBS3eTvQur5VLSZLPatnXGKzCXjci1lIQX2p/QIqADMGFWIN9OECAwEA +AaOBhzCBhDAPBgNVHREECDAGhwTAqADHMAwGA1UdEwEB/wQCMAAwDgYDVR0PAQH/ +BAQDAgeAMBMGA1UdJQQMMAoGCCsGAQUFBwMBMB0GA1UdDgQWBBS5qBLbIJcZt5l4 +eJfDhd9W7dmYZzAfBgNVHSMEGDAWgBRdHz+g2vVmqghDZi2kEHtZfoyHZzANBgkq +hkiG9w0BAQsFAAOCAYEAMGmRoddVWo5THkH39R40niqfX3pN3pCx9ywDG6rSm7pO +WYY8h91rWpecT40fxILVzdqYc6i7Ufinh1qQOnm/ExEUCPUEpMhn75W3lNhgUKmC +wGjYRVvMZrnLjhamiChWZx9kvYrURe9JHwwevwxS6V69TRaF0J0OOFaQcRqFE9iP +2/VdnmKk+OlbUdxBj3NarFWcqhMmxSoetWEQS5cTGzyrcZhnWPjgp8YNUhay7Lnx +HrYjiU8YJJNX163pBvKp37QaIcbs/w6o2IlzCSx5xb7qcMboFGP9mH41bnbuoki2 +leWSoDyXS348wZcfkCwd1mNsnPsUFnDTq6gdzS1gbsps309ZmilnTdqQJUY8dmEf +h4M6PoUikYLfKT+/kdqH8ywJ/Uycj3+I9xT8hiig+E9dywp+ExHMy8OzREMEXvt9 +KMLgNqGrqwtrU3OBBEnLEKrAJttNkuQZDUlR59HvZ5xERXzelnmBKwMvYQ0sut+y +6nYZUewPypQuEevdc5RR +-----END CERTIFICATE----- diff --git a/server/index.js b/server/index.js index e122b9f..b4014e2 100644 --- a/server/index.js +++ b/server/index.js @@ -16,12 +16,18 @@ import { } from './grpcClient.js'; import { attachAudioWebSocket } from './audioWebSocket.js'; import { streamCameraVideo } from './cameraVideoStream.js'; +import { startQuicServer } from './quicServer.js'; +import { consumeQuicEvent, listRobots } from './robotRegistry.js'; const app = express() app.use(express.static(path.join(__dirname, '../dist'))); app.use(cors()) app.use(express.json()) +app.get('/api/quic/robots', (req, res) => { + res.json({ code: 200, data: listRobots() }) +}) + // 将相机逐帧 gRPC 数据转码为浏览器可直接播放的 fragmented MP4。 // 示例:/api/camera/video?ip=192.168.0.28%3A50052&deviceId=camera1 app.get('/api/camera/video', streamCameraVideo) @@ -546,6 +552,7 @@ if (httpsEnabled) { } attachAudioWebSocket(httpServer) +startQuicServer({ onEvent: consumeQuicEvent }) httpServer.listen(PORT, '0.0.0.0', () => { const protocol = httpsEnabled ? 'https' : 'http'; diff --git a/server/proto/cmvr/quic_edge/v1/quic_edge.proto b/server/proto/cmvr/quic_edge/v1/quic_edge.proto new file mode 100644 index 0000000..48f3a32 --- /dev/null +++ b/server/proto/cmvr/quic_edge/v1/quic_edge.proto @@ -0,0 +1,232 @@ +syntax = "proto3"; + +package cmvr.quic_edge.v1; + +// Every application message on the edge-opened reliable bidirectional stream +// is carried by this envelope. message_sequence is strictly increasing per +// sender (gaps are allowed) and is independent from the heartbeat sequence used +// for liveness acknowledgement. +message EdgeControlEnvelope { + uint32 protocol_version = 1; + uint64 message_sequence = 2; + + oneof payload { + NodeRegisterRequest node_register_request = 10; + NodeRegisterResponse node_register_response = 11; + NodeHeartbeat node_heartbeat = 12; + NodeHeartbeatAck node_heartbeat_ack = 13; + + MediaSessionOpen media_session_open = 20; + MediaTrackDescriptor media_track_descriptor = 21; + MediaSessionClose media_session_close = 22; + + ProtocolError protocol_error = 30; + } +} + +message NetworkInterfaceAddress { + enum AddressFamily { + ADDRESS_FAMILY_UNSPECIFIED = 0; + ADDRESS_FAMILY_IPV4 = 1; + ADDRESS_FAMILY_IPV6 = 2; + } + + string interface_name = 1; + string ip_address = 2; + AddressFamily family = 3; + bool loopback = 4; +} + +// The existing cmvr-es gRPC server remains the robot-control endpoint. The +// edge advertises its current reachable address through the QUIC control plane. +message GrpcEndpoint { + string host = 1; + uint32 port = 2; + bool tls = 3; +} + +// Stable protocol-level categories for devices managed by cmvr-es. These +// values intentionally do not reuse the configuration or gRPC API enums: +// their zero values and supported categories have different semantics. +enum DeviceKind { + DEVICE_KIND_UNSPECIFIED = 0; + DEVICE_KIND_AGV = 1; + DEVICE_KIND_ARM = 2; + DEVICE_KIND_BATTERY = 3; + DEVICE_KIND_BIO_HEAD = 4; + DEVICE_KIND_CAMERA = 5; + DEVICE_KIND_CAN_BUS = 6; + DEVICE_KIND_DEX_HAND = 7; + DEVICE_KIND_GRIPPER = 8; + DEVICE_KIND_MICROPHONE = 9; + DEVICE_KIND_MOTOR = 10; + DEVICE_KIND_MOTOR_SYSTEM = 11; + DEVICE_KIND_ROBOT = 12; + DEVICE_KIND_SPEAKER = 13; +} + +// DeviceManager's view of a configured entry. REGISTERED means that the +// manager owns a device record but has no more specific lifecycle signal. +enum ManagedDeviceState { + MANAGED_DEVICE_STATE_UNSPECIFIED = 0; + MANAGED_DEVICE_STATE_DISABLED = 1; + MANAGED_DEVICE_STATE_INITIALIZING = 2; + MANAGED_DEVICE_STATE_REGISTERED = 3; + MANAGED_DEVICE_STATE_READY = 4; + MANAGED_DEVICE_STATE_RUNNING = 5; + MANAGED_DEVICE_STATE_STOPPED = 6; + MANAGED_DEVICE_STATE_ERROR = 7; +} + +// Health is independent of lifecycle. UNSPECIFIED means that no trustworthy +// health observation is available and must never be interpreted as healthy. +enum DeviceHealthStatus { + DEVICE_HEALTH_STATUS_UNSPECIFIED = 0; + DEVICE_HEALTH_STATUS_HEALTHY = 1; + DEVICE_HEALTH_STATUS_DEGRADED = 2; + DEVICE_HEALTH_STATUS_FAULT = 3; +} + +message ManagedDeviceStatus { + string device_id = 1; + DeviceKind kind = 2; + + // Concrete implementation name when a device object exists; otherwise a + // category label. It is for display/diagnostics only. Consumers use kind, + // rather than this free-form string, for machine decisions. + string type_name = 3; + + bool enabled = 4; + ManagedDeviceState manager_state = 5; + DeviceHealthStatus health = 6; + + // false means that no error is currently confirmed. It does not turn + // DEVICE_HEALTH_STATUS_UNSPECIFIED into a healthy observation. + bool has_error = 7; + string error_message = 8; + + // Time at which DeviceManager last changed the lifecycle/error record. + // DeviceManagerSnapshot.sampled_at_unix_ms is the freshness timestamp for + // the health observation carried by this heartbeat. + uint64 status_updated_at_unix_ms = 9; +} + +message DeviceManagerSnapshot { + string manager_name = 1; + string manager_version = 2; + string manager_description = 3; + + // Current cmvr-es senders include only enabled devices. The enabled field in + // each row and DISABLED enum value remain part of v1 for wire compatibility. + repeated ManagedDeviceStatus devices = 4; + uint64 sampled_at_unix_ms = 5; +} + +message NodeDescriptor { + string node_id = 1; + string boot_id = 2; + string software_version = 3; + repeated NetworkInterfaceAddress local_interfaces = 4; + GrpcEndpoint grpc_endpoint = 5; + string robot_id = 6; +} + +// This must be the first application message sent after each QUIC connection +// is established. A reconnect always creates a new registration session. +message NodeRegisterRequest { + NodeDescriptor node = 1; + uint64 sent_at_unix_ms = 2; +} + +message NodeRegisterResponse { + bool accepted = 1; + string session_id = 2; + string message = 3; + + // Zero tells the edge to keep its locally configured interval. + uint32 heartbeat_interval_ms = 4; + + // Derived by the receiver from the authenticated QUIC peer address. It is + // not copied from a client-supplied local interface. + string observed_source_ip = 5; +} + +// A heartbeat carries a fresh network snapshot so address changes are reported +// without opening a second protocol or connection. +message NodeHeartbeat { + string node_id = 1; + string boot_id = 2; + string session_id = 3; + uint64 sequence = 4; + uint64 sent_at_unix_ms = 5; + string software_version = 6; + repeated NetworkInterfaceAddress local_interfaces = 7; + GrpcEndpoint grpc_endpoint = 8; + DeviceManagerSnapshot device_manager = 9; + string robot_id = 10; +} + +message NodeHeartbeatAck { + bool accepted = 1; + uint64 acknowledged_sequence = 2; + string message = 3; + uint64 server_time_unix_ms = 4; + string observed_source_ip = 5; + string session_id = 6; +} + +message MediaSessionOpen { + string node_id = 1; + uint64 session_epoch = 2; + + // Binds the media epoch to the accepted node registration on this QUIC + // connection. Media must not start before this session is assigned. + string session_id = 3; +} + +enum MediaKind { + MEDIA_KIND_UNSPECIFIED = 0; + MEDIA_KIND_VIDEO = 1; + MEDIA_KIND_AUDIO = 2; +} + +message MediaTrackDescriptor { + uint32 track_id = 1; + MediaKind kind = 2; + string device_id = 3; + string codec = 4; + uint64 codec_generation = 5; + + // The exact MediaSourceHub track and the 32-bit token repeated in every + // DATAGRAM header. The full generation remains on the reliable stream. + string source_track_id = 6; + uint32 codec_generation_token = 7; + string payload_format = 8; + + // Video fields. They are zero for audio tracks. + uint32 width = 10; + uint32 height = 11; + uint32 frames_per_second = 12; + + // Audio fields. They are zero for video tracks. + uint32 sample_rate = 20; + uint32 channels = 21; + + // Decoder initialization bytes, for example AVCC/HVCC or AudioSpecificConfig. + // Existing cmvr-es sources may leave this empty when configuration NAL units + // are carried in-band. + bytes codec_config = 30; +} + +message MediaSessionClose { + string reason = 1; + string session_id = 2; + uint64 session_epoch = 3; +} + +message ProtocolError { + uint32 code = 1; + string message = 2; + uint64 related_message_sequence = 3; + bool fatal = 4; +} diff --git a/server/proto/cmvr/quic_gateway/v1/quic_gateway.proto b/server/proto/cmvr/quic_gateway/v1/quic_gateway.proto new file mode 100644 index 0000000..ef97aa9 --- /dev/null +++ b/server/proto/cmvr/quic_gateway/v1/quic_gateway.proto @@ -0,0 +1,104 @@ +syntax = "proto3"; + +package cmvr.quic_gateway.v1; + +option java_package = "com.cmvr.quic.gateway.v1"; +option java_multiple_files = true; +option java_outer_classname = "QuicGatewayProtocol"; + +enum NodeEventType { + NODE_EVENT_TYPE_UNSPECIFIED = 0; + NODE_EVENT_TYPE_REGISTERED = 1; + NODE_EVENT_TYPE_HEARTBEAT = 2; + NODE_EVENT_TYPE_OFFLINE = 3; + NODE_EVENT_TYPE_MEDIA_CHANGED = 4; +} + +message NodeEvent { + NodeEventType type = 1; + NodeSnapshot node = 2; + uint64 event_time_unix_ms = 3; +} + +message NodeSnapshot { + string node_id = 1; + string boot_id = 2; + string session_id = 3; + string software_version = 4; + bool online = 5; + uint64 registered_at_unix_ms = 6; + uint64 last_heartbeat_at_unix_ms = 7; + uint64 heartbeat_sequence = 8; + string observed_source_ip = 9; + GrpcEndpoint grpc_endpoint = 10; + repeated NetworkInterfaceAddress local_interfaces = 11; + DeviceManagerSnapshot device_manager = 12; + uint64 media_session_epoch = 13; + repeated MediaTrack tracks = 14; + string robot_id = 15; +} + +message GrpcEndpoint { + string host = 1; + uint32 port = 2; + bool tls = 3; +} + +message NetworkInterfaceAddress { + string interface_name = 1; + string ip_address = 2; + int32 family = 3; + bool loopback = 4; +} + +message DeviceManagerSnapshot { + string manager_name = 1; + string manager_version = 2; + string manager_description = 3; + uint64 sampled_at_unix_ms = 4; + repeated ManagedDeviceStatus devices = 5; +} + +message ManagedDeviceStatus { + string device_id = 1; + int32 kind = 2; + string type_name = 3; + bool enabled = 4; + int32 manager_state = 5; + int32 health = 6; + bool has_error = 7; + string error_message = 8; + uint64 status_updated_at_unix_ms = 9; +} + +message MediaTrack { + uint32 track_id = 1; + int32 kind = 2; + string device_id = 3; + string source_track_id = 4; + string codec = 5; + string payload_format = 6; + uint64 codec_generation = 7; + uint32 codec_generation_token = 8; + uint32 width = 9; + uint32 height = 10; + uint32 frames_per_second = 11; + uint32 sample_rate = 12; + uint32 channels = 13; + bytes codec_config = 14; +} + +message MediaFrame { + string node_id = 1; + string session_id = 2; + uint64 session_epoch = 3; + uint32 track_id = 4; + int32 kind = 5; + string codec = 6; + string payload_format = 7; + uint64 packet_sequence = 8; + uint64 frame_sequence = 9; + uint64 capture_timestamp_us = 10; + uint32 flags = 11; + bytes payload = 12; +} diff --git a/server/quic/README.md b/server/quic/README.md new file mode 100644 index 0000000..d52bc97 --- /dev/null +++ b/server/quic/README.md @@ -0,0 +1,21 @@ +# CMVR QUIC 接收服务 + +Node 服务启动时会同时启动 `server/quic/server.py`,在 UDP 5175 上接收 +CMVR edge 消息(原有 Web 服务仍使用 TCP 5175)。首次运行先安装 Python 依赖: + +```powershell +py -3 -m pip install -r server/quic/requirements.txt +``` + +可用环境变量: + +- `QUIC_ENABLED=false`:不随 Node 启动 QUIC 服务。 +- `QUIC_HOST` / `QUIC_PORT`:监听地址和 UDP 端口,默认 `0.0.0.0:5175`。 +- `QUIC_CERT_FILE` / `QUIC_KEY_FILE`:TLS 证书和私钥。 +- `QUIC_ALPN`:逗号分隔的 ALPN,默认 `cmvr-quic-edge/1`。 +- `QUIC_HEARTBEAT_INTERVAL_MS`:注册响应下发的心跳周期,默认 5000。 +- `QUIC_PYTHON`:Python 可执行文件路径。 + +控制消息以 JSON 行输出到服务日志。注册和心跳消息会分别得到 +`NodeRegisterResponse` 和 `NodeHeartbeatAck` 回复。媒体 DATAGRAM 的 64 字节 +CMQD 帧头会被解析,分片完整后以 `media_frame` JSON 事件和 Base64 数据输出。 diff --git a/server/quic/generated/__init__.py b/server/quic/generated/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/server/quic/generated/quic_edge_pb2.py b/server/quic/generated/quic_edge_pb2.py new file mode 100644 index 0000000..41e58aa --- /dev/null +++ b/server/quic/generated/quic_edge_pb2.py @@ -0,0 +1,72 @@ +# -*- coding: utf-8 -*- +# Generated by the protocol buffer compiler. DO NOT EDIT! +# NO CHECKED-IN PROTOBUF GENCODE +# source: quic_edge.proto +# Protobuf Python Version: 7.35.1 +"""Generated protocol buffer code.""" +from google.protobuf import descriptor as _descriptor +from google.protobuf import descriptor_pool as _descriptor_pool +from google.protobuf import runtime_version as _runtime_version +from google.protobuf import symbol_database as _symbol_database +from google.protobuf.internal import builder as _builder +_runtime_version.ValidateProtobufRuntimeVersion( + _runtime_version.Domain.PUBLIC, + 7, + 35, + 1, + '', + 'quic_edge.proto' +) +# @@protoc_insertion_point(imports) + +_sym_db = _symbol_database.Default() + + + + +DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n\x0fquic_edge.proto\x12\x11\x63mvr.quic_edge.v1\"\xf6\x04\n\x13\x45\x64geControlEnvelope\x12\x18\n\x10protocol_version\x18\x01 \x01(\r\x12\x18\n\x10message_sequence\x18\x02 \x01(\x04\x12G\n\x15node_register_request\x18\n \x01(\x0b\x32&.cmvr.quic_edge.v1.NodeRegisterRequestH\x00\x12I\n\x16node_register_response\x18\x0b \x01(\x0b\x32\'.cmvr.quic_edge.v1.NodeRegisterResponseH\x00\x12:\n\x0enode_heartbeat\x18\x0c \x01(\x0b\x32 .cmvr.quic_edge.v1.NodeHeartbeatH\x00\x12\x41\n\x12node_heartbeat_ack\x18\r \x01(\x0b\x32#.cmvr.quic_edge.v1.NodeHeartbeatAckH\x00\x12\x41\n\x12media_session_open\x18\x14 \x01(\x0b\x32#.cmvr.quic_edge.v1.MediaSessionOpenH\x00\x12I\n\x16media_track_descriptor\x18\x15 \x01(\x0b\x32\'.cmvr.quic_edge.v1.MediaTrackDescriptorH\x00\x12\x43\n\x13media_session_close\x18\x16 \x01(\x0b\x32$.cmvr.quic_edge.v1.MediaSessionCloseH\x00\x12:\n\x0eprotocol_error\x18\x1e \x01(\x0b\x32 .cmvr.quic_edge.v1.ProtocolErrorH\x00\x42\t\n\x07payload\"\x84\x02\n\x17NetworkInterfaceAddress\x12\x16\n\x0einterface_name\x18\x01 \x01(\t\x12\x12\n\nip_address\x18\x02 \x01(\t\x12H\n\x06\x66\x61mily\x18\x03 \x01(\x0e\x32\x38.cmvr.quic_edge.v1.NetworkInterfaceAddress.AddressFamily\x12\x10\n\x08loopback\x18\x04 \x01(\x08\"a\n\rAddressFamily\x12\x1e\n\x1a\x41\x44\x44RESS_FAMILY_UNSPECIFIED\x10\x00\x12\x17\n\x13\x41\x44\x44RESS_FAMILY_IPV4\x10\x01\x12\x17\n\x13\x41\x44\x44RESS_FAMILY_IPV6\x10\x02\"7\n\x0cGrpcEndpoint\x12\x0c\n\x04host\x18\x01 \x01(\t\x12\x0c\n\x04port\x18\x02 \x01(\r\x12\x0b\n\x03tls\x18\x03 \x01(\x08\"\xbb\x02\n\x13ManagedDeviceStatus\x12\x11\n\tdevice_id\x18\x01 \x01(\t\x12+\n\x04kind\x18\x02 \x01(\x0e\x32\x1d.cmvr.quic_edge.v1.DeviceKind\x12\x11\n\ttype_name\x18\x03 \x01(\t\x12\x0f\n\x07\x65nabled\x18\x04 \x01(\x08\x12<\n\rmanager_state\x18\x05 \x01(\x0e\x32%.cmvr.quic_edge.v1.ManagedDeviceState\x12\x35\n\x06health\x18\x06 \x01(\x0e\x32%.cmvr.quic_edge.v1.DeviceHealthStatus\x12\x11\n\thas_error\x18\x07 \x01(\x08\x12\x15\n\rerror_message\x18\x08 \x01(\t\x12!\n\x19status_updated_at_unix_ms\x18\t \x01(\x04\"\xb8\x01\n\x15\x44\x65viceManagerSnapshot\x12\x14\n\x0cmanager_name\x18\x01 \x01(\t\x12\x17\n\x0fmanager_version\x18\x02 \x01(\t\x12\x1b\n\x13manager_description\x18\x03 \x01(\t\x12\x37\n\x07\x64\x65vices\x18\x04 \x03(\x0b\x32&.cmvr.quic_edge.v1.ManagedDeviceStatus\x12\x1a\n\x12sampled_at_unix_ms\x18\x05 \x01(\x04\"\xdc\x01\n\x0eNodeDescriptor\x12\x0f\n\x07node_id\x18\x01 \x01(\t\x12\x0f\n\x07\x62oot_id\x18\x02 \x01(\t\x12\x18\n\x10software_version\x18\x03 \x01(\t\x12\x44\n\x10local_interfaces\x18\x04 \x03(\x0b\x32*.cmvr.quic_edge.v1.NetworkInterfaceAddress\x12\x36\n\rgrpc_endpoint\x18\x05 \x01(\x0b\x32\x1f.cmvr.quic_edge.v1.GrpcEndpoint\x12\x10\n\x08robot_id\x18\x06 \x01(\t\"_\n\x13NodeRegisterRequest\x12/\n\x04node\x18\x01 \x01(\x0b\x32!.cmvr.quic_edge.v1.NodeDescriptor\x12\x17\n\x0fsent_at_unix_ms\x18\x02 \x01(\x04\"\x88\x01\n\x14NodeRegisterResponse\x12\x10\n\x08\x61\x63\x63\x65pted\x18\x01 \x01(\x08\x12\x12\n\nsession_id\x18\x02 \x01(\t\x12\x0f\n\x07message\x18\x03 \x01(\t\x12\x1d\n\x15heartbeat_interval_ms\x18\x04 \x01(\r\x12\x1a\n\x12observed_source_ip\x18\x05 \x01(\t\"\xdc\x02\n\rNodeHeartbeat\x12\x0f\n\x07node_id\x18\x01 \x01(\t\x12\x0f\n\x07\x62oot_id\x18\x02 \x01(\t\x12\x12\n\nsession_id\x18\x03 \x01(\t\x12\x10\n\x08sequence\x18\x04 \x01(\x04\x12\x17\n\x0fsent_at_unix_ms\x18\x05 \x01(\x04\x12\x18\n\x10software_version\x18\x06 \x01(\t\x12\x44\n\x10local_interfaces\x18\x07 \x03(\x0b\x32*.cmvr.quic_edge.v1.NetworkInterfaceAddress\x12\x36\n\rgrpc_endpoint\x18\x08 \x01(\x0b\x32\x1f.cmvr.quic_edge.v1.GrpcEndpoint\x12@\n\x0e\x64\x65vice_manager\x18\t \x01(\x0b\x32(.cmvr.quic_edge.v1.DeviceManagerSnapshot\x12\x10\n\x08robot_id\x18\n \x01(\t\"\xa1\x01\n\x10NodeHeartbeatAck\x12\x10\n\x08\x61\x63\x63\x65pted\x18\x01 \x01(\x08\x12\x1d\n\x15\x61\x63knowledged_sequence\x18\x02 \x01(\x04\x12\x0f\n\x07message\x18\x03 \x01(\t\x12\x1b\n\x13server_time_unix_ms\x18\x04 \x01(\x04\x12\x1a\n\x12observed_source_ip\x18\x05 \x01(\t\x12\x12\n\nsession_id\x18\x06 \x01(\t\"N\n\x10MediaSessionOpen\x12\x0f\n\x07node_id\x18\x01 \x01(\t\x12\x15\n\rsession_epoch\x18\x02 \x01(\x04\x12\x12\n\nsession_id\x18\x03 \x01(\t\"\xd8\x02\n\x14MediaTrackDescriptor\x12\x10\n\x08track_id\x18\x01 \x01(\r\x12*\n\x04kind\x18\x02 \x01(\x0e\x32\x1c.cmvr.quic_edge.v1.MediaKind\x12\x11\n\tdevice_id\x18\x03 \x01(\t\x12\r\n\x05\x63odec\x18\x04 \x01(\t\x12\x18\n\x10\x63odec_generation\x18\x05 \x01(\x04\x12\x17\n\x0fsource_track_id\x18\x06 \x01(\t\x12\x1e\n\x16\x63odec_generation_token\x18\x07 \x01(\r\x12\x16\n\x0epayload_format\x18\x08 \x01(\t\x12\r\n\x05width\x18\n \x01(\r\x12\x0e\n\x06height\x18\x0b \x01(\r\x12\x19\n\x11\x66rames_per_second\x18\x0c \x01(\r\x12\x13\n\x0bsample_rate\x18\x14 \x01(\r\x12\x10\n\x08\x63hannels\x18\x15 \x01(\r\x12\x14\n\x0c\x63odec_config\x18\x1e \x01(\x0c\"N\n\x11MediaSessionClose\x12\x0e\n\x06reason\x18\x01 \x01(\t\x12\x12\n\nsession_id\x18\x02 \x01(\t\x12\x15\n\rsession_epoch\x18\x03 \x01(\x04\"_\n\rProtocolError\x12\x0c\n\x04\x63ode\x18\x01 \x01(\r\x12\x0f\n\x07message\x18\x02 \x01(\t\x12 \n\x18related_message_sequence\x18\x03 \x01(\x04\x12\r\n\x05\x66\x61tal\x18\x04 \x01(\x08*\xeb\x02\n\nDeviceKind\x12\x1b\n\x17\x44\x45VICE_KIND_UNSPECIFIED\x10\x00\x12\x13\n\x0f\x44\x45VICE_KIND_AGV\x10\x01\x12\x13\n\x0f\x44\x45VICE_KIND_ARM\x10\x02\x12\x17\n\x13\x44\x45VICE_KIND_BATTERY\x10\x03\x12\x18\n\x14\x44\x45VICE_KIND_BIO_HEAD\x10\x04\x12\x16\n\x12\x44\x45VICE_KIND_CAMERA\x10\x05\x12\x17\n\x13\x44\x45VICE_KIND_CAN_BUS\x10\x06\x12\x18\n\x14\x44\x45VICE_KIND_DEX_HAND\x10\x07\x12\x17\n\x13\x44\x45VICE_KIND_GRIPPER\x10\x08\x12\x1a\n\x16\x44\x45VICE_KIND_MICROPHONE\x10\t\x12\x15\n\x11\x44\x45VICE_KIND_MOTOR\x10\n\x12\x1c\n\x18\x44\x45VICE_KIND_MOTOR_SYSTEM\x10\x0b\x12\x15\n\x11\x44\x45VICE_KIND_ROBOT\x10\x0c\x12\x17\n\x13\x44\x45VICE_KIND_SPEAKER\x10\r*\xad\x02\n\x12ManagedDeviceState\x12$\n MANAGED_DEVICE_STATE_UNSPECIFIED\x10\x00\x12!\n\x1dMANAGED_DEVICE_STATE_DISABLED\x10\x01\x12%\n!MANAGED_DEVICE_STATE_INITIALIZING\x10\x02\x12#\n\x1fMANAGED_DEVICE_STATE_REGISTERED\x10\x03\x12\x1e\n\x1aMANAGED_DEVICE_STATE_READY\x10\x04\x12 \n\x1cMANAGED_DEVICE_STATE_RUNNING\x10\x05\x12 \n\x1cMANAGED_DEVICE_STATE_STOPPED\x10\x06\x12\x1e\n\x1aMANAGED_DEVICE_STATE_ERROR\x10\x07*\x9f\x01\n\x12\x44\x65viceHealthStatus\x12$\n DEVICE_HEALTH_STATUS_UNSPECIFIED\x10\x00\x12 \n\x1c\x44\x45VICE_HEALTH_STATUS_HEALTHY\x10\x01\x12!\n\x1d\x44\x45VICE_HEALTH_STATUS_DEGRADED\x10\x02\x12\x1e\n\x1a\x44\x45VICE_HEALTH_STATUS_FAULT\x10\x03*S\n\tMediaKind\x12\x1a\n\x16MEDIA_KIND_UNSPECIFIED\x10\x00\x12\x14\n\x10MEDIA_KIND_VIDEO\x10\x01\x12\x14\n\x10MEDIA_KIND_AUDIO\x10\x02\x62\x06proto3') + +_globals = globals() +_builder.BuildMessageAndEnumDescriptors(DESCRIPTOR, _globals) +_builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'quic_edge_pb2', _globals) +if not _descriptor._USE_C_DESCRIPTORS: + DESCRIPTOR._loaded_options = None + _globals['_DEVICEKIND']._serialized_start=3075 + _globals['_DEVICEKIND']._serialized_end=3438 + _globals['_MANAGEDDEVICESTATE']._serialized_start=3441 + _globals['_MANAGEDDEVICESTATE']._serialized_end=3742 + _globals['_DEVICEHEALTHSTATUS']._serialized_start=3745 + _globals['_DEVICEHEALTHSTATUS']._serialized_end=3904 + _globals['_MEDIAKIND']._serialized_start=3906 + _globals['_MEDIAKIND']._serialized_end=3989 + _globals['_EDGECONTROLENVELOPE']._serialized_start=39 + _globals['_EDGECONTROLENVELOPE']._serialized_end=669 + _globals['_NETWORKINTERFACEADDRESS']._serialized_start=672 + _globals['_NETWORKINTERFACEADDRESS']._serialized_end=932 + _globals['_NETWORKINTERFACEADDRESS_ADDRESSFAMILY']._serialized_start=835 + _globals['_NETWORKINTERFACEADDRESS_ADDRESSFAMILY']._serialized_end=932 + _globals['_GRPCENDPOINT']._serialized_start=934 + _globals['_GRPCENDPOINT']._serialized_end=989 + _globals['_MANAGEDDEVICESTATUS']._serialized_start=992 + _globals['_MANAGEDDEVICESTATUS']._serialized_end=1307 + _globals['_DEVICEMANAGERSNAPSHOT']._serialized_start=1310 + _globals['_DEVICEMANAGERSNAPSHOT']._serialized_end=1494 + _globals['_NODEDESCRIPTOR']._serialized_start=1497 + _globals['_NODEDESCRIPTOR']._serialized_end=1717 + _globals['_NODEREGISTERREQUEST']._serialized_start=1719 + _globals['_NODEREGISTERREQUEST']._serialized_end=1814 + _globals['_NODEREGISTERRESPONSE']._serialized_start=1817 + _globals['_NODEREGISTERRESPONSE']._serialized_end=1953 + _globals['_NODEHEARTBEAT']._serialized_start=1956 + _globals['_NODEHEARTBEAT']._serialized_end=2304 + _globals['_NODEHEARTBEATACK']._serialized_start=2307 + _globals['_NODEHEARTBEATACK']._serialized_end=2468 + _globals['_MEDIASESSIONOPEN']._serialized_start=2470 + _globals['_MEDIASESSIONOPEN']._serialized_end=2548 + _globals['_MEDIATRACKDESCRIPTOR']._serialized_start=2551 + _globals['_MEDIATRACKDESCRIPTOR']._serialized_end=2895 + _globals['_MEDIASESSIONCLOSE']._serialized_start=2897 + _globals['_MEDIASESSIONCLOSE']._serialized_end=2975 + _globals['_PROTOCOLERROR']._serialized_start=2977 + _globals['_PROTOCOLERROR']._serialized_end=3072 +# @@protoc_insertion_point(module_scope) diff --git a/server/quic/requirements.txt b/server/quic/requirements.txt new file mode 100644 index 0000000..dde6b4e --- /dev/null +++ b/server/quic/requirements.txt @@ -0,0 +1,3 @@ +aioquic==1.3.0 +protobuf==7.35.1 + diff --git a/server/quic/server.py b/server/quic/server.py new file mode 100644 index 0000000..3cbf30d --- /dev/null +++ b/server/quic/server.py @@ -0,0 +1,258 @@ +"""CMVR edge QUIC gateway. + +Control stream records are protobuf messages prefixed by a four-byte, +network-byte-order unsigned length. +""" + +import asyncio +import json +import os +import signal +import struct +import time +import uuid +from base64 import b64encode +from pathlib import Path + +from aioquic.asyncio import QuicConnectionProtocol +from aioquic.asyncio.server import QuicServer +from aioquic.quic.configuration import QuicConfiguration +from aioquic.quic.events import ( + ConnectionTerminated, + DatagramFrameReceived, + HandshakeCompleted, + StreamDataReceived, +) +from google.protobuf.json_format import MessageToDict + +from generated import quic_edge_pb2 + + +PROTOCOL_VERSION = 1 + + +def log(event: str, **fields): + print(json.dumps({"source": "quic", "event": event, **fields}, ensure_ascii=False), flush=True) + + +MAX_CONTROL_FRAME_BYTES = int(os.getenv("QUIC_MAX_CONTROL_FRAME_BYTES", "1048576")) +DATAGRAM_HEADER_BYTES = 64 +DATAGRAM_MAGIC = b"CMQD" +MEDIA_REASSEMBLY_TIMEOUT_SECONDS = 5 + + +def pop_delimited(buffer: bytearray): + if len(buffer) < 4: + return None + length = int.from_bytes(buffer[:4], "big") + if length == 0 or length > MAX_CONTROL_FRAME_BYTES: + raise ValueError(f"invalid control frame length: {length}") + if len(buffer) < 4 + length: + return None + payload = bytes(buffer[4:4 + length]) + del buffer[:4 + length] + return payload + + +class CmvrQuicProtocol(QuicConnectionProtocol): + def __init__(self, *args, **kwargs): + super().__init__(*args, **kwargs) + self.buffers = {} + self.sessions = {} + self.out_sequence = 0 + self.media_frames = {} + + def peer(self): + paths = getattr(self._quic, "_network_paths", []) + address = paths[0].addr if paths else ("unknown", 0) + return {"ip": str(address[0]), "port": int(address[1])} + + def quic_event_received(self, event): + if isinstance(event, HandshakeCompleted): + log("connected", peer=self.peer(), alpn=event.alpn_protocol) + elif isinstance(event, StreamDataReceived): + self.receive_stream(event.stream_id, event.data, event.end_stream) + elif isinstance(event, DatagramFrameReceived): + self.receive_datagram(event.data) + elif isinstance(event, ConnectionTerminated): + log("disconnected", peer=self.peer(), error_code=event.error_code, reason=event.reason_phrase) + + def receive_stream(self, stream_id: int, data: bytes, end_stream: bool): + buffer = self.buffers.setdefault(stream_id, bytearray()) + buffer.extend(data) + try: + while payload := pop_delimited(buffer): + envelope = quic_edge_pb2.EdgeControlEnvelope() + envelope.ParseFromString(payload) + self.handle_envelope(stream_id, envelope) + except Exception as error: + log("protocol_error", peer=self.peer(), stream_id=stream_id, message=str(error)) + self.send_error(stream_id, 1, str(error), fatal=True) + if end_stream: + if buffer: + log("truncated_message", peer=self.peer(), stream_id=stream_id, remaining_bytes=len(buffer)) + self.buffers.pop(stream_id, None) + + def handle_envelope(self, stream_id, envelope): + payload_type = envelope.WhichOneof("payload") + message = getattr(envelope, payload_type) if payload_type else None + log( + "control_message", + peer=self.peer(), + stream_id=stream_id, + protocol_version=envelope.protocol_version, + message_sequence=str(envelope.message_sequence), + payload_type=payload_type, + payload=MessageToDict(message, preserving_proto_field_name=True) if message else None, + ) + + if envelope.protocol_version != PROTOCOL_VERSION: + self.send_error(stream_id, 2, f"unsupported protocol version: {envelope.protocol_version}", envelope.message_sequence, True) + return + if payload_type == "node_register_request": + request = envelope.node_register_request + session_id = str(uuid.uuid4()) + self.sessions[request.node.node_id] = session_id + response = quic_edge_pb2.EdgeControlEnvelope(protocol_version=PROTOCOL_VERSION) + response.node_register_response.CopyFrom(quic_edge_pb2.NodeRegisterResponse( + accepted=True, + session_id=session_id, + message="registered", + heartbeat_interval_ms=int(os.getenv("QUIC_HEARTBEAT_INTERVAL_MS", "5000")), + observed_source_ip=self.peer()["ip"], + )) + self.send_envelope(stream_id, response) + elif payload_type == "node_heartbeat": + heartbeat = envelope.node_heartbeat + expected = self.sessions.get(heartbeat.node_id) + accepted = bool(expected and expected == heartbeat.session_id) + response = quic_edge_pb2.EdgeControlEnvelope(protocol_version=PROTOCOL_VERSION) + response.node_heartbeat_ack.CopyFrom(quic_edge_pb2.NodeHeartbeatAck( + accepted=accepted, + acknowledged_sequence=heartbeat.sequence, + message="ok" if accepted else "unknown or mismatched session", + server_time_unix_ms=int(time.time() * 1000), + observed_source_ip=self.peer()["ip"], + session_id=expected or "", + )) + self.send_envelope(stream_id, response) + + def receive_datagram(self, data: bytes): + now = time.monotonic() + for key, frame in list(self.media_frames.items()): + if now - frame["updated_at"] > MEDIA_REASSEMBLY_TIMEOUT_SECONDS: + del self.media_frames[key] + + if len(data) < DATAGRAM_HEADER_BYTES or data[:4] != DATAGRAM_MAGIC: + log("invalid_datagram", peer=self.peer(), size=len(data), message="short packet or magic mismatch") + return + values = struct.unpack(">4sBBHHHHHIIQQQQII", data[:DATAGRAM_HEADER_BYTES]) + (_, version, kind, flags, header_size, fragment_index, fragment_count, + payload_size, track_id, generation_token, session_epoch, + packet_sequence, frame_sequence, capture_timestamp_us, frame_size, + fragment_offset) = values + payload = data[DATAGRAM_HEADER_BYTES:] + if (version != PROTOCOL_VERSION or header_size != DATAGRAM_HEADER_BYTES or + kind not in (1, 2) or fragment_count == 0 or + fragment_index >= fragment_count or payload_size != len(payload) or + frame_size == 0 or fragment_offset + payload_size > frame_size): + log("invalid_datagram", peer=self.peer(), size=len(data), message="invalid CMQD header fields") + return + + header = { + "kind": "video" if kind == 1 else "audio", + "flags": flags, + "key_frame": bool(flags & 1), + "discontinuity": bool(flags & 4), + "fragment_index": fragment_index, + "fragment_count": fragment_count, + "payload_size": payload_size, + "track_id": track_id, + "codec_generation_token": generation_token, + "session_epoch": str(session_epoch), + "packet_sequence": str(packet_sequence), + "frame_sequence": str(frame_sequence), + "capture_timestamp_us": str(capture_timestamp_us), + "frame_size": frame_size, + "fragment_offset": fragment_offset, + } + log("media_fragment", peer=self.peer(), **header) + + key = (session_epoch, track_id, frame_sequence) + frame = self.media_frames.setdefault(key, { + "data": bytearray(frame_size), "received": set(), "header": header, + "fragment_count": fragment_count, "updated_at": now, + }) + if frame["fragment_count"] != fragment_count or len(frame["data"]) != frame_size: + del self.media_frames[key] + log("invalid_datagram", peer=self.peer(), message="inconsistent fragments for frame", **header) + return + frame["data"][fragment_offset:fragment_offset + payload_size] = payload + frame["received"].add(fragment_index) + frame["updated_at"] = now + if len(frame["received"]) == fragment_count: + log("media_frame", peer=self.peer(), **frame["header"], data_base64=b64encode(frame["data"]).decode()) + del self.media_frames[key] + + def send_envelope(self, stream_id, envelope): + self.out_sequence += 1 + envelope.message_sequence = self.out_sequence + data = envelope.SerializeToString() + self._quic.send_stream_data(stream_id, len(data).to_bytes(4, "big") + data, end_stream=False) + self.transmit() + + def send_error(self, stream_id, code, message, related_sequence=0, fatal=False): + envelope = quic_edge_pb2.EdgeControlEnvelope(protocol_version=PROTOCOL_VERSION) + envelope.protocol_error.CopyFrom(quic_edge_pb2.ProtocolError( + code=code, message=message, related_message_sequence=related_sequence, fatal=fatal + )) + self.send_envelope(stream_id, envelope) + + +class LoggingQuicServer(QuicServer): + """Expose packet arrival even when QUIC/TLS negotiation cannot complete.""" + + def __init__(self, *args, **kwargs): + super().__init__(*args, **kwargs) + self._last_packet_log = {} + + def datagram_received(self, data, addr): + now = time.monotonic() + if now - self._last_packet_log.get(addr, 0) >= 5: + log("udp_packet_received", peer={"ip": str(addr[0]), "port": int(addr[1])}, size=len(data)) + self._last_packet_log[addr] = now + super().datagram_received(data, addr) + + +async def main(): + host = os.getenv("QUIC_HOST", "0.0.0.0") + port = int(os.getenv("QUIC_PORT", "5175")) + cert = Path(os.getenv("QUIC_CERT_FILE", "server/certs/quic-gateway.crt")) + key = Path(os.getenv("QUIC_KEY_FILE", "server/certs/quic-gateway.key")) + alpn = [item.strip() for item in os.getenv("QUIC_ALPN", "cmvr-quic-edge/1").split(",") if item.strip()] + + configuration = QuicConfiguration(is_client=False, alpn_protocols=alpn, max_datagram_frame_size=65536) + configuration.load_cert_chain(cert, key) + loop = asyncio.get_running_loop() + _, server = await loop.create_datagram_endpoint( + lambda: LoggingQuicServer( + configuration=configuration, + create_protocol=CmvrQuicProtocol, + ), + local_addr=(host, port), + ) + log("listening", address=f"{host}:{port}", alpn=alpn, cert=str(cert)) + + stopped = asyncio.Event() + loop = asyncio.get_running_loop() + for sig in (signal.SIGINT, signal.SIGTERM): + try: + loop.add_signal_handler(sig, stopped.set) + except NotImplementedError: + signal.signal(sig, lambda *_: loop.call_soon_threadsafe(stopped.set)) + await stopped.wait() + server.close() + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/server/quicServer.js b/server/quicServer.js new file mode 100644 index 0000000..51c82a6 --- /dev/null +++ b/server/quicServer.js @@ -0,0 +1,60 @@ +import { spawn } from 'node:child_process' +import path from 'node:path' +import { fileURLToPath } from 'node:url' + +const serverDir = path.dirname(fileURLToPath(import.meta.url)) + +export function startQuicServer({ onEvent } = {}) { + if (process.env.QUIC_ENABLED === 'false') return null + + const script = path.join(serverDir, 'quic', 'server.py') + const python = process.env.QUIC_PYTHON || (process.platform === 'win32' ? 'py' : 'python3') + const args = process.platform === 'win32' && !process.env.QUIC_PYTHON + ? ['-3', script] + : [script] + + const child = spawn(python, args, { + cwd: path.resolve(serverDir, '..'), + env: { + ...process.env, + QUIC_HOST: process.env.QUIC_HOST || '0.0.0.0', + QUIC_PORT: process.env.QUIC_PORT || '5175', + QUIC_CERT_FILE: process.env.QUIC_CERT_FILE || path.join(serverDir, 'certs', 'quic-gateway.crt'), + QUIC_KEY_FILE: process.env.QUIC_KEY_FILE || path.join(serverDir, 'certs', 'quic-gateway.key') + }, + stdio: ['ignore', 'pipe', 'pipe'] + }) + + child.stdout.setEncoding('utf8') + child.stderr.setEncoding('utf8') + let stdoutBuffer = '' + child.stdout.on('data', data => { + stdoutBuffer += data + const lines = stdoutBuffer.split(/\r?\n/) + stdoutBuffer = lines.pop() || '' + for (const line of lines) { + if (!line) continue + console.log(line) + try { + onEvent?.(JSON.parse(line)) + } catch { + // Keep non-JSON Python diagnostics visible without feeding the registry. + } + } + }) + child.stderr.on('data', data => process.stderr.write(data)) + child.on('error', error => console.error('[quic] 无法启动 Python 服务:', error.message)) + child.on('exit', (code, signal) => { + if (code !== 0 && signal !== 'SIGTERM') { + console.error(`[quic] 服务异常退出 (code=${code}, signal=${signal || 'none'})`) + } + }) + + const stop = () => { + if (!child.killed) child.kill() + } + process.once('SIGINT', stop) + process.once('SIGTERM', stop) + process.once('exit', stop) + return child +} diff --git a/server/robotRegistry.js b/server/robotRegistry.js new file mode 100644 index 0000000..b6a5462 --- /dev/null +++ b/server/robotRegistry.js @@ -0,0 +1,84 @@ +const robots = new Map() + +const KIND_NAMES = { + DEVICE_KIND_AGV: 'agv', + DEVICE_KIND_ARM: 'arm', + DEVICE_KIND_BATTERY: 'battery', + DEVICE_KIND_BIO_HEAD: 'bioHead', + DEVICE_KIND_CAMERA: 'camera', + DEVICE_KIND_CAN_BUS: 'canBus', + DEVICE_KIND_DEX_HAND: 'dexHand', + DEVICE_KIND_GRIPPER: 'gripper', + DEVICE_KIND_MICROPHONE: 'microphone', + DEVICE_KIND_MOTOR: 'motor', + DEVICE_KIND_MOTOR_SYSTEM: 'motorSystem', + DEVICE_KIND_ROBOT: 'robot', + DEVICE_KIND_SPEAKER: 'speaker' +} + +function groupDevices(devices = []) { + return devices.reduce((groups, device) => { + const kind = KIND_NAMES[device.kind] || 'other' + if (!groups[kind]) groups[kind] = [] + groups[kind].push({ + deviceId: device.device_id, + typeName: device.type_name || '', + enabled: device.enabled !== false, + state: device.manager_state || 'MANAGED_DEVICE_STATE_UNSPECIFIED', + health: device.health || 'DEVICE_HEALTH_STATUS_UNSPECIFIED', + hasError: Boolean(device.has_error), + errorMessage: device.error_message || '' + }) + return groups + }, {}) +} + +function endpoint(peer, grpcEndpoint = {}) { + return peer?.ip && grpcEndpoint.port ? `${peer.ip}:${grpcEndpoint.port}` : '' +} + +export function consumeQuicEvent(event) { + if (event?.event !== 'control_message') return + const payload = event.payload || {} + + if (event.payload_type === 'node_register_request') { + const node = payload.node || {} + const robotId = node.robot_id || node.node_id + if (!robotId) return + const previous = robots.get(robotId) || {} + robots.set(robotId, { + ...previous, + robotId, + nodeId: node.node_id || previous.nodeId || '', + ip: endpoint(event.peer, node.grpc_endpoint) || previous.ip || '', + peerIp: event.peer?.ip || previous.peerIp || '', + grpcEndpoint: node.grpc_endpoint || previous.grpcEndpoint || {}, + devices: previous.devices || {}, + online: true, + lastSeenAt: new Date().toISOString() + }) + } + + if (event.payload_type === 'node_heartbeat') { + const robotId = payload.robot_id || payload.node_id + if (!robotId) return + const previous = robots.get(robotId) || {} + robots.set(robotId, { + ...previous, + robotId, + nodeId: payload.node_id || previous.nodeId || '', + ip: endpoint(event.peer, payload.grpc_endpoint) || previous.ip || '', + peerIp: event.peer?.ip || previous.peerIp || '', + grpcEndpoint: payload.grpc_endpoint || previous.grpcEndpoint || {}, + devices: groupDevices(payload.device_manager?.devices), + online: true, + heartbeatSequence: payload.sequence || '0', + lastSeenAt: new Date().toISOString() + }) + } +} + +export function listRobots() { + return [...robots.values()].sort((a, b) => a.robotId.localeCompare(b.robotId)) +} + diff --git a/src/api/robot.js b/src/api/robot.js new file mode 100644 index 0000000..2281ff2 --- /dev/null +++ b/src/api/robot.js @@ -0,0 +1,5 @@ +import request from '@/utils/request' + +export function getQuicRobots() { + return request('/api/quic/robots') +} diff --git a/src/utils/request.js b/src/utils/request.js index 9ae6466..10f0333 100644 --- a/src/utils/request.js +++ b/src/utils/request.js @@ -9,12 +9,15 @@ const request = async (url, params = {}, method = 'GET', headers = { 'Content-Ty }, 5000) try { - const res = await fetch(url, { + const options = { method, headers, - body: JSON.stringify(params), signal - }) + } + if (!['GET', 'HEAD'].includes(method.toUpperCase())) { + options.body = JSON.stringify(params) + } + const res = await fetch(url, options) const data = await res.json() if (!res.ok) { diff --git a/src/views/Address.vue b/src/views/Address.vue index cd44177..cdbc4b0 100644 --- a/src/views/Address.vue +++ b/src/views/Address.vue @@ -5,6 +5,28 @@
填写设备通讯参数,建立与AGV底盘和机械臂的连接
请操持设备与上位机在同一局域网内
+
+ + + + 刷新列表 +
+ + + + + @@ -38,18 +60,26 @@