1use std::net::SocketAddr;
4use std::sync::Arc;
5use std::sync::atomic::{AtomicBool, Ordering};
6
7use bytes::{BufMut, BytesMut};
8use tokio::net::UdpSocket;
9use tokio::sync::{Mutex, Notify};
10use tracing::{debug, trace, warn};
11
12use crate::registers::RegisterMap;
13
14const GVCP_CMD_KEY: u8 = 0x42;
16
17const DISCOVERY_CMD: u16 = 0x0002;
19const FORCEIP_CMD: u16 = 0x0004;
20const FORCEIP_ACK: u16 = 0x0005;
21const READREG_CMD: u16 = 0x0080;
22const WRITEREG_CMD: u16 = 0x0082;
23const READMEM_CMD: u16 = 0x0084;
24const WRITEMEM_CMD: u16 = 0x0086;
25const ACTION_CMD: u16 = 0x0100;
29
30const DISCOVERY_ACK: u16 = 0x0003;
32const READREG_ACK: u16 = 0x0081;
33const WRITEREG_ACK: u16 = 0x0083;
34const READMEM_ACK: u16 = 0x0085;
35const WRITEMEM_ACK: u16 = 0x0087;
36const ACTION_ACK: u16 = 0x0101;
37
38pub const FAKE_DEVICE_KEY: u32 = 0x0000_0042;
42pub const FAKE_GROUP_KEY: u32 = 0x0000_0001;
44pub const FAKE_GROUP_MASK: u32 = 0x0000_0001;
46
47pub const FAKE_MAC: [u8; 6] = [0xDE, 0xAD, 0xBE, 0xEF, 0xCA, 0xFE];
50pub const FAKE_MANUFACTURER: &str = "viva-genicam";
52pub const FAKE_MODEL: &str = "FakeGigE";
54pub const FAKE_VERSION: &str = "1.0.0";
56pub const FAKE_SERIAL: &str = "FAKE-001";
58pub const FAKE_USER_NAME: &str = "FakeCamera";
60
61const STATUS_SUCCESS: u16 = 0x0000;
63const STATUS_INVALID_PARAMETER: u16 = 0x8002;
65
66pub async fn run(
71 socket: Arc<UdpSocket>,
72 regs: Arc<Mutex<RegisterMap>>,
73 acq_start_notify: Arc<Notify>,
74 acq_stop_flag: Arc<AtomicBool>,
75 bind_ip: std::net::Ipv4Addr,
76) {
77 let mut buf = [0u8; 2048];
78 loop {
79 let (len, peer) = match socket.recv_from(&mut buf).await {
80 Ok(r) => r,
81 Err(e) => {
82 warn!(error = %e, "GVCP recv error");
83 continue;
84 }
85 };
86 let pkt = &buf[..len];
87 if len < 8 || pkt[0] != GVCP_CMD_KEY {
88 trace!(len, "ignoring non-GVCP packet");
89 continue;
90 }
91
92 let flags = pkt[1];
93 let command = u16::from_be_bytes([pkt[2], pkt[3]]);
94 let _length = u16::from_be_bytes([pkt[4], pkt[5]]);
95 let request_id = u16::from_be_bytes([pkt[6], pkt[7]]);
96 let payload = &pkt[8..];
97
98 if matches!(
103 command,
104 READREG_CMD | WRITEREG_CMD | READMEM_CMD | WRITEMEM_CMD
105 ) && regs.lock().await.note_register_command()
106 {
107 warn!(%peer, "heartbeat expired; control privilege released");
108 }
109
110 match command {
111 DISCOVERY_CMD => {
112 let resp = build_discovery_ack(request_id, bind_ip);
113 let _ = socket.send_to(&resp, peer).await;
114 debug!(%peer, "discovery response sent");
115 }
116 FORCEIP_CMD => {
117 handle_forceip(&socket, peer, request_id, payload, bind_ip).await;
118 }
119 READREG_CMD => {
120 handle_readreg(&socket, peer, request_id, payload, ®s).await;
121 }
122 WRITEREG_CMD => {
123 handle_writereg(
124 &socket,
125 peer,
126 request_id,
127 payload,
128 ®s,
129 &acq_start_notify,
130 &acq_stop_flag,
131 )
132 .await;
133 }
134 READMEM_CMD => {
135 handle_readmem(&socket, peer, request_id, payload, ®s).await;
136 }
137 WRITEMEM_CMD => {
138 handle_writemem(
139 &socket,
140 peer,
141 request_id,
142 payload,
143 ®s,
144 &acq_start_notify,
145 &acq_stop_flag,
146 )
147 .await;
148 }
149 ACTION_CMD => {
150 handle_action(&socket, peer, request_id, flags, payload).await;
151 }
152 _ => {
153 debug!(command, "unsupported GVCP command");
154 }
155 }
156 }
157}
158
159fn build_discovery_ack(request_id: u16, ip: std::net::Ipv4Addr) -> Vec<u8> {
161 let payload_len: u16 = 248;
163 let mut buf = BytesMut::with_capacity(8 + payload_len as usize);
164
165 buf.put_u16(STATUS_SUCCESS);
167 buf.put_u16(DISCOVERY_ACK);
168 buf.put_u16(payload_len);
169 buf.put_u16(request_id);
170
171 buf.put_u16(2); buf.put_u16(0); buf.put_u32(0); buf.put_slice(&[0u8; 2]);
182
183 buf.put_slice(&FAKE_MAC);
185
186 buf.put_u32(0x0000_0007); buf.put_u32(0x0000_0005); buf.put_slice(&[0u8; 12]);
191
192 buf.put_slice(&ip.octets());
194
195 buf.put_slice(&[0u8; 12]);
197
198 buf.put_slice(&[255, 255, 255, 0]);
200
201 buf.put_slice(&[0u8; 12]);
203
204 buf.put_slice(&[0, 0, 0, 0]);
206
207 put_fixed_string(&mut buf, FAKE_MANUFACTURER, 32); put_fixed_string(&mut buf, FAKE_MODEL, 32); put_fixed_string(&mut buf, FAKE_VERSION, 32); put_fixed_string(&mut buf, "Fake camera for testing", 48); put_fixed_string(&mut buf, FAKE_SERIAL, 16); put_fixed_string(&mut buf, FAKE_USER_NAME, 16); buf.to_vec()
221}
222
223fn put_fixed_string(buf: &mut BytesMut, s: &str, len: usize) {
224 let bytes = s.as_bytes();
225 let copy_len = bytes.len().min(len);
226 buf.put_slice(&bytes[..copy_len]);
227 for _ in copy_len..len {
228 buf.put_u8(0);
229 }
230}
231
232fn build_ack(ack_cmd: u16, request_id: u16, payload: &[u8]) -> Vec<u8> {
234 let mut buf = BytesMut::with_capacity(8 + payload.len());
235 buf.put_u16(STATUS_SUCCESS);
236 buf.put_u16(ack_cmd);
237 buf.put_u16(payload.len() as u16);
238 buf.put_u16(request_id);
239 buf.put_slice(payload);
240 buf.to_vec()
241}
242
243fn build_error_ack(ack_cmd: u16, request_id: u16, status: u16) -> Vec<u8> {
245 let mut buf = BytesMut::with_capacity(8);
246 buf.put_u16(status);
247 buf.put_u16(ack_cmd);
248 buf.put_u16(0);
249 buf.put_u16(request_id);
250 buf.to_vec()
251}
252
253async fn handle_readreg(
254 socket: &UdpSocket,
255 peer: SocketAddr,
256 request_id: u16,
257 payload: &[u8],
258 regs: &Mutex<RegisterMap>,
259) {
260 if payload.len() < 4 || !payload.len().is_multiple_of(4) {
262 return;
263 }
264 let store = regs.lock().await;
265 let mut resp_payload = BytesMut::new();
266 for chunk in payload.chunks(4) {
267 let addr = u32::from_be_bytes([chunk[0], chunk[1], chunk[2], chunk[3]]) as u64;
268 let data = store.read(addr, 4);
269 resp_payload.put_slice(&data);
270 }
271 let resp = build_ack(READREG_ACK, request_id, &resp_payload);
272 let _ = socket.send_to(&resp, peer).await;
273 trace!(%peer, regs = payload.len() / 4, "READREG response");
274}
275
276async fn handle_writereg(
277 socket: &UdpSocket,
278 peer: SocketAddr,
279 request_id: u16,
280 payload: &[u8],
281 regs: &Mutex<RegisterMap>,
282 acq_start: &Notify,
283 acq_stop_flag: &AtomicBool,
284) {
285 if payload.len() < 8 || !payload.len().is_multiple_of(8) {
287 return;
288 }
289 let mut store = regs.lock().await;
290 for chunk in payload.chunks(8) {
291 let addr = u32::from_be_bytes([chunk[0], chunk[1], chunk[2], chunk[3]]) as u64;
292 let value = &chunk[4..8];
293 store.write(addr, value);
294 store.handle_special_write(addr);
295 check_acquisition(addr, value, acq_start, acq_stop_flag);
296 }
297 let resp = build_ack(WRITEREG_ACK, request_id, &[0, 0, 0, 0]);
299 let _ = socket.send_to(&resp, peer).await;
300 trace!(%peer, "WRITEREG response");
301}
302
303async fn handle_readmem(
304 socket: &UdpSocket,
305 peer: SocketAddr,
306 request_id: u16,
307 payload: &[u8],
308 regs: &Mutex<RegisterMap>,
309) {
310 if payload.len() < 8 {
312 return;
313 }
314 let addr = u32::from_be_bytes([payload[0], payload[1], payload[2], payload[3]]) as u64;
315 let count = u16::from_be_bytes([payload[6], payload[7]]) as usize;
316
317 if !addr.is_multiple_of(4) || count == 0 || !count.is_multiple_of(4) {
322 let resp = build_error_ack(READMEM_ACK, request_id, STATUS_INVALID_PARAMETER);
323 let _ = socket.send_to(&resp, peer).await;
324 debug!(%peer, addr = format!("0x{addr:x}"), count, "READMEM rejected: unaligned");
325 return;
326 }
327
328 let store = regs.lock().await;
329 let data = store.read(addr, count);
330
331 let mut resp_payload = BytesMut::with_capacity(4 + data.len());
333 resp_payload.put_u32(addr as u32);
334 resp_payload.put_slice(&data);
335 let resp = build_ack(READMEM_ACK, request_id, &resp_payload);
336 let _ = socket.send_to(&resp, peer).await;
337 trace!(%peer, addr = format!("0x{addr:x}"), count, "READMEM response");
338}
339
340async fn handle_writemem(
341 socket: &UdpSocket,
342 peer: SocketAddr,
343 request_id: u16,
344 payload: &[u8],
345 regs: &Mutex<RegisterMap>,
346 acq_start: &Notify,
347 acq_stop_flag: &AtomicBool,
348) {
349 if payload.len() < 4 {
351 return;
352 }
353 let addr = u32::from_be_bytes([payload[0], payload[1], payload[2], payload[3]]) as u64;
354 let data = &payload[4..];
355 let (test_packet, dest_ip, dest_port, max_on_wire) = {
356 let mut store = regs.lock().await;
357 store.write(addr, data);
358 store.handle_special_write(addr);
359 check_acquisition(addr, data, acq_start, acq_stop_flag);
360 (
361 store.take_pending_test_packet(),
362 store.stream_dest_ip(),
363 store.stream_dest_port(),
364 store.max_on_wire(),
365 )
366 };
367
368 let mut resp_payload = BytesMut::with_capacity(4);
370 resp_payload.put_u32(addr as u32);
371 let resp = build_ack(WRITEMEM_ACK, request_id, &resp_payload);
372 let _ = socket.send_to(&resp, peer).await;
373
374 if let Some(size) = test_packet {
378 fire_test_packet(size, dest_ip, dest_port, max_on_wire).await;
379 }
380 trace!(%peer, addr = format!("0x{addr:x}"), len = data.len(), "WRITEMEM response");
381}
382
383async fn fire_test_packet(
396 size: u32,
397 dest_ip: std::net::Ipv4Addr,
398 dest_port: u16,
399 max_on_wire: Option<u32>,
400) {
401 if dest_port == 0 {
402 debug!("test packet requested before a stream destination was set");
403 return;
404 }
405 if max_on_wire.is_some_and(|max| size > max) {
406 debug!(
407 size,
408 "test packet exceeds the path ceiling; dropped in flight"
409 );
410 return;
411 }
412
413 const IP_AND_UDP_HEADERS: u32 = 28;
414 let payload_len = size.saturating_sub(IP_AND_UDP_HEADERS) as usize;
415 let Ok(sock) = UdpSocket::bind("0.0.0.0:0").await else {
416 return;
417 };
418 let mut pkt = BytesMut::with_capacity(payload_len.max(8));
421 pkt.put_u16(0); pkt.put_u16(0); pkt.put_u8(0x03); pkt.put_u8(0);
425 pkt.put_u16(0); pkt.resize(payload_len.max(8), 0);
427 let _ = sock.send_to(&pkt, (dest_ip, dest_port)).await;
428 debug!(size, %dest_ip, dest_port, "test packet sent");
429}
430
431async fn handle_action(
438 socket: &UdpSocket,
439 peer: SocketAddr,
440 request_id: u16,
441 flags: u8,
442 payload: &[u8],
443) {
444 const SCHEDULED: u8 = 0x80;
445 const ACK_REQUIRED: u8 = 0x01;
446
447 let scheduled = flags & SCHEDULED != 0;
448 let expected = if scheduled { 20 } else { 12 };
449 if payload.len() < expected {
450 warn!(
451 len = payload.len(),
452 expected, scheduled, "malformed action command payload"
453 );
454 return;
455 }
456
457 let device_key = u32::from_be_bytes([payload[0], payload[1], payload[2], payload[3]]);
458 let group_key = u32::from_be_bytes([payload[4], payload[5], payload[6], payload[7]]);
459 let group_mask = u32::from_be_bytes([payload[8], payload[9], payload[10], payload[11]]);
460
461 if device_key != FAKE_DEVICE_KEY
462 || group_key != FAKE_GROUP_KEY
463 || group_mask & FAKE_GROUP_MASK == 0
464 {
465 debug!(
466 device_key,
467 group_key, group_mask, "action command not addressed to this device"
468 );
469 return;
470 }
471
472 debug!(%peer, scheduled, "action command accepted");
473 if flags & ACK_REQUIRED != 0 {
474 let mut buf = BytesMut::with_capacity(8);
475 buf.put_u16(STATUS_SUCCESS);
476 buf.put_u16(ACTION_ACK);
477 buf.put_u16(0);
478 buf.put_u16(request_id);
479 let _ = socket.send_to(&buf, peer).await;
480 }
481}
482
483async fn handle_forceip(
484 socket: &UdpSocket,
485 peer: SocketAddr,
486 request_id: u16,
487 payload: &[u8],
488 bind_ip: std::net::Ipv4Addr,
489) {
490 if payload.len() < 56 {
500 warn!(len = payload.len(), "FORCEIP payload too short");
501 return;
502 }
503
504 let target_mac = &payload[2..8];
505 let fake_mac: [u8; 6] = FAKE_MAC;
506 if target_mac != fake_mac {
507 debug!(
508 target = ?target_mac,
509 "FORCEIP: MAC mismatch, ignoring"
510 );
511 return;
512 }
513
514 let ip = std::net::Ipv4Addr::new(payload[20], payload[21], payload[22], payload[23]);
515 let subnet = std::net::Ipv4Addr::new(payload[36], payload[37], payload[38], payload[39]);
516 let gateway = std::net::Ipv4Addr::new(payload[52], payload[53], payload[54], payload[55]);
517
518 debug!(
519 %bind_ip,
520 %ip,
521 %subnet,
522 %gateway,
523 "FORCEIP accepted (fake camera ignores IP change)"
524 );
525
526 let resp = build_ack(FORCEIP_ACK, request_id, &[]);
528 let _ = socket.send_to(&resp, peer).await;
529}
530
531fn check_acquisition(addr: u64, data: &[u8], acq_start: &Notify, acq_stop_flag: &AtomicBool) {
533 use crate::registers::{REG_ACQ_START, REG_ACQ_STOP};
534
535 if addr == REG_ACQ_START && data.len() >= 4 {
536 let val = u32::from_be_bytes([data[0], data[1], data[2], data[3]]);
537 if val != 0 {
538 debug!("AcquisitionStart triggered");
539 acq_stop_flag.store(false, Ordering::SeqCst);
540 acq_start.notify_one();
541 }
542 } else if addr == REG_ACQ_STOP && data.len() >= 4 {
543 let val = u32::from_be_bytes([data[0], data[1], data[2], data[3]]);
544 if val != 0 {
545 debug!("AcquisitionStop triggered");
546 acq_stop_flag.store(true, Ordering::SeqCst);
547 }
548 }
549}