Compare commits

...

3 Commits

27 changed files with 1144 additions and 245 deletions

7
.gitignore vendored
View File

@ -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

4
components.d.ts vendored
View File

@ -13,8 +13,6 @@ declare module 'vue' {
export interface GlobalComponents {
ElButton: typeof import('element-plus/es')['ElButton']
ElContainer: typeof import('element-plus/es')['ElContainer']
ElForm: typeof import('element-plus/es')['ElForm']
ElFormItem: typeof import('element-plus/es')['ElFormItem']
ElHeader: typeof import('element-plus/es')['ElHeader']
ElIcon: typeof import('element-plus/es')['ElIcon']
ElInput: typeof import('element-plus/es')['ElInput']
@ -24,6 +22,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']

View File

@ -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"

View File

@ -23,7 +23,7 @@ function waitForGrpcReady(client, timeoutMs = 5000) {
})
}
export function attachAudioWebSocket(httpServer) {
export function attachAudioWebSocket(httpServer, { resolveRobotAddress, resolveRobotDeviceId } = {}) {
const wss = new WebSocketServer({ server: httpServer, path: '/audio' })
wss.on('connection', (socket) => {
@ -101,14 +101,27 @@ export function attachAudioWebSocket(httpServer) {
}
if (message.type !== 'init' || initialized) return
const { microphoneId, robotAddress } = message
speakerId = message.speakerId?.trim()
if (![microphoneId, speakerId, robotAddress].every((value) => typeof value === 'string' && value.trim())) {
const { robotId } = message
if (typeof robotId !== 'string' || !robotId.trim()) {
sendJson(socket, { type: 'error', message: '机器人地址、麦克风 ID 和扬声器 ID 不能为空' })
socket.close(1008, 'invalid init')
return
}
const robotAddress = resolveRobotAddress?.(robotId.trim())
if (!robotAddress) {
sendJson(socket, { type: 'error', message: `机器人 ${robotId} 不在线或尚未注册` })
socket.close(1008, 'robot unavailable')
return
}
const microphoneId = resolveRobotDeviceId?.(robotId.trim(), 'microphone')
speakerId = resolveRobotDeviceId?.(robotId.trim(), 'speaker')
if (!microphoneId || !speakerId) {
sendJson(socket, { type: 'error', message: `机器人 ${robotId} 缺少麦克风或扬声器设备` })
socket.close(1008, 'audio device unavailable')
return
}
initialized = true
microphoneClient = createMicrophoneGrpcClient(robotAddress.trim())
speakerClient = createSpeakerGrpcClient(robotAddress.trim())

View File

@ -84,14 +84,9 @@ function message(error) {
}
export function streamCameraVideo(req, res) {
const ip = typeof req.query.ip === 'string' ? req.query.ip.trim() : ''
const deviceId = typeof req.query.deviceId === 'string' ? req.query.deviceId.trim() : ''
if (!ip || !deviceId) {
res.status(400).json({ code: 400, message: 'ip 和 deviceId 不能为空' })
return
}
const deviceId = req.deviceId
const client = createCameraGrpcClient(ip)
const client = createCameraGrpcClient(req.robotAddress)
const grpcStream = client.getRgbImageStream()
let transcoder
let closed = false

View File

@ -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-----

View File

@ -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-----

View File

@ -16,25 +16,61 @@ import {
} from './grpcClient.js';
import { attachAudioWebSocket } from './audioWebSocket.js';
import { streamCameraVideo } from './cameraVideoStream.js';
import { startQuicServer } from './quicServer.js';
import { consumeQuicEvent, getRobotAddress, getRobotDeviceId, 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() })
})
app.use('/api', (req, res, next) => {
const robotId = String(req.body?.robotId || req.query?.robotId || '').trim()
if (!robotId) {
res.status(400).json({ code: 400, message: 'robotId 不能为空' })
return
}
const address = getRobotAddress(robotId)
if (!address) {
res.status(404).json({ code: 404, message: `机器人 ${robotId} 不在线或尚未注册` })
return
}
req.robotId = robotId
req.robotAddress = address
next()
})
app.use('/api', (req, res, next) => {
let kind = ''
if (req.path.startsWith('/agv/')) kind = 'agv'
else if (req.path.startsWith('/arm/')) kind = 'arm'
else if (req.path.startsWith('/audio/microphone/')) kind = 'microphone'
else if (req.path.startsWith('/audio/speaker/')) kind = 'speaker'
else if (req.path.startsWith('/camera/')) kind = 'camera'
if (!kind) return next()
const deviceId = getRobotDeviceId(req.robotId, kind)
if (!deviceId) {
res.status(404).json({ code: 404, message: `机器人 ${req.robotId} 没有可用的 ${kind} 设备` })
return
}
req.deviceId = deviceId
next()
})
// 将相机逐帧 gRPC 数据转码为浏览器可直接播放的 fragmented MP4。
// 示例:/api/camera/video?ip=192.168.0.28%3A50052&deviceId=camera1
// 示例:/api/camera/video?robotId=robot-1
app.get('/api/camera/video', streamCameraVideo)
// 获取机器人麦克风音量(0-100)。
app.post('/api/audio/microphone/volume/get', (req, res) => {
const { ip, deviceId } = req.body
if (!ip || !deviceId) {
res.status(400).json({ code: 400, message: '机器人地址和麦克风设备 ID 不能为空' })
return
}
const deviceId = req.deviceId
const grpcClient = createMicrophoneGrpcClient(ip)
const grpcClient = createMicrophoneGrpcClient(req.robotAddress)
grpcClient.getVolume(
{ header: { device_id: deviceId } },
{ deadline: Date.now() + 5000 },
@ -59,18 +95,14 @@ app.post('/api/audio/microphone/volume/get', (req, res) => {
// 设置机器人麦克风音量(0-100)。
app.post('/api/audio/microphone/volume/set', (req, res) => {
const { ip, deviceId } = req.body
const deviceId = req.deviceId
const volume = Number(req.body.volume)
if (!ip || !deviceId) {
res.status(400).json({ code: 400, message: '机器人地址和麦克风设备 ID 不能为空' })
return
}
if (!Number.isInteger(volume) || volume < 0 || volume > 100) {
res.status(400).json({ code: 400, message: '音量必须是 0 到 100 的整数' })
return
}
const grpcClient = createMicrophoneGrpcClient(ip)
const grpcClient = createMicrophoneGrpcClient(req.robotAddress)
grpcClient.setVolume({
header: { device_id: deviceId },
volume
@ -94,13 +126,9 @@ app.post('/api/audio/microphone/volume/set', (req, res) => {
// 获取机器人扬声器音量(0-100)。
app.post('/api/audio/speaker/volume/get', (req, res) => {
const { ip, deviceId } = req.body
if (!ip || !deviceId) {
res.status(400).json({ code: 400, message: '机器人地址和扬声器设备 ID 不能为空' })
return
}
const deviceId = req.deviceId
const grpcClient = createSpeakerGrpcClient(ip)
const grpcClient = createSpeakerGrpcClient(req.robotAddress)
grpcClient.getVolume(
{ header: { device_id: deviceId } },
{ deadline: Date.now() + 5000 },
@ -125,18 +153,14 @@ app.post('/api/audio/speaker/volume/get', (req, res) => {
// 设置机器人扬声器音量(0-100)。
app.post('/api/audio/speaker/volume/set', (req, res) => {
const { ip, deviceId } = req.body
const deviceId = req.deviceId
const volume = Number(req.body.volume)
if (!ip || !deviceId) {
res.status(400).json({ code: 400, message: '机器人地址和扬声器设备 ID 不能为空' })
return
}
if (!Number.isInteger(volume) || volume < 0 || volume > 100) {
res.status(400).json({ code: 400, message: '音量必须是 0 到 100 的整数' })
return
}
const grpcClient = createSpeakerGrpcClient(ip)
const grpcClient = createSpeakerGrpcClient(req.robotAddress)
grpcClient.setVolume({
header: { device_id: deviceId },
volume
@ -162,11 +186,11 @@ app.post('/api/audio/speaker/volume/set', (req, res) => {
app.post('/api/agv/getRuntimeState', async (req, res) => {
const request = {
header: {
device_id: req.body.deviceId || 'agv_src1100'
device_id: req.deviceId
}
};
const grpcClient = createAGVGrpcClient(req.body.ip)
const grpcClient = createAGVGrpcClient(req.robotAddress)
const resp = await new Promise((resolve, reject) => {
grpcClient.getRuntimeState(request, (err, res) => {
@ -193,10 +217,10 @@ app.post('/api/agv/getRuntimeState', async (req, res) => {
// 获取机器人当前地图
app.post('/api/agv/getCurrentMap', async (req, res) => {
const grpcClient = createAGVGrpcClient(req.body.ip)
const grpcClient = createAGVGrpcClient(req.robotAddress)
const request = {
header: {
device_id: req.body.deviceId || 'agv_src1100'
device_id: req.deviceId
},
map_name: req.body.mapName
};
@ -226,10 +250,10 @@ app.post('/api/agv/getCurrentMap', async (req, res) => {
// 机器人运动控制
app.post('/api/agv/moveRobot', async (req, res) => {
const grpcClient = createAGVGrpcClient(req.body.ip)
const grpcClient = createAGVGrpcClient(req.robotAddress)
const request = {
header: {
device_id: req.body.deviceId || 'agv_src1100'
device_id: req.deviceId
},
velocity: {
vx: req.body.vx,
@ -266,10 +290,10 @@ app.post('/api/agv/moveRobot', async (req, res) => {
// 机器人停止运动
app.post('/api/agv/stopRobot', async (req, res) => {
const grpcClient = createAGVGrpcClient(req.body.ip)
const grpcClient = createAGVGrpcClient(req.robotAddress)
const request = {
header: {
device_id: req.body.deviceId || 'agv_src1100'
device_id: req.deviceId
}
};
@ -306,11 +330,11 @@ app.post('/api/agv/stopRobot', async (req, res) => {
app.post('/api/arm/getPose', async (req, res) => {
const request = {
header: {
device_id: req.body.deviceId || 'huayan_arm'
device_id: req.deviceId
}
};
const grpcClient = createARMGrpcClient(req.body.ip)
const grpcClient = createARMGrpcClient(req.robotAddress)
const resp = await new Promise((resolve, reject) => {
grpcClient.getPose(request, (err, res) => {
@ -339,11 +363,11 @@ app.post('/api/arm/getPose', async (req, res) => {
app.post('/api/arm/getJointState', async (req, res) => {
const request = {
header: {
device_id: req.body.deviceId || 'huayan_arm'
device_id: req.deviceId
}
};
const grpcClient = createARMGrpcClient(req.body.ip)
const grpcClient = createARMGrpcClient(req.robotAddress)
const resp = await new Promise((resolve, reject) => {
grpcClient.getJointState(request, (err, res) => {
@ -372,7 +396,7 @@ app.post('/api/arm/getJointState', async (req, res) => {
app.post('/api/arm/speedJ', async (req, res) => {
const request = {
header: {
device_id: req.body.deviceId || 'huayan_arm'
device_id: req.deviceId
},
velocity: {
velocity: req.body.velocities
@ -380,7 +404,7 @@ app.post('/api/arm/speedJ', async (req, res) => {
acceleration: req.body.acceleration,
duration: req.body.duration
};
const grpcClient = createARMGrpcClient(req.body.ip)
const grpcClient = createARMGrpcClient(req.robotAddress)
const resp = await new Promise((resolve, reject) => {
grpcClient.speedJ(request, (err, res) => {
@ -408,13 +432,13 @@ app.post('/api/arm/speedJ', async (req, res) => {
app.post('/api/arm/speedL', async (req, res) => {
const request = {
header: {
device_id: req.body.deviceId || 'huayan_arm'
device_id: req.deviceId
},
velocity: req.body.velocity,
acceleration: req.body.acceleration,
duration: req.body.duration
};
const grpcClient = createARMGrpcClient(req.body.ip)
const grpcClient = createARMGrpcClient(req.robotAddress)
const resp = await new Promise((resolve, reject) => {
grpcClient.speedL(request, (err, res) => {
@ -443,10 +467,10 @@ app.post('/api/arm/speedL', async (req, res) => {
app.post('/api/arm/stopMotion', async (req, res) => {
const request = {
// header: {
device_id: req.body.deviceId || 'huayan_arm',
device_id: req.deviceId,
// },
};
const grpcClient = createARMGrpcClient(req.body.ip)
const grpcClient = createARMGrpcClient(req.robotAddress)
const resp = await new Promise((resolve, reject) => {
grpcClient.stopMotion(request, (err, res) => {
@ -473,9 +497,9 @@ app.post('/api/arm/stopMotion', async (req, res) => {
app.post('/api/arm/torqueOn', async (req, res) => {
const request = {
device_id: req.body.deviceId || 'huayan_arm',
device_id: req.deviceId,
};
const grpcClient = createARMGrpcClient(req.body.ip)
const grpcClient = createARMGrpcClient(req.robotAddress)
const resp = await new Promise((resolve, reject) => {
grpcClient.torqueOn(request, (err, res) => {
@ -545,7 +569,11 @@ if (httpsEnabled) {
httpServer = http.createServer(app);
}
attachAudioWebSocket(httpServer)
attachAudioWebSocket(httpServer, {
resolveRobotAddress: getRobotAddress,
resolveRobotDeviceId: getRobotDeviceId
})
startQuicServer({ onEvent: consumeQuicEvent })
httpServer.listen(PORT, '0.0.0.0', () => {
const protocol = httpsEnabled ? 'https' : 'http';

View File

@ -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;
}

View File

@ -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;
}

21
server/quic/README.md Normal file
View File

@ -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 数据输出。

View File

File diff suppressed because one or more lines are too long

View File

@ -0,0 +1,3 @@
aioquic==1.3.0
protobuf==7.35.1

258
server/quic/server.py Normal file
View File

@ -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())

60
server/quicServer.js Normal file
View File

@ -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
}

93
server/robotRegistry.js Normal file
View File

@ -0,0 +1,93 @@
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))
}
export function getRobotAddress(robotId) {
if (!robotId) return ''
return robots.get(robotId)?.ip || ''
}
export function getRobotDeviceId(robotId, kind) {
const devices = robots.get(robotId)?.devices?.[kind] || []
return (devices.find(device => device.enabled) || devices[0])?.deviceId || ''
}

View File

@ -3,7 +3,7 @@ import request from '@/utils/request'
/**
* 机器人前往目标站点
* @param {Object} params
* @param {string} [params.deviceId] - 设备ID,默认 agv_src1100
* robotId 由统一请求层从 Pinia 自动附加,设备 ID 由后端注册表解析。
*/
export function getRuntimeState(params) {
return request('/api/agv/getRuntimeState', params, 'POST')

5
src/api/robot.js Normal file
View File

@ -0,0 +1,5 @@
import request from '@/utils/request'
export function getQuicRobots() {
return request('/api/quic/robots')
}

View File

@ -3,13 +3,8 @@ import { defineStore } from 'pinia'
// 唯一id:robot
export const useRobotStore = defineStore('robot', {
state: () => ({
ip: '192.168.0.28:50052',
robotId: '',
robotType: 'inspection',
agvDeviceId:'src1100',
almDeviceId: 'huayan_arm',
spkDeviceId: 'spk1',
micDeviceId: 'mic1',
cameraDeviceId: 'real_cam1',
position: { x: 0, y: 0, theta: 0 },
battery: {
percentage: 0.99,
@ -30,27 +25,12 @@ export const useRobotStore = defineStore('robot', {
}
},
actions: {
setIp(str) {
this.ip = str
setRobotId(str) {
this.robotId = str
},
setRobotType(str) {
this.robotType = str
},
setAgvDeviceId(str) {
this.agvDeviceId = str
},
setAlmDeviceId(str) {
this.almDeviceId = str
},
setSpkDeviceId(str) {
this.spkDeviceId = str
},
setMicDeviceId(str) {
this.micDeviceId = str
},
setCameraDeviceId(str) {
this.cameraDeviceId = str
},
// 更新机器人位置
setPosition(data = {}) {
this.position = data

View File

@ -1,4 +1,5 @@
import { ElMessage } from 'element-plus'
import { useRobotStore } from '@/stores/robot'
const request = async (url, params = {}, method = 'GET', headers = { 'Content-Type': 'application/json' }) => {
const controller = new AbortController()
@ -9,12 +10,17 @@ const request = async (url, params = {}, method = 'GET', headers = { 'Content-Ty
}, 5000)
try {
const res = await fetch(url, {
const robotId = useRobotStore().robotId
const requestParams = robotId && !params.robotId ? { ...params, robotId } : params
const options = {
method,
headers,
body: JSON.stringify(params),
signal
})
}
if (!['GET', 'HEAD'].includes(method.toUpperCase())) {
options.body = JSON.stringify(requestParams)
}
const res = await fetch(url, options)
const data = await res.json()
if (!res.ok) {

View File

@ -5,32 +5,28 @@
<div class="secondaryTitle">填写设备通讯参数,建立与AGV底盘和机械臂的连接</div>
<div class="tips">请操持设备与上位机在同一局域网内</div>
<div class="form">
<el-form ref="ruleFormRef" :model="ruleForm" :rules="rules" label-width="auto">
<el-form-item label="IP和端口" prop="ip">
<el-input v-model="ruleForm.ip" placeholder="请输入IP和端口" />
</el-form-item>
<el-form-item label="机器人类型" prop="robotType">
<el-select v-model="ruleForm.robotType" placeholder="请选择机器人类型">
<el-option label="巡检机器人" value="inspection"></el-option>
<el-option label="讲解机器人" value="explanation"></el-option>
<div class="robot-list">
<el-select
v-model="selectedRobotId"
placeholder="请选择已连接机器人"
clearable
filterable
@change="selectRobot"
>
<el-option
v-for="robot in robots"
:key="robot.robotId"
:label="robot.robotId"
:value="robot.robotId"
/>
</el-select>
</el-form-item>
<!-- <el-form-item label="底盘设备ID" prop="agvDeviceId">
<el-input v-model="ruleForm.agvDeviceId" placeholder="请输入底盘设备ID" />
</el-form-item>
<el-form-item label="机械臂设备ID" prop="almDeviceId">
<el-input v-model="ruleForm.almDeviceId" placeholder="请输入机械臂设备ID" />
</el-form-item>
<el-form-item label="扬声器设备ID" prop="spkDeviceId">
<el-input v-model="ruleForm.spkDeviceId" placeholder="请输入扬声器设备ID" />
</el-form-item>
<el-form-item label="麦克风设备ID" prop="micDeviceId">
<el-input v-model="ruleForm.micDeviceId" placeholder="请输入麦克风设备ID" />
</el-form-item>
<el-form-item label="相机设备ID" prop="cameraDeviceId">
<el-input v-model="ruleForm.cameraDeviceId" placeholder="请输入相机设备ID" />
</el-form-item> -->
</el-form>
<el-button :loading="robotListLoading" @click="loadRobots">刷新列表</el-button>
</div>
<el-table v-if="selectedRobot" :data="selectedDevices" size="small" max-height="220">
<el-table-column prop="kind" label="设备类型" width="120" />
<el-table-column prop="deviceId" label="设备 ID" />
<el-table-column prop="health" label="健康状态" width="190" />
</el-table>
<el-button :loading="loading" @click="testConnect" class="connect">进入上位机控制页面</el-button>
</div>
<div class="error-tips">连接异常排查:检查机器人状态、局域网联通、设备ID调谐是否匹配</div>
@ -38,78 +34,87 @@
</div>
</template>
<script setup>
import { ref } from 'vue'
import { computed, onMounted, onUnmounted, ref } from 'vue'
import { useRouter } from 'vue-router'
import { getRuntimeState } from '@/api/agv.js'
import { useRobotStore } from '@/stores/robot'
import { ElMessage } from 'element-plus'
import { ca } from 'element-plus/es/locales.mjs'
import { getQuicRobots } from '@/api/robot.js'
const router = useRouter()
const robotStore = useRobotStore()
const loading = ref(false)
const ruleFormRef = ref()
const ruleForm = ref({
ip: '',
robotType: 'inspection',
// agvDeviceId:'src1100',
// almDeviceId: 'huayan_arm',
// spkDeviceId: 'spk1',
// micDeviceId: 'mic1',
// cameraDeviceId: 'real_cam1'
})
const robotListLoading = ref(false)
const robots = ref([])
const selectedRobotId = ref(robotStore.robotId || '')
const selectedRobot = computed(() => robots.value.find(robot => robot.robotId === selectedRobotId.value))
const selectedDevices = computed(() => Object.entries(selectedRobot.value?.devices || {}).flatMap(
([kind, devices]) => devices.map(device => ({ kind, ...device }))
))
let refreshTimer
const rules = ref({
ip: [
{ required: true, message: '请输入机器人IP和端口', trigger: 'blur' }
],
robotType: [
{ required: true, message: '请选择机器人类型', trigger: 'blur' }
],
// agvDeviceId: [
// { required: true, message: '请输入底盘设备ID', trigger: 'blur' }
// ],
// almDeviceId: [
// { required: true, message: '请输入机械臂设备ID', trigger: 'blur' }
// ],
// spkDeviceId: [
// { required: true, message: '请输入扬声器设备ID', trigger: 'blur' }
// ],
// micDeviceId: [
// { required: true, message: '请输入麦克风设备ID', trigger: 'blur' }
// ],
// cameraDeviceId: [
// { required: true, message: '请输入相机设备ID', trigger: 'blur' }
// ],
})
const testConnect = async () => {
const valid = await ruleFormRef.value.validate();
if (valid) {
loading.value = true
robotStore.setIp(ruleForm.value.ip)
robotStore.setRobotType(ruleForm.value.robotType)
if (ruleForm.value.robotType === 'inspection') {
robotStore.setAgvDeviceId('src1100')
robotStore.setAlmDeviceId('huayan_arm')
robotStore.setSpkDeviceId('spk1')
robotStore.setMicDeviceId('mic1')
robotStore.setCameraDeviceId('real_cam1')
} else if (ruleForm.value.robotType === 'explanation') {
robotStore.setAgvDeviceId('src1100')
robotStore.setAlmDeviceId('huayan_arm')
robotStore.setSpkDeviceId('spk1')
robotStore.setMicDeviceId('mic1')
robotStore.setCameraDeviceId('real_cam1')
/**
* 通过机械臂数量判断机器人类型
* 1 个机械臂为巡检机器人,2 个机械臂为讲解机器人
* @param robot
*/
const getRobotType = (robot) => {
const armCount = robot?.devices?.arm?.length || 0
if (armCount === 1) return 'inspection'
if (armCount === 2) return 'explanation'
return ''
}
/**
* 选择机器人并且设置机器人类型
*/
const selectRobot = () => {
const robot = selectedRobot.value
if (!robot) return
robotStore.setRobotId(robot.robotId)
const robotType = getRobotType(robot)
if (robotType) robotStore.setRobotType(robotType)
}
/**
* 加载机器人列表
*/
const loadRobots = async () => {
robotListLoading.value = true
try {
const result = await getRuntimeState({
ip: robotStore.ip,
deviceId: robotStore.agvDeviceId
const result = await getQuicRobots()
robots.value = result.data || []
if (selectedRobotId.value) selectRobot()
} finally {
robotListLoading.value = false
}
}
onMounted(() => {
loadRobots()
refreshTimer = setInterval(loadRobots, 5000)
})
onUnmounted(() => clearInterval(refreshTimer))
const testConnect = async () => {
if (!selectedRobot.value) {
ElMessage.warning('请先选择已连接机器人')
return
}
const robotType = getRobotType(selectedRobot.value)
if (!robotType) {
ElMessage.warning('机器人机械臂数量异常,无法判断机器人类型')
return
}
loading.value = true
robotStore.setRobotId(selectedRobot.value.robotId)
robotStore.setRobotType(robotType)
try {
const result = await getRuntimeState()
loading.value = false
if (result.code === 200) {
robotStore.setPosition(result.data.pose)
@ -127,7 +132,6 @@ const testConnect = async () => {
ElMessage.error(error?.message || '连接服务错误')
}
}
}
</script>
<style lang="scss" scoped>
@ -174,11 +178,6 @@ const testConnect = async () => {
margin-bottom: 12px;
}
:deep(.el-form-item__label) {
font-size: 16px;
line-height: 40px;
}
:deep(.el-input__inner) {
height: 40px;
// 避免 iPad Safari 聚焦小字号输入框时自动缩放页面。
@ -202,6 +201,20 @@ const testConnect = async () => {
margin-top: 20px;
}
.robot-list {
display: flex;
gap: 12px;
margin-bottom: 16px;
.el-select {
flex: 1;
}
}
.el-table {
margin-bottom: 16px;
}
.error-tips {
text-align: center;
font-size: 14px;

View File

@ -303,7 +303,6 @@ import { Remove, CirclePlus } from '@element-plus/icons-vue'
import { useRobotStore } from '@/stores/robot'
const robotStore = useRobotStore()
const deviceId = robotStore.almDeviceId
const urdfViewerRef = ref()
@ -343,17 +342,11 @@ const tcp = ref({
})
const torqueOn = async () => {
await torqueOnApi({
ip: robotStore.ip,
deviceId: deviceId
})
await torqueOnApi()
}
const getPose = async () => {
const res = await getPoseApi({
ip: robotStore.ip,
deviceId: deviceId
})
const res = await getPoseApi()
if (res.code === 200) {
const { x, y, z, rx, ry, rz } = res.data
tcp.value = { x, y, z, rx, ry, rz }
@ -375,10 +368,7 @@ const jointNamsMap = {
*/
const getJointState = async (robot) => {
try {
const res = await getJointStateApi({
deviceId: deviceId,
ip: robotStore.ip,
})
const res = await getJointStateApi()
if (res.code === 200) {
setTimeout(() => {
const names = res.data.name
@ -404,8 +394,6 @@ const updateJoint = async (jointName, _dir) => {
const index = jointControls.value.findIndex(joint => joint.name === jointName)
velocities[index] = jointSpeed.value * (_dir === '+' ? 1 : -1)
const res = await speedJApi({
deviceId: deviceId,
ip: robotStore.ip,
acceleration: jointSpeed.value * (_dir === '+' ? 1 : -1) + 0.1,
duration: 60,
velocities: velocities
@ -475,8 +463,6 @@ const move = async (axis, direction, type) => {
}
speedLApi({
deviceId: deviceId,
ip: robotStore.ip,
duration: 60,
acceleration: type === 'posture' ? postureSpeed.value + 0.1 : endSpeed.value + 0.1,
velocity: {
@ -489,10 +475,7 @@ const move = async (axis, direction, type) => {
* 停止机械臂移动
*/
const stop = async () => {
await stopMotionApi({
deviceId: deviceId,
ip: robotStore.ip
})
await stopMotionApi()
}
const activePointerIds = new Set()

View File

@ -151,14 +151,11 @@ const statusMeta = computed(() => {
})
function microphoneVolumeRequestParams() {
return {
ip: robotStore.ip,
deviceId: robotStore.micDeviceId
}
return {}
}
async function loadMicrophoneVolume({ silent = false } = {}) {
if (!robotStore.ip || !robotStore.micDeviceId) return
if (!robotStore.robotId) return
isMicrophoneVolumeLoading.value = true
try {
const response = await getMicrophoneVolume(microphoneVolumeRequestParams())
@ -194,14 +191,11 @@ async function updateMicrophoneVolume(volume) {
}
function speakerVolumeRequestParams() {
return {
ip: robotStore.ip,
deviceId: robotStore.spkDeviceId
}
return {}
}
async function loadSpeakerVolume({ silent = false } = {}) {
if (!robotStore.ip || !robotStore.spkDeviceId) return
if (!robotStore.robotId) return
isSpeakerVolumeLoading.value = true
try {
const response = await getSpeakerVolume(speakerVolumeRequestParams())
@ -391,9 +385,7 @@ function connectAudioServer(sessionToken) {
const timeout = window.setTimeout(() => reject(new Error('连接音频服务超时')), JOIN_TIMEOUT)
socket.onopen = () => sendControlMessage({
type: 'init',
microphoneId: robotStore.micDeviceId,
speakerId: robotStore.spkDeviceId,
robotAddress: robotStore.ip
robotId: robotStore.robotId
})
socket.onmessage = (event) => {
if (typeof event.data !== 'string') {
@ -455,7 +447,7 @@ async function startCall() {
errorMessage.value = ''
const sessionToken = ++callSessionId
try {
if (!robotStore.ip || !robotStore.micDeviceId || !robotStore.spkDeviceId) {
if (!robotStore.robotId) {
throw new Error('请先配置机器人地址、麦克风 ID 和扬声器 ID')
}
// 必须在按钮点击的用户手势中创建,以符合浏览器自动播放策略。

View File

@ -14,20 +14,10 @@ import VoiceConversation from './VoiceConversation.vue';
import { computed } from 'vue';
import { useRobotStore } from '@/stores/robot';
const VIDEO_SERVER_HOST = import.meta.env.VITE_VIDEO_SERVER_HOST || 'http://localhost:8080/'
const TerminalId = import.meta.env.VITE_TERMINAL_ID || 'c2f8a06826baf03843925c1a2a13bcfd'
const robotStore = useRobotStore()
const videoUrl = `${VIDEO_SERVER_HOST}/api/edge/camera/rgb-stream/public?terminalId=${TerminalId}&deviceId=${robotStore.cameraDeviceId}`
// const videoUrl = computed(() => {
// if (!robotStore.ip || !robotStore.cameraDeviceId) return ''
// const params = new URLSearchParams({
// ip: robotStore.ip,
// deviceId: robotStore.cameraDeviceId
// })
// return `/api/camera/video?${params}`
// })
const videoUrl = computed(() => robotStore.robotId
? `/api/camera/video?robotId=${encodeURIComponent(robotStore.robotId)}`
: '')
</script>
<style lang="scss" scoped>
.page-container {

View File

@ -194,8 +194,6 @@ const run = async (data) => {
return
}
const res = await moveRobot({
ip: robotStore.ip,
deviceId: robotStore.agvDeviceId,
...data
})
if (res.code === 200) {
@ -210,8 +208,7 @@ function onRelease() {
moveRobotData.value.vx = 0;
moveRobotData.value.vy = 0;
// stopRobot({
// ip: robotStore.ip,
// deviceId: robotStore.agvDeviceId,
// robotId: robotStore.robotId,
// })
}

View File

@ -487,8 +487,6 @@ const loadSourceMap = async (mapName) => {
resizeCanvasToContainer()
const res = await getCurrentMap({
ip: robotStore.ip,
deviceId: robotStore.agvDeviceId,
mapName: mapName
})
if (res.code === 200) {
@ -540,7 +538,7 @@ watch(() => props.mapName,
async (newVal) => {
if (newVal) {
await loadSourceMap(newVal)
addRobot(robotStore.ip, robotStore.ip, robotStore.position.x, robotStore.position.y, robotStore.position.theta, getThemeColor('--color-canvas-robot'))
addRobot(robotStore.robotId, robotStore.robotId, robotStore.position.x, robotStore.position.y, robotStore.position.theta, getThemeColor('--color-canvas-robot'))
startRendering()
}
}, {