Skip to main content

spvirit_server/
handler.rs

1//! PVA protocol handler — the core TCP connection processor.
2//!
3//! [`handle_connection`] uses a [`SourceRegistry`] to resolve PV names across
4//! multiple registered sources, serving PVs over the EPICS PVAccess protocol.
5
6use 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// ---------------------------------------------------------------------------
45// PvListMode — controls virtual PV listing behaviour
46// ---------------------------------------------------------------------------
47
48/// Controls how the server exposes its PV directory.
49#[derive(Debug, Clone, Copy, PartialEq, Eq)]
50pub enum PvListMode {
51    /// No PV listing at all.
52    Off,
53    /// Respond to UDP search for known PVs only; no GET_FIELD listing.
54    Discover,
55    /// Full pvlist & server-RPC listing support.
56    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
73// ---------------------------------------------------------------------------
74// Server shared state
75// ---------------------------------------------------------------------------
76
77/// Shared server state that is passed to every connection handler.
78pub 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
123// ---------------------------------------------------------------------------
124// Virtual PV helpers
125// ---------------------------------------------------------------------------
126
127pub 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
150// ---------------------------------------------------------------------------
151// Pattern / wildcard utilities
152// ---------------------------------------------------------------------------
153
154pub 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
240// ---------------------------------------------------------------------------
241// Network helpers
242// ---------------------------------------------------------------------------
243
244pub 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
252/// Bind the fixed UDP search port with `SO_REUSEADDR` (and `SO_REUSEPORT`
253/// on Unix) so other local PVA consumers such as `p4p` can also listen on
254/// the same well-known port. On macOS in particular, a plain
255/// `UdpSocket::bind(5076)` prevents any subsequent binder from joining the
256/// port, which broke co-located clients.
257pub 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
303// ---------------------------------------------------------------------------
304// Debug utilities
305// ---------------------------------------------------------------------------
306
307pub 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
362// ---------------------------------------------------------------------------
363// Store-based snapshot/writable helpers (delegate to SourceRegistry + virtual PVs)
364// ---------------------------------------------------------------------------
365
366async 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
400// ---------------------------------------------------------------------------
401// Notify helpers
402// ---------------------------------------------------------------------------
403
404async 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
411// ---------------------------------------------------------------------------
412// GET_FIELD handler
413// ---------------------------------------------------------------------------
414
415async 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
551// ---------------------------------------------------------------------------
552// Server RPC handler
553// ---------------------------------------------------------------------------
554
555async 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
596// ---------------------------------------------------------------------------
597// Control message handler (inside segmented stream)
598// ---------------------------------------------------------------------------
599
600async 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
617// ---------------------------------------------------------------------------
618// UDP search handler
619// ---------------------------------------------------------------------------
620
621/// Run the UDP search responder.
622pub 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
743// ---------------------------------------------------------------------------
744// TCP server
745// ---------------------------------------------------------------------------
746
747/// Accept TCP connections and spawn a handler for each.
748///
749/// Callers must bind the `TcpListener` before spawning any other tasks so that
750/// an `EADDRINUSE` failure is detected eagerly and the beacon is never started.
751pub 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
771// ---------------------------------------------------------------------------
772// Core TCP connection handler
773// ---------------------------------------------------------------------------
774
775/// Handle a single PVA TCP connection.
776///
777/// This is the main protocol loop: handshake, then dispatch each command
778/// (CreateChannel, GET, PUT, PUT_GET, MONITOR, RPC, etc.) using the
779/// [`SourceRegistry`] abstraction.
780pub 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    // Per EPICS PVA protocol: send SET_BYTE_ORDER control message before validation.
805    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    // Server sends Connection Validation (cmd=1) next.
818    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    // One reassembler for the whole connection: the segments of a message may
833    // be separated by control frames, which are handled in between.
834    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        // Segment reassembly is delegated to the shared codec state machine.
881        // Control frames come back verbatim and are answered here; they never
882        // reach command dispatch below.
883        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        // Connection Validation (cmd=1): respond with CONNECTION_VALIDATED (cmd=9).
907        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                        // GET
1019                        if is_init {
1020                            // Init only needs the type descriptor, not the data.
1021                            // Use get_descriptor first; fall back to snapshot.
1022                            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                        // PUT
1092                        if is_init {
1093                            // Init only needs the type descriptor, not current data.
1094                            // Use get_descriptor first; fall back to snapshot so PUT
1095                            // can target PVs that do not yet have any data.
1096                            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                        // PUT_GET
1247                        if is_init {
1248                            // Init only needs the type descriptor, not current data.
1249                            // Use get_descriptor first; fall back to snapshot so that
1250                            // clients can initiate PUT_GET before any data exists.
1251                            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                        // MONITOR
1374                        if is_init {
1375                            // Init only needs the type descriptor, not the data.
1376                            // Use get_descriptor first; fall back to snapshot. This
1377                            // lets clients subscribe to PVs before any data has been
1378                            // produced (e.g. NTNDArray before acquire). Real monitor
1379                            // updates are pushed once data arrives. Strict clients
1380                            // like p4p treat a MONITOR init error as fatal, so we
1381                            // must not error out when only the descriptor is known.
1382                            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                            // Monitor destroy
1471                            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                            // Monitor start/stop/pipeline-ack
1487                            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                        // RPC
1570                        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                                // RPC INIT — acknowledge with status OK
1584                                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                                // RPC EXEC — decode self-describing PVD args and
1594                                // delegate to the source registry
1595                                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                // Silently accept — pvxs and pvAccessCPP ignore AUTHNZ.
1696                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}