1use std::net::{IpAddr, Ipv4Addr, SocketAddr};
7use std::sync::Arc;
8use std::sync::atomic::{AtomicU16, AtomicU32, Ordering};
9use std::time::{Duration, Instant, SystemTime};
10
11use regex::Regex;
12use tokio::io::{AsyncReadExt, AsyncWriteExt};
13use tokio::net::{TcpListener, TcpStream, UdpSocket};
14use tokio::sync::mpsc;
15use tracing::{debug, error, info, warn};
16
17use spvirit_codec::epics_decode::{PvaHeader, PvaPacket, PvaPacketCommand};
18use spvirit_codec::spvd_decode::{PvdDecoder, StructureDesc, extract_subfield_desc};
19use spvirit_codec::spvd_encode::{
20 decode_pv_request_fields, filter_structure_desc, nt_payload_desc,
21};
22use spvirit_codec::spvirit_encode::{
23 encode_connection_validation, encode_control_message, encode_create_channel_error,
24 encode_create_channel_response, encode_get_field_error, encode_get_field_response,
25 encode_header, encode_message_error, encode_monitor_data_response_payload,
26 encode_op_data_response_filtered, encode_op_error, encode_op_get_data_response_payload,
27 encode_op_init_response_desc, encode_op_put_get_data_error_response,
28 encode_op_put_get_data_response_payload, encode_op_put_get_init_error_response,
29 encode_op_put_get_init_response, encode_op_put_getput_response_payload, encode_op_put_response,
30 encode_op_put_status_response, encode_op_rpc_data_response_payload,
31 encode_op_status_error_response, encode_op_status_response, encode_search_response,
32 ip_from_bytes, ip_to_bytes,
33};
34
35use spvirit_codec::{SegmentOutcome, SegmentReassembler};
36
37use spvirit_types::{NtPayload, NtScalar, NtScalarArray, ScalarArrayValue, ScalarValue};
38
39use crate::decode::decode_put_body;
40use crate::monitor::MonitorRegistry;
41use crate::pvstore::SourceRegistry;
42use crate::state::{ConnState, MonitorState, MonitorSub};
43
44#[derive(Debug, Clone, Copy, PartialEq, Eq)]
50pub enum PvListMode {
51 Off,
53 Discover,
55 List,
57}
58
59impl PvListMode {
60 pub fn parse(raw: &str) -> Result<Self, String> {
61 match raw.trim().to_ascii_lowercase().as_str() {
62 "off" => Ok(Self::Off),
63 "discover" => Ok(Self::Discover),
64 "list" => Ok(Self::List),
65 other => Err(format!(
66 "Invalid pvlist-mode '{}'; expected off|discover|list",
67 other
68 )),
69 }
70 }
71}
72
73pub struct ServerState {
79 pub sources: Arc<SourceRegistry>,
80 pub registry: Arc<MonitorRegistry>,
81 pub sid_counter: AtomicU32,
82 pub beacon_change: Arc<AtomicU16>,
83 pub compute_alarms: bool,
84 pub pvlist_mode: PvListMode,
85 pub pvlist_max: usize,
86 pub pvlist_allow_pattern: Option<Regex>,
87 pub guid: [u8; 12],
88 pub tcp_port: u16,
89 pub advertise_ip: Option<IpAddr>,
90 pub listen_ip: IpAddr,
91}
92
93impl ServerState {
94 pub fn new(
95 sources: Arc<SourceRegistry>,
96 registry: Arc<MonitorRegistry>,
97 compute_alarms: bool,
98 pvlist_mode: PvListMode,
99 pvlist_max: usize,
100 pvlist_allow_pattern: Option<Regex>,
101 guid: [u8; 12],
102 tcp_port: u16,
103 advertise_ip: Option<IpAddr>,
104 listen_ip: IpAddr,
105 ) -> Self {
106 Self {
107 sources,
108 registry,
109 sid_counter: AtomicU32::new(1),
110 beacon_change: Arc::new(AtomicU16::new(0)),
111 compute_alarms,
112 pvlist_mode,
113 pvlist_max,
114 pvlist_allow_pattern,
115 guid,
116 tcp_port,
117 advertise_ip,
118 listen_ip,
119 }
120 }
121}
122
123pub fn is_pvlist_virtual_pv(pv_name: &str) -> bool {
128 pv_name == "__pvlist"
129}
130
131pub fn is_server_rpc_pv(pv_name: &str) -> bool {
132 pv_name == "server"
133}
134
135pub fn is_virtual_event_pv(pv_name: &str) -> bool {
136 pv_name.starts_with("__event:")
137}
138
139pub fn virtual_event_nt(pv_name: &str) -> NtPayload {
140 NtPayload::Scalar(
141 NtScalar::from_value(ScalarValue::Bool(false))
142 .with_description(format!("Virtual event trigger for {}", pv_name)),
143 )
144}
145
146pub fn virtual_pvlist_nt(entries: Vec<String>) -> NtPayload {
147 NtPayload::ScalarArray(NtScalarArray::from_value(ScalarArrayValue::Str(entries)))
148}
149
150pub fn is_pattern_query(raw: &str) -> bool {
155 raw.contains('*') || raw.contains('?')
156}
157
158pub fn wildcard_match(pattern: &str, text: &str) -> bool {
159 let p = pattern.as_bytes();
160 let t = text.as_bytes();
161 let mut i = 0usize;
162 let mut j = 0usize;
163 let mut star: Option<usize> = None;
164 let mut match_j = 0usize;
165
166 while j < t.len() {
167 if i < p.len() && (p[i] == b'?' || p[i] == t[j]) {
168 i += 1;
169 j += 1;
170 } else if i < p.len() && p[i] == b'*' {
171 star = Some(i);
172 i += 1;
173 match_j = j;
174 } else if let Some(star_idx) = star {
175 i = star_idx + 1;
176 match_j += 1;
177 j = match_j;
178 } else {
179 return false;
180 }
181 }
182
183 while i < p.len() && p[i] == b'*' {
184 i += 1;
185 }
186 i == p.len()
187}
188
189pub fn collect_visible_pv_names(
190 all_names: &[String],
191 mode: PvListMode,
192 allow_pattern: Option<&Regex>,
193 max_items: usize,
194) -> Vec<String> {
195 let mut names: Vec<String> = all_names
196 .iter()
197 .filter(|name| {
198 allow_pattern
199 .as_ref()
200 .map(|re| re.is_match(name))
201 .unwrap_or(true)
202 })
203 .cloned()
204 .collect();
205 names.sort();
206 if names.len() > max_items {
207 names.truncate(max_items);
208 }
209 if mode == PvListMode::List && names.len() < max_items {
210 names.push("__pvlist".to_string());
211 }
212 names
213}
214
215fn build_pvlist_structure(names: &[String]) -> StructureDesc {
216 use spvirit_codec::spvd_decode::{FieldDesc, FieldType, TypeCode};
217 StructureDesc {
218 struct_id: Some("epics:pva/pvlist:1.0".to_string()),
219 fields: names
220 .iter()
221 .map(|name| FieldDesc {
222 name: name.clone(),
223 field_type: FieldType::Scalar(TypeCode::Boolean),
224 })
225 .collect(),
226 }
227}
228
229fn requested_pvlist_pattern(field_name: Option<&str>) -> Option<&str> {
230 let raw = field_name.map(str::trim).unwrap_or("");
231 if raw.is_empty() || raw == "*" || raw == "__pvlist" || raw.eq_ignore_ascii_case("pvlist") {
232 return Some("*");
233 }
234 if is_pattern_query(raw) {
235 return Some(raw);
236 }
237 None
238}
239
240pub fn search_reply_target(addr: &[u8; 16], port: u16, peer: SocketAddr) -> SocketAddr {
245 let target_port = if port != 0 { port } else { peer.port() };
246 let target_ip = ip_from_bytes(addr)
247 .filter(|ip| !ip.is_unspecified())
248 .unwrap_or_else(|| peer.ip());
249 SocketAddr::new(target_ip, target_port)
250}
251
252pub fn bind_udp_search_socket(addr: SocketAddr) -> std::io::Result<UdpSocket> {
258 use socket2::{Domain, Protocol, Socket, Type};
259
260 let domain = if addr.is_ipv4() {
261 Domain::IPV4
262 } else {
263 Domain::IPV6
264 };
265 let socket = Socket::new(domain, Type::DGRAM, Some(Protocol::UDP))?;
266 socket.set_reuse_address(true)?;
267 #[cfg(unix)]
268 socket.set_reuse_port(true)?;
269 socket.set_nonblocking(true)?;
270 socket.bind(&addr.into())?;
271 UdpSocket::from_std(socket.into())
272}
273
274pub fn infer_udp_response_ip(peer: SocketAddr) -> Option<IpAddr> {
275 let bind_addr = if peer.is_ipv4() {
276 "0.0.0.0:0"
277 } else {
278 "[::]:0"
279 };
280 let sock = std::net::UdpSocket::bind(bind_addr).ok()?;
281 sock.connect(peer).ok()?;
282 let local = sock.local_addr().ok()?;
283 if local.ip().is_unspecified() {
284 None
285 } else {
286 Some(local.ip())
287 }
288}
289
290pub fn rand_guid() -> [u8; 12] {
291 let pid = std::process::id().to_le_bytes();
292 let nanos = SystemTime::now()
293 .duration_since(SystemTime::UNIX_EPOCH)
294 .unwrap_or_default()
295 .as_nanos()
296 .to_le_bytes();
297 let mut guid = [0u8; 12];
298 guid[..4].copy_from_slice(&pid);
299 guid[4..12].copy_from_slice(&nanos[..8]);
300 guid
301}
302
303pub fn validate_encoded_packet(conn_id: u64, label: &str, bytes: &[u8]) {
308 let mut pkt = PvaPacket::new(bytes);
309 let decoded = pkt.decode_payload();
310 match decoded {
311 Some(PvaPacketCommand::ConnectionValidation(payload)) => {
312 debug!(
313 "Conn {}: {} decoded as cmd=1 buffer_size={} qos={} authz={:?}",
314 conn_id, label, payload.buffer_size, payload.qos, payload.authz
315 );
316 }
317 Some(PvaPacketCommand::ConnectionValidated(_)) => {
318 debug!("Conn {}: {} decoded as cmd=9", conn_id, label);
319 }
320 Some(other) => {
321 debug!("Conn {}: {} decoded as {:?}", conn_id, label, other);
322 }
323 None => {
324 debug!("Conn {}: {} failed to decode", conn_id, label);
325 }
326 }
327}
328
329pub fn dump_hex_packet(
330 conn_id: u64,
331 dir: &str,
332 label: &str,
333 version: u8,
334 is_be: bool,
335 bytes: &[u8],
336) {
337 debug!(
338 "Conn {}: {} {} ver={} be={} len={}",
339 conn_id,
340 dir,
341 label,
342 version,
343 is_be,
344 bytes.len()
345 );
346 let mut offset = 0usize;
347 while offset < bytes.len() {
348 let end = usize::min(offset + 16, bytes.len());
349 let chunk = &bytes[offset..end];
350 let mut line = String::new();
351 for (i, b) in chunk.iter().enumerate() {
352 if i > 0 {
353 line.push(' ');
354 }
355 line.push_str(&format!("{:02x}", b));
356 }
357 debug!("Conn {}: {:04x} {}", conn_id, offset, line);
358 offset += 16;
359 }
360}
361
362async fn get_nt_snapshot(state: &ServerState, pv_name: &str) -> Option<NtPayload> {
367 if is_pvlist_virtual_pv(pv_name) {
368 if state.pvlist_mode != PvListMode::List {
369 return None;
370 }
371 let all_names = state.sources.names().await;
372 let names = collect_visible_pv_names(
373 &all_names,
374 state.pvlist_mode,
375 state.pvlist_allow_pattern.as_ref(),
376 state.pvlist_max,
377 );
378 return Some(virtual_pvlist_nt(names));
379 }
380 if is_virtual_event_pv(pv_name) {
381 return Some(virtual_event_nt(pv_name));
382 }
383 state.sources.get(pv_name).await
384}
385
386async fn is_writable_pv(state: &ServerState, pv_name: &str) -> bool {
387 if is_virtual_event_pv(pv_name) {
388 return true;
389 }
390 state.sources.is_writable(pv_name).await
391}
392
393async fn has_pv(state: &ServerState, pv_name: &str) -> bool {
394 state.sources.has_pv(pv_name).await
395 || is_virtual_event_pv(pv_name)
396 || (is_pvlist_virtual_pv(pv_name) && state.pvlist_mode == PvListMode::List)
397 || (is_server_rpc_pv(pv_name) && state.pvlist_mode != PvListMode::Off)
398}
399
400async fn notify_changed_records(state: &ServerState, changed: Vec<(String, NtPayload)>) {
405 for (name, payload) in changed {
406 state.beacon_change.fetch_add(1, Ordering::SeqCst);
407 state.registry.notify_monitors(&name, &payload).await;
408 }
409}
410
411async fn handle_get_field_request(
416 state: &ServerState,
417 conn_state: &ConnState,
418 conn_id: u64,
419 payload: spvirit_codec::epics_decode::PvaGetFieldPayload,
420 version: u8,
421 is_be: bool,
422) {
423 if payload.is_server {
424 let resp = encode_get_field_error(
425 payload.cid,
426 "Unexpected server GET_FIELD payload",
427 version,
428 is_be,
429 );
430 state.registry.send_msg(conn_id, resp).await;
431 return;
432 }
433
434 let request_id = payload.ioid.unwrap_or(payload.cid);
435
436 let sid = payload
437 .sid
438 .or_else(|| conn_state.cid_to_sid.get(&payload.cid).copied())
439 .or_else(|| {
440 conn_state
441 .sid_to_pv
442 .contains_key(&payload.cid)
443 .then_some(payload.cid)
444 })
445 .or_else(|| {
446 (payload.cid == 0 && conn_state.sid_to_pv.len() == 1)
447 .then(|| conn_state.sid_to_pv.keys().copied().next())
448 .flatten()
449 });
450
451 if let Some(sid) = sid {
452 if let Some(pv_name) = conn_state.sid_to_pv.get(&sid) {
453 if let Some(nt) = get_nt_snapshot(state, pv_name).await {
454 let full_desc = nt_payload_desc(&nt);
455 let sub = payload.field_name.as_deref().filter(|s| !s.is_empty());
456 let desc = if let Some(field_path) = sub {
457 match extract_subfield_desc(&full_desc, field_path) {
458 Some(sub_desc) => sub_desc,
459 None => {
460 let resp = encode_get_field_error(
461 request_id,
462 &format!("sub-field '{}' not found", field_path),
463 version,
464 is_be,
465 );
466 state.registry.send_msg(conn_id, resp).await;
467 return;
468 }
469 }
470 } else {
471 full_desc
472 };
473 let resp = encode_get_field_response(request_id, &desc, version, is_be);
474 dump_hex_packet(conn_id, "tx", "cmd=17 get_field", version, is_be, &resp);
475 state.registry.send_msg(conn_id, resp).await;
476 debug!(
477 "Conn {}: get_field cid={} sid={:?} ioid={:?} resolved_sid={} pv='{}' field={:?}",
478 conn_id,
479 payload.cid,
480 payload.sid,
481 payload.ioid,
482 sid,
483 pv_name,
484 payload.field_name
485 );
486 return;
487 }
488 let resp = encode_get_field_error(request_id, "PV not found", version, is_be);
489 state.registry.send_msg(conn_id, resp).await;
490 return;
491 }
492 }
493
494 if state.pvlist_mode != PvListMode::List {
495 let resp = encode_get_field_error(
496 request_id,
497 "GET_FIELD listing is disabled (set --pvlist-mode=list)",
498 version,
499 is_be,
500 );
501 state.registry.send_msg(conn_id, resp).await;
502 return;
503 }
504
505 let Some(pattern) = requested_pvlist_pattern(payload.field_name.as_deref()) else {
506 let resp = encode_get_field_error(
507 request_id,
508 "GET_FIELD requires a valid list pattern",
509 version,
510 is_be,
511 );
512 state.registry.send_msg(conn_id, resp).await;
513 return;
514 };
515
516 let all_names = state.sources.names().await;
517 let mut names = collect_visible_pv_names(
518 &all_names,
519 state.pvlist_mode,
520 state.pvlist_allow_pattern.as_ref(),
521 state.pvlist_max,
522 );
523 if pattern != "*" {
524 names.retain(|name| wildcard_match(pattern, name));
525 }
526 if names.is_empty() {
527 let resp =
528 encode_get_field_error(request_id, "No PVs matched list request", version, is_be);
529 state.registry.send_msg(conn_id, resp).await;
530 return;
531 }
532 let desc = build_pvlist_structure(&names);
533 let resp = encode_get_field_response(request_id, &desc, version, is_be);
534 dump_hex_packet(
535 conn_id,
536 "tx",
537 "cmd=17 get_field_list",
538 version,
539 is_be,
540 &resp,
541 );
542 state.registry.send_msg(conn_id, resp).await;
543 debug!(
544 "Conn {}: get_field list pattern='{}' returned {} entries",
545 conn_id,
546 pattern,
547 names.len()
548 );
549}
550
551async fn handle_server_rpc(
556 state: &ServerState,
557 conn_id: u64,
558 ioid: u32,
559 subcmd: u8,
560 version: u8,
561 is_be: bool,
562) {
563 if state.pvlist_mode != PvListMode::List {
564 let resp = encode_op_status_error_response(
565 20,
566 ioid,
567 subcmd,
568 "RPC list endpoint disabled (set --pvlist-mode=list)",
569 version,
570 is_be,
571 );
572 state.registry.send_msg(conn_id, resp).await;
573 return;
574 }
575
576 let all_names = state.sources.names().await;
577 let names = collect_visible_pv_names(
578 &all_names,
579 state.pvlist_mode,
580 state.pvlist_allow_pattern.as_ref(),
581 state.pvlist_max,
582 );
583 let payload = NtPayload::ScalarArray(NtScalarArray::from_value(ScalarArrayValue::Str(names)));
584
585 let is_init = (subcmd & 0x08) != 0;
586 if is_init {
587 let resp = encode_op_status_response(20, ioid, subcmd, version, is_be);
588 state.registry.send_msg(conn_id, resp).await;
589 return;
590 }
591
592 let resp = encode_op_rpc_data_response_payload(ioid, subcmd, &payload, version, is_be);
593 state.registry.send_msg(conn_id, resp).await;
594}
595
596async fn handle_control_message(state: &ServerState, conn_id: u64, header: &PvaHeader) {
601 debug!(
602 "Conn {}: control (segmented) cmd={} data={}",
603 conn_id, header.command, header.payload_length
604 );
605 if header.command == 3 {
606 let resp = encode_control_message(
607 true,
608 header.flags.is_msb,
609 header.version,
610 4,
611 header.payload_length,
612 );
613 state.registry.send_msg(conn_id, resp).await;
614 }
615}
616
617pub async fn run_udp_search(
623 state: Arc<ServerState>,
624 addr: SocketAddr,
625 tcp_port: u16,
626 guid: [u8; 12],
627 advertise_ip: Option<IpAddr>,
628) -> Result<(), Box<dyn std::error::Error>> {
629 let socket = bind_udp_search_socket(addr)?;
630 socket.set_broadcast(true)?;
631 let mut buf = vec![0u8; 4096];
632
633 loop {
634 let (len, peer) = socket.recv_from(&mut buf).await?;
635 let data = &buf[..len];
636 let header = PvaHeader::new(data);
637 if header.flags.is_control || header.command != 3 {
638 continue;
639 }
640 let mut pkt = PvaPacket::new(data);
641 let Some(cmd) = pkt.decode_payload() else {
642 continue;
643 };
644 let version = pkt.header.version;
645 let is_be = pkt.header.flags.is_msb;
646 match cmd {
647 PvaPacketCommand::Search(payload) => {
648 debug!(
649 "UDP search from {}: pv_count={} mask=0x{:02x}",
650 peer,
651 payload.pv_requests.len(),
652 payload.mask
653 );
654 let accepts_tcp = payload.protocols.is_empty()
655 || payload
656 .protocols
657 .iter()
658 .any(|p| p.eq_ignore_ascii_case("tcp"));
659 if !accepts_tcp {
660 debug!("UDP search: no compatible protocol (tcp not accepted)");
661 continue;
662 }
663 let all_names = state.sources.names().await;
664 let visible_names = collect_visible_pv_names(
665 &all_names,
666 state.pvlist_mode,
667 state.pvlist_allow_pattern.as_ref(),
668 state.pvlist_max,
669 );
670 let mut cids = Vec::new();
671 for (cid, name) in &payload.pv_requests {
672 if state.sources.has_pv(name).await
673 || is_virtual_event_pv(name)
674 || (is_pvlist_virtual_pv(name) && state.pvlist_mode == PvListMode::List)
675 || (is_server_rpc_pv(name) && state.pvlist_mode != PvListMode::Off)
676 {
677 cids.push(*cid);
678 continue;
679 }
680 if state.pvlist_mode != PvListMode::Off
681 && is_pattern_query(name)
682 && visible_names.iter().any(|pv| wildcard_match(name, pv))
683 {
684 cids.push(*cid);
685 }
686 }
687 let response_required = (payload.mask & 0x01) != 0;
688 let server_discovery_ping = payload.pv_requests.is_empty();
689 let found = server_discovery_ping || !cids.is_empty();
690 if !found && !response_required {
691 debug!("UDP search: no matches and response not required");
692 continue;
693 }
694 let resp_ip = if let Some(ip) = advertise_ip {
695 ip
696 } else if !addr.ip().is_unspecified() {
697 addr.ip()
698 } else if let Some(ip) = infer_udp_response_ip(peer) {
699 debug!("UDP search: inferred response address {}", ip);
700 ip
701 } else {
702 IpAddr::V4(Ipv4Addr::UNSPECIFIED)
703 };
704 let addr_bytes = if resp_ip.is_unspecified() {
705 debug!("UDP search: responding with zero address (unspecified listen)");
706 [0u8; 16]
707 } else {
708 ip_to_bytes(resp_ip)
709 };
710 let response = encode_search_response(
711 guid,
712 payload.seq,
713 addr_bytes,
714 tcp_port,
715 "tcp",
716 found,
717 &cids,
718 version,
719 is_be,
720 );
721 let reply_target = search_reply_target(&payload.addr, payload.port, peer);
722 if let Err(e) = socket.send_to(&response, reply_target).await {
723 debug!(
724 "UDP search: failed sending {} matches to {}: {}",
725 cids.len(),
726 reply_target,
727 e
728 );
729 continue;
730 }
731 debug!(
732 "UDP search: responded found={} with {} matches to {}",
733 found,
734 cids.len(),
735 reply_target
736 );
737 }
738 _ => {}
739 }
740 }
741}
742
743pub async fn run_tcp_server(
752 state: Arc<ServerState>,
753 listener: TcpListener,
754 conn_timeout: Duration,
755) -> Result<(), Box<dyn std::error::Error>> {
756 let conn_id = Arc::new(std::sync::atomic::AtomicU64::new(1));
757
758 loop {
759 let (stream, peer) = listener.accept().await?;
760 let id = conn_id.fetch_add(1, Ordering::SeqCst);
761 info!("TCP connection {} from {}", id, peer);
762 let state_clone = state.clone();
763 tokio::spawn(async move {
764 if let Err(e) = handle_connection(state_clone, stream, id, conn_timeout).await {
765 error!("Connection {} error: {}", id, e);
766 }
767 });
768 }
769}
770
771pub async fn handle_connection(
781 state: Arc<ServerState>,
782 stream: TcpStream,
783 conn_id: u64,
784 conn_timeout: Duration,
785) -> Result<(), Box<dyn std::error::Error>> {
786 let (mut reader, mut writer) = stream.into_split();
787 let (tx, mut rx) = mpsc::channel::<Vec<u8>>(128);
788
789 {
790 let mut conns = state.registry.conns.lock().await;
791 conns.insert(conn_id, tx);
792 }
793
794 let writer_task = tokio::spawn(async move {
795 while let Some(msg) = rx.recv().await {
796 if writer.write_all(&msg).await.is_err() {
797 break;
798 }
799 }
800 });
801
802 let mut conn_state = ConnState::default();
803
804 let set_byte_order = encode_control_message(true, false, 2, 2, 0);
806 validate_encoded_packet(conn_id, "set_byte_order", &set_byte_order);
807 dump_hex_packet(
808 conn_id,
809 "tx",
810 "ctrl=2 set_byte_order",
811 2,
812 false,
813 &set_byte_order,
814 );
815 state.registry.send_msg(conn_id, set_byte_order).await;
816
817 let server_validation =
819 encode_connection_validation(16_384, 512, &["anonymous", "ca"], 2, false);
820 validate_encoded_packet(conn_id, "server_validation_init", &server_validation);
821 dump_hex_packet(
822 conn_id,
823 "tx",
824 "cmd=1 server_validation_init",
825 2,
826 false,
827 &server_validation,
828 );
829 state.registry.send_msg(conn_id, server_validation).await;
830
831 let mut last_activity = Instant::now();
832 let mut reassembler = SegmentReassembler::new();
835
836 loop {
837 let mut header = [0u8; 8];
838 let elapsed = last_activity.elapsed();
839 if elapsed >= conn_timeout {
840 info!("Conn {} idle timeout", conn_id);
841 break;
842 }
843 let remaining = conn_timeout - elapsed;
844 let read_header = tokio::time::timeout(remaining, reader.read_exact(&mut header)).await;
845 match read_header {
846 Ok(Ok(_)) => {}
847 Ok(Err(_)) => break,
848 Err(_) => {
849 info!("Conn {} idle timeout", conn_id);
850 break;
851 }
852 }
853 let header_pkt = PvaPacket::new(&header);
854 let payload_len = if header_pkt.header.flags.is_control {
855 0usize
856 } else {
857 header_pkt.header.payload_length as usize
858 };
859 let mut payload = vec![0u8; payload_len];
860 if payload_len > 0 {
861 let elapsed = last_activity.elapsed();
862 if elapsed >= conn_timeout {
863 info!("Conn {} idle timeout", conn_id);
864 break;
865 }
866 let remaining = conn_timeout - elapsed;
867 let read_payload =
868 tokio::time::timeout(remaining, reader.read_exact(&mut payload)).await;
869 match read_payload {
870 Ok(Ok(_)) => {}
871 Ok(Err(_)) => break,
872 Err(_) => {
873 info!("Conn {} idle timeout", conn_id);
874 break;
875 }
876 }
877 }
878 last_activity = Instant::now();
879
880 let full = match reassembler.push(header, payload) {
884 Ok(SegmentOutcome::Complete(msg)) => msg,
885 Ok(SegmentOutcome::Pending) => continue,
886 Ok(SegmentOutcome::Control(msg)) => {
887 let ctrl = PvaPacket::new(&msg);
888 handle_control_message(&state, conn_id, &ctrl.header).await;
889 continue;
890 }
891 Err(e) => {
892 warn!("Conn {}: segmentation error: {}", conn_id, e);
893 break;
894 }
895 };
896
897 let mut pkt = PvaPacket::new(&full);
898 let Some(cmd) = pkt.decode_payload() else {
899 continue;
900 };
901 let version = pkt.header.version;
902 let is_be = pkt.header.flags.is_msb;
903 let cmd_code = pkt.header.command;
904 let payload_slice = if full.len() >= 8 { &full[8..] } else { &[] };
905
906 if cmd_code == 1 {
908 dump_hex_packet(conn_id, "rx", "cmd=1 validation", version, is_be, &full);
909 let validation = spvirit_codec::epics_decode::PvaConnectionValidationPayload::new(
910 payload_slice,
911 is_be,
912 false,
913 );
914 if let Some(val) = validation {
915 debug!(
916 "Conn {}: validation request (cmd=1) ver={} be={} buf={} qos={} authz={:?}",
917 conn_id, version, is_be, val.buffer_size, val.qos, val.authz
918 );
919 let resp = spvirit_codec::spvirit_encode::encode_connection_validated(
920 true, version, is_be,
921 );
922 validate_encoded_packet(conn_id, "conn_validated_resp", &resp);
923 dump_hex_packet(
924 conn_id,
925 "tx",
926 "cmd=9 connection_validated",
927 version,
928 is_be,
929 &resp,
930 );
931 state.registry.send_msg(conn_id, resp).await;
932 continue;
933 }
934 }
935 if cmd_code == 17 {
936 dump_hex_packet(conn_id, "rx", "cmd=17 get_field", version, is_be, &full);
937 }
938
939 match cmd {
940 PvaPacketCommand::Control(payload) => {
941 debug!("Conn {}: control {}", conn_id, payload);
942 if payload.command == 3 {
943 let resp = encode_control_message(true, is_be, version, 4, payload.data);
944 state.registry.send_msg(conn_id, resp).await;
945 }
946 continue;
947 }
948 PvaPacketCommand::ConnectionValidation(_) => {
949 debug!("Conn {}: validation request (decoded)", conn_id);
950 }
951 PvaPacketCommand::ConnectionValidated(_) => {
952 debug!("Conn {}: validation confirmed (decoded)", conn_id);
953 }
954 PvaPacketCommand::CreateChannel(payload) => {
955 debug!(
956 "Conn {}: create_channel count={}",
957 conn_id,
958 payload.channels.len()
959 );
960 for (cid, pv_name) in payload.channels {
961 if has_pv(&state, &pv_name).await {
962 let sid = state.sid_counter.fetch_add(1, Ordering::SeqCst);
963 conn_state.cid_to_sid.insert(cid, sid);
964 conn_state.sid_to_pv.insert(sid, pv_name.clone());
965 let resp = encode_create_channel_response(cid, sid, version, is_be);
966 state.registry.send_msg(conn_id, resp).await;
967 info!(
968 "Conn {}: channel '{}' cid={} sid={}",
969 conn_id, pv_name, cid, sid
970 );
971 } else {
972 let resp = encode_create_channel_error(cid, "PV not found", version, is_be);
973 state.registry.send_msg(conn_id, resp).await;
974 info!(
975 "Conn {}: channel '{}' not found (cid={})",
976 conn_id, pv_name, cid
977 );
978 }
979 }
980 }
981 PvaPacketCommand::Op(payload) => {
982 if payload.is_server {
983 continue;
984 }
985 let sid = payload.sid_or_cid;
986 let ioid = payload.ioid;
987 debug!(
988 "Conn {}: op cmd={} ioid={} sid={} sub=0x{:02x} body_len={}",
989 conn_id,
990 payload.command,
991 ioid,
992 sid,
993 payload.subcmd,
994 payload.body.len()
995 );
996 let Some(pv_name) = conn_state.sid_to_pv.get(&sid).cloned() else {
997 state
998 .registry
999 .send_msg(
1000 conn_id,
1001 encode_op_error(
1002 payload.command,
1003 payload.subcmd,
1004 ioid,
1005 "Unknown SID",
1006 version,
1007 is_be,
1008 ),
1009 )
1010 .await;
1011 continue;
1012 };
1013
1014 let is_init = (payload.subcmd & 0x08) != 0;
1015
1016 match payload.command {
1017 10 => {
1018 if is_init {
1020 let full_desc =
1023 if let Some(desc) = state.sources.get_descriptor(&pv_name).await {
1024 desc
1025 } else if let Some(nt) = get_nt_snapshot(&state, &pv_name).await {
1026 nt_payload_desc(&nt)
1027 } else {
1028 state
1029 .registry
1030 .send_msg(
1031 conn_id,
1032 encode_op_error(
1033 payload.command,
1034 payload.subcmd,
1035 ioid,
1036 "PV not found",
1037 version,
1038 is_be,
1039 ),
1040 )
1041 .await;
1042 continue;
1043 };
1044 let pv_req_fields = decode_pv_request_fields(&payload.body, is_be);
1045 let desc = match &pv_req_fields {
1046 Some(fields) => filter_structure_desc(&full_desc, fields),
1047 None => full_desc,
1048 };
1049 conn_state.ioid_to_desc.insert(ioid, desc.clone());
1050 conn_state.ioid_to_pv.insert(ioid, pv_name.clone());
1051 let resp = encode_op_init_response_desc(
1052 payload.command,
1053 ioid,
1054 0x08,
1055 &desc,
1056 version,
1057 is_be,
1058 );
1059 state.registry.send_msg(conn_id, resp).await;
1060 info!("Conn {}: get init pv='{}' ioid={}", conn_id, pv_name, ioid);
1061 } else {
1062 let Some(nt) = get_nt_snapshot(&state, &pv_name).await else {
1063 state
1064 .registry
1065 .send_msg(
1066 conn_id,
1067 encode_op_error(
1068 payload.command,
1069 payload.subcmd,
1070 ioid,
1071 "PV has no data yet",
1072 version,
1073 is_be,
1074 ),
1075 )
1076 .await;
1077 continue;
1078 };
1079 let resp = if let Some(desc) = conn_state.ioid_to_desc.get(&ioid) {
1080 encode_op_data_response_filtered(
1081 10, ioid, &nt, desc, version, is_be,
1082 )
1083 } else {
1084 encode_op_get_data_response_payload(ioid, &nt, version, is_be)
1085 };
1086 state.registry.send_msg(conn_id, resp).await;
1087 debug!("Conn {}: get data pv='{}' ioid={}", conn_id, pv_name, ioid);
1088 }
1089 }
1090 11 => {
1091 if is_init {
1093 let desc =
1097 if let Some(desc) = state.sources.get_descriptor(&pv_name).await {
1098 desc
1099 } else if let Some(nt) = get_nt_snapshot(&state, &pv_name).await {
1100 nt_payload_desc(&nt)
1101 } else {
1102 state
1103 .registry
1104 .send_msg(
1105 conn_id,
1106 encode_op_error(
1107 payload.command,
1108 payload.subcmd,
1109 ioid,
1110 "PV not found",
1111 version,
1112 is_be,
1113 ),
1114 )
1115 .await;
1116 continue;
1117 };
1118 if !is_virtual_event_pv(&pv_name)
1119 && !is_writable_pv(&state, &pv_name).await
1120 {
1121 let resp = encode_op_put_status_response(
1122 ioid,
1123 0x08,
1124 "Write access denied",
1125 version,
1126 is_be,
1127 );
1128 state.registry.send_msg(conn_id, resp).await;
1129 continue;
1130 }
1131 conn_state.ioid_to_desc.insert(ioid, desc.clone());
1132 conn_state.ioid_to_pv.insert(ioid, pv_name.clone());
1133 let resp = encode_op_init_response_desc(
1134 payload.command,
1135 ioid,
1136 0x08,
1137 &desc,
1138 version,
1139 is_be,
1140 );
1141 state.registry.send_msg(conn_id, resp).await;
1142 info!("Conn {}: put init pv='{}' ioid={}", conn_id, pv_name, ioid);
1143 } else {
1144 if (payload.subcmd & 0x40) != 0 {
1145 if !is_virtual_event_pv(&pv_name)
1146 && !is_writable_pv(&state, &pv_name).await
1147 {
1148 let resp = encode_op_put_status_response(
1149 ioid,
1150 0x40,
1151 "Write access denied",
1152 version,
1153 is_be,
1154 );
1155 state.registry.send_msg(conn_id, resp).await;
1156 continue;
1157 }
1158 if let Some(nt) = get_nt_snapshot(&state, &pv_name).await {
1159 let resp = encode_op_put_getput_response_payload(
1160 ioid, &nt, version, is_be,
1161 );
1162 state.registry.send_msg(conn_id, resp).await;
1163 debug!(
1164 "Conn {}: put get-put pv='{}' ioid={}",
1165 conn_id, pv_name, ioid
1166 );
1167 } else {
1168 state
1169 .registry
1170 .send_msg(
1171 conn_id,
1172 encode_op_error(
1173 payload.command,
1174 payload.subcmd,
1175 ioid,
1176 "PV not found",
1177 version,
1178 is_be,
1179 ),
1180 )
1181 .await;
1182 }
1183 continue;
1184 }
1185 let desc = match conn_state.ioid_to_desc.get(&ioid) {
1186 Some(d) => d.clone(),
1187 None => {
1188 state
1189 .registry
1190 .send_msg(
1191 conn_id,
1192 encode_op_error(
1193 payload.command,
1194 payload.subcmd,
1195 ioid,
1196 "PUT without init",
1197 version,
1198 is_be,
1199 ),
1200 )
1201 .await;
1202 continue;
1203 }
1204 };
1205 let decoded = decode_put_body(&payload.body, &desc, is_be);
1206 if let Some(value) = decoded.as_ref() {
1207 match state.sources.put(&pv_name, value).await {
1208 Ok(changed) => {
1209 notify_changed_records(&state, changed).await;
1210 }
1211 Err(msg) => {
1212 let resp = encode_op_put_status_response(
1213 ioid,
1214 payload.subcmd,
1215 &msg,
1216 version,
1217 is_be,
1218 );
1219 state.registry.send_msg(conn_id, resp).await;
1220 continue;
1221 }
1222 }
1223 } else {
1224 debug!(
1225 "Conn {}: put decode failed ioid={} body_len={}",
1226 conn_id,
1227 ioid,
1228 payload.body.len()
1229 );
1230 let resp = encode_op_put_status_response(
1231 ioid,
1232 payload.subcmd,
1233 "cannot decode PUT body",
1234 version,
1235 is_be,
1236 );
1237 state.registry.send_msg(conn_id, resp).await;
1238 continue;
1239 }
1240 let resp = encode_op_put_response(ioid, payload.subcmd, version, is_be);
1241 state.registry.send_msg(conn_id, resp).await;
1242 debug!("Conn {}: put data pv='{}' ioid={}", conn_id, pv_name, ioid);
1243 }
1244 }
1245 12 => {
1246 if is_init {
1248 let desc =
1252 if let Some(desc) = state.sources.get_descriptor(&pv_name).await {
1253 desc
1254 } else if let Some(nt) = get_nt_snapshot(&state, &pv_name).await {
1255 nt_payload_desc(&nt)
1256 } else {
1257 state
1258 .registry
1259 .send_msg(
1260 conn_id,
1261 encode_op_error(
1262 payload.command,
1263 payload.subcmd,
1264 ioid,
1265 "PV not found",
1266 version,
1267 is_be,
1268 ),
1269 )
1270 .await;
1271 continue;
1272 };
1273 if !is_virtual_event_pv(&pv_name)
1274 && !is_writable_pv(&state, &pv_name).await
1275 {
1276 let resp = encode_op_put_get_init_error_response(
1277 ioid,
1278 "Write access denied",
1279 version,
1280 is_be,
1281 );
1282 state.registry.send_msg(conn_id, resp).await;
1283 continue;
1284 }
1285 conn_state.ioid_to_desc.insert(ioid, desc.clone());
1286 conn_state.ioid_to_pv.insert(ioid, pv_name.clone());
1287 let resp =
1288 encode_op_put_get_init_response(ioid, &desc, &desc, version, is_be);
1289 state.registry.send_msg(conn_id, resp).await;
1290 info!(
1291 "Conn {}: put_get init pv='{}' ioid={}",
1292 conn_id, pv_name, ioid
1293 );
1294 } else {
1295 let desc = match conn_state.ioid_to_desc.get(&ioid) {
1296 Some(d) => d.clone(),
1297 None => {
1298 state
1299 .registry
1300 .send_msg(
1301 conn_id,
1302 encode_op_error(
1303 payload.command,
1304 payload.subcmd,
1305 ioid,
1306 "PUT_GET without init",
1307 version,
1308 is_be,
1309 ),
1310 )
1311 .await;
1312 continue;
1313 }
1314 };
1315 let decoded = decode_put_body(&payload.body, &desc, is_be);
1316 if let Some(value) = decoded.as_ref() {
1317 match state.sources.put(&pv_name, value).await {
1318 Ok(changed) => {
1319 notify_changed_records(&state, changed).await;
1320 }
1321 Err(msg) => {
1322 let resp = encode_op_put_get_data_error_response(
1323 ioid, &msg, version, is_be,
1324 );
1325 state.registry.send_msg(conn_id, resp).await;
1326 continue;
1327 }
1328 }
1329 } else {
1330 debug!(
1331 "Conn {}: put_get decode failed ioid={} body_len={}",
1332 conn_id,
1333 ioid,
1334 payload.body.len()
1335 );
1336 let resp = encode_op_put_get_data_error_response(
1337 ioid,
1338 "cannot decode PUT body",
1339 version,
1340 is_be,
1341 );
1342 state.registry.send_msg(conn_id, resp).await;
1343 continue;
1344 }
1345 if let Some(nt) = get_nt_snapshot(&state, &pv_name).await {
1346 let resp = encode_op_put_get_data_response_payload(
1347 ioid, &nt, version, is_be,
1348 );
1349 state.registry.send_msg(conn_id, resp).await;
1350 } else {
1351 state
1352 .registry
1353 .send_msg(
1354 conn_id,
1355 encode_op_error(
1356 payload.command,
1357 payload.subcmd,
1358 ioid,
1359 "PV not found",
1360 version,
1361 is_be,
1362 ),
1363 )
1364 .await;
1365 }
1366 debug!(
1367 "Conn {}: put_get data pv='{}' ioid={}",
1368 conn_id, pv_name, ioid
1369 );
1370 }
1371 }
1372 13 => {
1373 if is_init {
1375 let full_desc =
1383 if let Some(desc) = state.sources.get_descriptor(&pv_name).await {
1384 desc
1385 } else if let Some(nt) = get_nt_snapshot(&state, &pv_name).await {
1386 nt_payload_desc(&nt)
1387 } else {
1388 state
1389 .registry
1390 .send_msg(
1391 conn_id,
1392 encode_op_error(
1393 payload.command,
1394 payload.subcmd,
1395 ioid,
1396 "PV not found",
1397 version,
1398 is_be,
1399 ),
1400 )
1401 .await;
1402 continue;
1403 };
1404 let pv_req_fields = decode_pv_request_fields(&payload.body, is_be);
1405 let desc = match &pv_req_fields {
1406 Some(fields) => filter_structure_desc(&full_desc, fields),
1407 None => full_desc,
1408 };
1409 conn_state.ioid_to_desc.insert(ioid, desc.clone());
1410 conn_state.ioid_to_pv.insert(ioid, pv_name.clone());
1411 let pipeline_enabled = (payload.subcmd & 0x80) != 0;
1412 let mut nfree = 0u32;
1413 if pipeline_enabled && payload.body.len() >= 4 {
1414 let start = payload.body.len() - 4;
1415 nfree = if is_be {
1416 u32::from_be_bytes([
1417 payload.body[start],
1418 payload.body[start + 1],
1419 payload.body[start + 2],
1420 payload.body[start + 3],
1421 ])
1422 } else {
1423 u32::from_le_bytes([
1424 payload.body[start],
1425 payload.body[start + 1],
1426 payload.body[start + 2],
1427 payload.body[start + 3],
1428 ])
1429 };
1430 }
1431 let resp = encode_op_init_response_desc(
1432 payload.command,
1433 ioid,
1434 0x08,
1435 &desc,
1436 version,
1437 is_be,
1438 );
1439 state.registry.send_msg(conn_id, resp).await;
1440 conn_state.ioid_to_monitor.insert(
1441 ioid,
1442 MonitorState {
1443 running: false,
1444 pipeline_enabled,
1445 nfree,
1446 },
1447 );
1448 {
1449 let mut monitors = state.registry.monitors.lock().await;
1450 monitors
1451 .entry(pv_name.clone())
1452 .or_default()
1453 .push(MonitorSub {
1454 conn_id,
1455 ioid,
1456 version,
1457 is_be,
1458 running: false,
1459 pipeline_enabled,
1460 nfree,
1461 filtered_desc: conn_state.ioid_to_desc.get(&ioid).cloned(),
1462 last_snapshot: None,
1463 });
1464 }
1465 info!(
1466 "Conn {}: monitor init pv='{}' ioid={}",
1467 conn_id, pv_name, ioid
1468 );
1469 } else if (payload.subcmd & 0x10) != 0 {
1470 if let Some(nt) = get_nt_snapshot(&state, &pv_name).await {
1472 let resp = encode_monitor_data_response_payload(
1473 ioid, 0x10, &nt, version, is_be,
1474 );
1475 state.registry.send_msg(conn_id, resp).await;
1476 }
1477 state
1478 .registry
1479 .remove_monitor_subscription(conn_id, ioid, &pv_name)
1480 .await;
1481 conn_state.ioid_to_monitor.remove(&ioid);
1482 conn_state.ioid_to_pv.remove(&ioid);
1483 conn_state.ioid_to_desc.remove(&ioid);
1484 info!("Conn {}: monitor end ioid={}", conn_id, ioid);
1485 } else if (payload.subcmd & 0x04) != 0 || (payload.subcmd & 0x80) != 0 {
1486 let start = (payload.subcmd & 0x44) == 0x44;
1488 let stop = (payload.subcmd & 0x44) == 0x04;
1489 let pipeline_ack = (payload.subcmd & 0x80) != 0;
1490 let mut nfree = None;
1491 if pipeline_ack && payload.body.len() >= 4 {
1492 let v = if is_be {
1493 u32::from_be_bytes([
1494 payload.body[0],
1495 payload.body[1],
1496 payload.body[2],
1497 payload.body[3],
1498 ])
1499 } else {
1500 u32::from_le_bytes([
1501 payload.body[0],
1502 payload.body[1],
1503 payload.body[2],
1504 payload.body[3],
1505 ])
1506 };
1507 nfree = Some(v);
1508 }
1509 let running = if start {
1510 true
1511 } else if stop {
1512 false
1513 } else {
1514 conn_state
1515 .ioid_to_monitor
1516 .get(&ioid)
1517 .map(|m| m.running)
1518 .unwrap_or(true)
1519 };
1520 state
1521 .registry
1522 .update_monitor_subscription(
1523 conn_id,
1524 ioid,
1525 &pv_name,
1526 running,
1527 nfree,
1528 Some(pipeline_ack),
1529 )
1530 .await;
1531 if let Some(mon) = conn_state.ioid_to_monitor.get_mut(&ioid) {
1532 mon.running = running;
1533 if pipeline_ack {
1534 mon.pipeline_enabled = true;
1535 }
1536 if let Some(v) = nfree {
1537 if pipeline_ack {
1538 mon.nfree = mon.nfree.saturating_add(v);
1539 } else {
1540 mon.nfree = v;
1541 }
1542 }
1543 }
1544 info!(
1545 "Conn {}: monitor {} ioid={} ack={} nfree={:?}",
1546 conn_id,
1547 if start {
1548 "start"
1549 } else if stop {
1550 "stop"
1551 } else {
1552 "ack"
1553 },
1554 ioid,
1555 pipeline_ack,
1556 nfree
1557 );
1558 if start {
1559 if let Some(nt) = get_nt_snapshot(&state, &pv_name).await {
1560 state
1561 .registry
1562 .send_monitor_update_for(&pv_name, conn_id, ioid, &nt)
1563 .await;
1564 }
1565 }
1566 }
1567 }
1568 20 => {
1569 if is_server_rpc_pv(&pv_name) {
1571 handle_server_rpc(
1572 &state,
1573 conn_id,
1574 ioid,
1575 payload.subcmd,
1576 version,
1577 is_be,
1578 )
1579 .await;
1580 } else {
1581 let is_init = (payload.subcmd & 0x08) != 0;
1582 if is_init {
1583 let resp = encode_op_status_response(
1585 20,
1586 ioid,
1587 payload.subcmd,
1588 version,
1589 is_be,
1590 );
1591 state.registry.send_msg(conn_id, resp).await;
1592 } else {
1593 let decoder = PvdDecoder::new(is_be);
1596 let args = if !payload.body.is_empty() {
1597 decoder
1598 .parse_introspection_with_len(&payload.body)
1599 .ok()
1600 .and_then(|(desc, consumed)| {
1601 decoder
1602 .decode_structure(&payload.body[consumed..], &desc)
1603 .ok()
1604 .map(|(val, _)| val)
1605 })
1606 } else {
1607 None
1608 };
1609 let empty =
1610 spvirit_codec::spvd_decode::DecodedValue::Structure(vec![]);
1611 let args_ref = args.as_ref().unwrap_or(&empty);
1612 match state.sources.rpc(&pv_name, args_ref).await {
1613 Ok(result) => {
1614 let resp = encode_op_rpc_data_response_payload(
1615 ioid,
1616 payload.subcmd,
1617 &result,
1618 version,
1619 is_be,
1620 );
1621 state.registry.send_msg(conn_id, resp).await;
1622 }
1623 Err(msg) => {
1624 let resp = encode_op_status_error_response(
1625 20,
1626 ioid,
1627 payload.subcmd,
1628 &msg,
1629 version,
1630 is_be,
1631 );
1632 state.registry.send_msg(conn_id, resp).await;
1633 }
1634 }
1635 }
1636 }
1637 }
1638 14 | 16 => {
1639 state
1640 .registry
1641 .send_msg(
1642 conn_id,
1643 encode_op_error(
1644 payload.command,
1645 payload.subcmd,
1646 ioid,
1647 "Operation not supported",
1648 version,
1649 is_be,
1650 ),
1651 )
1652 .await;
1653 }
1654 _ => {
1655 state
1656 .registry
1657 .send_msg(
1658 conn_id,
1659 encode_op_error(
1660 payload.command,
1661 payload.subcmd,
1662 ioid,
1663 "Operation not supported",
1664 version,
1665 is_be,
1666 ),
1667 )
1668 .await;
1669 }
1670 }
1671 }
1672 PvaPacketCommand::DestroyChannel(payload) => {
1673 let sid = payload.sid;
1674 let cid = payload.cid;
1675 conn_state.cid_to_sid.remove(&cid);
1676 conn_state.sid_to_pv.remove(&sid);
1677 info!(
1678 "Conn {}: channel destroyed sid={} cid={}",
1679 conn_id, sid, cid
1680 );
1681 }
1682 PvaPacketCommand::DestroyRequest(payload) => {
1683 let ioid = payload.request_id;
1684 if let Some(pv_name) = conn_state.ioid_to_pv.remove(&ioid) {
1685 state
1686 .registry
1687 .remove_monitor_subscription(conn_id, ioid, &pv_name)
1688 .await;
1689 conn_state.ioid_to_desc.remove(&ioid);
1690 conn_state.ioid_to_monitor.remove(&ioid);
1691 info!("Conn {}: monitor unsubscribed ioid={}", conn_id, ioid);
1692 }
1693 }
1694 PvaPacketCommand::AuthNZ(_) => {
1695 debug!("Conn {}: ignoring AUTHNZ", conn_id);
1697 }
1698 PvaPacketCommand::AclChange(_) => {
1699 let resp =
1700 encode_message_error("ACL_CHANGE command is not supported", version, is_be);
1701 state.registry.send_msg(conn_id, resp).await;
1702 }
1703 PvaPacketCommand::GetField(payload) => {
1704 handle_get_field_request(&state, &conn_state, conn_id, payload, version, is_be)
1705 .await;
1706 }
1707 PvaPacketCommand::Echo(payload_bytes) => {
1708 let mut resp =
1709 encode_header(true, is_be, false, version, 2, payload_bytes.len() as u32);
1710 resp.extend_from_slice(&payload_bytes);
1711 state.registry.send_msg(conn_id, resp).await;
1712 }
1713 PvaPacketCommand::Message(_) => {
1714 let resp = encode_message_error("MESSAGE command is not supported", version, is_be);
1715 state.registry.send_msg(conn_id, resp).await;
1716 }
1717 PvaPacketCommand::MultipleData(_) => {
1718 let resp =
1719 encode_message_error("MULTIPLE_DATA command is not supported", version, is_be);
1720 state.registry.send_msg(conn_id, resp).await;
1721 }
1722 PvaPacketCommand::CancelRequest(_) => {
1723 let resp =
1724 encode_message_error("CANCEL_REQUEST command is not supported", version, is_be);
1725 state.registry.send_msg(conn_id, resp).await;
1726 }
1727 PvaPacketCommand::OriginTag(_) => {
1728 let resp =
1729 encode_message_error("ORIGIN_TAG command is not supported", version, is_be);
1730 state.registry.send_msg(conn_id, resp).await;
1731 }
1732 PvaPacketCommand::Search(payload) => {
1733 debug!(
1734 "Conn {}: TCP search: pv_count={} mask=0x{:02x}",
1735 conn_id,
1736 payload.pv_requests.len(),
1737 payload.mask
1738 );
1739 let accepts_tcp = payload.protocols.is_empty()
1740 || payload
1741 .protocols
1742 .iter()
1743 .any(|p| p.eq_ignore_ascii_case("tcp"));
1744 if accepts_tcp {
1745 let all_names = state.sources.names().await;
1746 let visible_names = collect_visible_pv_names(
1747 &all_names,
1748 state.pvlist_mode,
1749 state.pvlist_allow_pattern.as_ref(),
1750 state.pvlist_max,
1751 );
1752 let mut cids = Vec::new();
1753 for (cid, name) in &payload.pv_requests {
1754 if state.sources.has_pv(name).await
1755 || is_virtual_event_pv(name)
1756 || (is_pvlist_virtual_pv(name) && state.pvlist_mode == PvListMode::List)
1757 || (is_server_rpc_pv(name) && state.pvlist_mode != PvListMode::Off)
1758 {
1759 cids.push(*cid);
1760 continue;
1761 }
1762 if state.pvlist_mode != PvListMode::Off
1763 && is_pattern_query(name)
1764 && visible_names.iter().any(|pv| wildcard_match(name, pv))
1765 {
1766 cids.push(*cid);
1767 }
1768 }
1769 let server_discovery_ping = payload.pv_requests.is_empty();
1770 let found = server_discovery_ping || !cids.is_empty();
1771 let resp_ip = state.advertise_ip.unwrap_or(state.listen_ip);
1772 let addr_bytes = if resp_ip.is_unspecified() {
1773 [0u8; 16]
1774 } else {
1775 ip_to_bytes(resp_ip)
1776 };
1777 let response = encode_search_response(
1778 state.guid,
1779 payload.seq,
1780 addr_bytes,
1781 state.tcp_port,
1782 "tcp",
1783 found,
1784 &cids,
1785 version,
1786 is_be,
1787 );
1788 state.registry.send_msg(conn_id, response).await;
1789 debug!(
1790 "Conn {}: TCP search responded found={} matches={}",
1791 conn_id,
1792 found,
1793 cids.len()
1794 );
1795 } else {
1796 debug!("Conn {}: TCP search: no compatible protocol", conn_id);
1797 }
1798 }
1799 PvaPacketCommand::SearchResponse(_) | PvaPacketCommand::Beacon(_) => {
1800 let resp =
1801 encode_message_error("Unexpected command for server endpoint", version, is_be);
1802 state.registry.send_msg(conn_id, resp).await;
1803 }
1804 PvaPacketCommand::Unknown(payload) => {
1805 let resp = encode_message_error(
1806 &format!("Unknown command {}", payload.command),
1807 version,
1808 is_be,
1809 );
1810 state.registry.send_msg(conn_id, resp).await;
1811 }
1812 }
1813 }
1814
1815 state.registry.cleanup_connection(conn_id).await;
1816 let _ = writer_task.await;
1817 Ok(())
1818}