Skip to main content

spvirit_client/
pvlist.rs

1//! PV listing — discover available PV names from a PVA server.
2//!
3//! Provides multiple discovery strategies (tried in order by
4//! [`pvlist_with_fallback`]):
5//!
6//! 1. `__pvlist` GET (preferred — spvirit / EPICS7 servers)
7//! 2. Connection-level GET_FIELD (legacy, opt-in via env var)
8//! 3. Server RPC `op=channels`
9//! 4. Server GET with heuristic parsing
10
11use std::net::SocketAddr;
12
13use tokio::io::AsyncWriteExt;
14use tokio::net::TcpStream;
15use tokio::time::timeout;
16
17use spvirit_codec::epics_decode::{DecodeMode, PvaPacket, PvaPacketCommand};
18use spvirit_codec::spvd_decode::{
19    DecodedValue, FieldDesc, FieldType, PvdDecoder, StructureDesc, extract_nt_scalar_value,
20};
21use spvirit_codec::spvd_encode::{encode_string_pvd, encode_structure_desc};
22use spvirit_codec::spvirit_encode::encode_rpc_request;
23
24use crate::client::{
25    ChannelConn, build_client_validation, encode_get_field_request, encode_get_request,
26    establish_channel, pvget,
27};
28use crate::transport::{read_packet, read_until};
29use crate::types::{PvGetError, PvOptions};
30use spvirit_codec::SegmentReassembler;
31
32/// Which discovery strategy succeeded.
33#[derive(Clone, Copy, Debug, PartialEq, Eq)]
34pub enum PvListSource {
35    PvList,
36    GetField,
37    ServerRpc,
38    ServerGet,
39}
40
41// ─── Helpers ─────────────────────────────────────────────────────────────────
42
43const PV_REQUEST_EMPTY: [u8; 6] = [0xfd, 0x02, 0x00, 0x80, 0x00, 0x00];
44
45/// Sort, deduplicate, and remove empty entries.
46pub fn normalize_pv_names(mut names: Vec<String>) -> Vec<String> {
47    names.retain(|name| !name.trim().is_empty());
48    names.sort();
49    names.dedup();
50    names
51}
52
53/// Extract PV names from a decoded `__pvlist` value (NTScalarArray of strings).
54pub fn parse_pvlist_value(value: &DecodedValue) -> Option<Vec<String>> {
55    let root = extract_nt_scalar_value(value).unwrap_or(value);
56    let DecodedValue::Array(items) = root else {
57        return None;
58    };
59    let mut out = Vec::with_capacity(items.len());
60    for item in items {
61        if let DecodedValue::String(name) = item {
62            out.push(name.clone());
63        } else {
64            return None;
65        }
66    }
67    Some(out)
68}
69
70fn candidate_server_addrs(opts: &PvOptions, server_addr: SocketAddr) -> Vec<SocketAddr> {
71    let mut out = vec![server_addr];
72    let default_addr = SocketAddr::new(server_addr.ip(), opts.tcp_port);
73    if default_addr != server_addr {
74        out.push(default_addr);
75    }
76    out
77}
78
79fn is_get_field_fallback_enabled() -> bool {
80    match std::env::var("EPICS_PVA_ENABLE_GET_FIELD_FALLBACK") {
81        Ok(v) => {
82            let v = v.trim().to_ascii_uppercase();
83            v == "YES" || v == "Y" || v == "1" || v == "TRUE"
84        }
85        Err(_) => false,
86    }
87}
88
89fn collect_strings_from_decoded(value: &DecodedValue, out: &mut Vec<String>) {
90    match value {
91        DecodedValue::String(s) => out.push(s.clone()),
92        DecodedValue::Array(items) => {
93            for item in items {
94                collect_strings_from_decoded(item, out);
95            }
96        }
97        DecodedValue::Structure(fields) => {
98            for (_, item) in fields {
99                collect_strings_from_decoded(item, out);
100            }
101        }
102        _ => {}
103    }
104}
105
106fn looks_like_pv_name(candidate: &str) -> bool {
107    if candidate.is_empty() || candidate.len() > 128 {
108        return false;
109    }
110    if candidate.chars().any(|c| c.is_whitespace()) {
111        return false;
112    }
113    let lower = candidate.to_ascii_lowercase();
114    let deny = [
115        "value",
116        "alarm",
117        "timestamp",
118        "display",
119        "control",
120        "severity",
121        "message",
122        "seconds",
123        "nanoseconds",
124        "units",
125    ];
126    if deny.iter().any(|d| lower == *d) {
127        return false;
128    }
129    if lower.starts_with("epics:") {
130        return false;
131    }
132    true
133}
134
135fn extract_ascii_candidates(raw: &[u8], out: &mut Vec<String>) {
136    let mut i = 0usize;
137    while i < raw.len() {
138        if raw[i].is_ascii_alphanumeric() {
139            let start = i;
140            i += 1;
141            while i < raw.len() {
142                let b = raw[i];
143                if b.is_ascii_alphanumeric()
144                    || b == b':'
145                    || b == b'.'
146                    || b == b'_'
147                    || b == b'-'
148                    || b == b'/'
149                {
150                    i += 1;
151                } else {
152                    break;
153                }
154            }
155            let len = i - start;
156            if (3..=128).contains(&len) {
157                if let Ok(s) = std::str::from_utf8(&raw[start..start + len]) {
158                    out.push(s.to_string());
159                }
160            }
161        } else {
162            i += 1;
163        }
164    }
165}
166
167fn encode_server_rpc_channels_request(is_be: bool) -> Vec<u8> {
168    let desc = StructureDesc {
169        struct_id: Some("epics:nt/NTURI:1.0".to_string()),
170        fields: vec![
171            FieldDesc {
172                name: "scheme".to_string(),
173                field_type: FieldType::String,
174            },
175            FieldDesc {
176                name: "path".to_string(),
177                field_type: FieldType::String,
178            },
179            FieldDesc {
180                name: "query".to_string(),
181                field_type: FieldType::Structure(StructureDesc {
182                    struct_id: None,
183                    fields: vec![FieldDesc {
184                        name: "op".to_string(),
185                        field_type: FieldType::String,
186                    }],
187                }),
188            },
189        ],
190    };
191
192    let mut out = Vec::new();
193    out.push(0x80);
194    out.extend_from_slice(&encode_structure_desc(&desc, is_be));
195    out.extend_from_slice(&encode_string_pvd("pva", is_be));
196    out.extend_from_slice(&encode_string_pvd("server", is_be));
197    out.extend_from_slice(&encode_string_pvd("channels", is_be));
198    out
199}
200
201// ─── Listing strategies ──────────────────────────────────────────────────────
202
203/// List PVs via the `__pvlist` channel (preferred).
204async fn list_pvs_via_pvlist(
205    opts: &PvOptions,
206    server_addr: SocketAddr,
207) -> Result<Vec<String>, PvGetError> {
208    let mut get_opts = opts.clone();
209    get_opts.pv_name = "__pvlist".to_string();
210    get_opts.server_addr = Some(server_addr);
211    let result = pvget(&get_opts).await?;
212    let names = parse_pvlist_value(&result.value)
213        .ok_or_else(|| PvGetError::Decode("failed to decode __pvlist value".to_string()))?;
214    Ok(normalize_pv_names(names))
215}
216
217/// List PVs via connection-level GET_FIELD.
218pub async fn list_pvs_via_get_field(
219    opts: &PvOptions,
220    server_addr: SocketAddr,
221    field_pattern: Option<&str>,
222) -> Result<Vec<String>, PvGetError> {
223    let mut stream = timeout(opts.timeout, TcpStream::connect(server_addr))
224        .await
225        .map_err(|_| PvGetError::Timeout("connect"))??;
226
227    // One reassembler for this connection, shared by every read below.
228    let mut reassembler = SegmentReassembler::new();
229
230    let mut version = 2u8;
231    let mut is_be = false;
232
233    for _ in 0..2 {
234        if let Ok(bytes) = read_packet(&mut stream, opts.timeout, &mut reassembler).await {
235            let mut pkt = PvaPacket::new(&bytes);
236            if let Some(cmd) = pkt.decode_payload() {
237                match cmd {
238                    PvaPacketCommand::Control(payload) => {
239                        if payload.command == 2 {
240                            is_be = pkt.header.flags.is_msb;
241                        }
242                    }
243                    PvaPacketCommand::ConnectionValidation(_) => {
244                        version = pkt.header.version;
245                        is_be = pkt.header.flags.is_msb;
246                    }
247                    _ => {}
248                }
249            }
250        }
251    }
252
253    let validation = build_client_validation(opts, version, is_be);
254    stream.write_all(&validation).await?;
255
256    let _ = read_until(&mut stream, opts.timeout, &mut reassembler, |cmd| {
257        matches!(cmd, PvaPacketCommand::ConnectionValidated(_))
258    })
259    .await?;
260
261    let get_field = encode_get_field_request(0, 0, field_pattern, version, is_be);
262    stream.write_all(&get_field).await?;
263
264    let field_resp = read_until(
265        &mut stream,
266        opts.timeout,
267        &mut reassembler,
268        |cmd| matches!(cmd, PvaPacketCommand::GetField(payload) if payload.is_server),
269    )
270    .await?;
271    let mut pkt = PvaPacket::new(&field_resp);
272    let cmd = pkt.decode_payload().ok_or(PvGetError::Protocol(
273        "get_field listing decode failed".to_string(),
274    ))?;
275    let PvaPacketCommand::GetField(payload) = cmd else {
276        return Err(PvGetError::Protocol(
277            "unexpected GET_FIELD response".to_string(),
278        ));
279    };
280
281    if payload.status.as_ref().is_some_and(|s| s.is_error()) {
282        let detail = payload
283            .status
284            .as_ref()
285            .map(ToString::to_string)
286            .unwrap_or_default();
287        return Err(PvGetError::Protocol(format!(
288            "get_field listing error: {}",
289            detail
290        )));
291    }
292
293    let desc = payload
294        .introspection
295        .ok_or_else(|| PvGetError::Decode("missing GET_FIELD introspection".to_string()))?;
296
297    let names = desc.fields.into_iter().map(|f| f.name).collect::<Vec<_>>();
298    Ok(normalize_pv_names(names))
299}
300
301/// List PVs via server RPC on a specific channel name.
302async fn list_pvs_via_server_rpc_channel(
303    opts: &PvOptions,
304    server_addr: SocketAddr,
305    rpc_channel: &str,
306) -> Result<Vec<String>, PvGetError> {
307    let mut rpc_opts = opts.clone();
308    rpc_opts.pv_name = rpc_channel.to_string();
309    let ChannelConn {
310        mut stream,
311        sid,
312        version,
313        is_be,
314        mut reassembler,
315        ..
316    } = establish_channel(server_addr, &rpc_opts).await?;
317
318    let ioid = 1u32;
319    let rpc_init = encode_rpc_request(sid, ioid, 0x08, &PV_REQUEST_EMPTY, version, is_be);
320    stream.write_all(&rpc_init).await?;
321
322    let init_resp = read_until(
323        &mut stream,
324        opts.timeout,
325        &mut reassembler,
326        |cmd| match cmd {
327            PvaPacketCommand::Op(op) => {
328                op.command == 20 && op.ioid == ioid && (op.subcmd & 0x08) != 0
329            }
330            _ => false,
331        },
332    )
333    .await?;
334    let mut pkt = PvaPacket::new(&init_resp);
335    let init_cmd = pkt
336        .decode_payload()
337        .ok_or(PvGetError::Protocol("rpc init decode failed".to_string()))?;
338    if let PvaPacketCommand::Op(op) = init_cmd {
339        if op.status.as_ref().is_some_and(|s| s.is_error()) {
340            let detail = op
341                .status
342                .as_ref()
343                .map(ToString::to_string)
344                .unwrap_or_default();
345            return Err(PvGetError::Protocol(format!("rpc init failed: {}", detail)));
346        }
347    }
348
349    let rpc_payload = encode_server_rpc_channels_request(is_be);
350    let rpc_req = encode_rpc_request(sid, ioid, 0x00, &rpc_payload, version, is_be);
351    stream.write_all(&rpc_req).await?;
352
353    let rpc_resp = read_until(
354        &mut stream,
355        opts.timeout,
356        &mut reassembler,
357        |cmd| match cmd {
358            PvaPacketCommand::Op(op) => op.command == 20 && op.ioid == ioid && op.subcmd == 0x00,
359            _ => false,
360        },
361    )
362    .await?;
363    let mut pkt = PvaPacket::new(&rpc_resp);
364    let rpc_cmd = pkt.decode_payload().ok_or(PvGetError::Protocol(
365        "rpc response decode failed".to_string(),
366    ))?;
367    let PvaPacketCommand::Op(op) = rpc_cmd else {
368        return Err(PvGetError::Protocol("unexpected RPC response".to_string()));
369    };
370    if op.status.as_ref().is_some_and(|s| s.is_error()) {
371        let detail = op
372            .status
373            .as_ref()
374            .map(ToString::to_string)
375            .unwrap_or_default();
376        return Err(PvGetError::Protocol(format!(
377            "rpc execute failed: {}",
378            detail
379        )));
380    }
381
382    if op.body.is_empty() {
383        return Err(PvGetError::Decode("empty RPC response".to_string()));
384    }
385
386    let decoder = PvdDecoder::new(is_be);
387    let (desc, consumed) = decoder
388        .parse_introspection_with_len(&op.body)
389        .map_err(|_| PvGetError::Decode("RPC missing introspection".to_string()))?;
390    let value_raw = op
391        .body
392        .get(consumed..)
393        .ok_or_else(|| PvGetError::Decode("RPC malformed payload".to_string()))?;
394    let (decoded, _) = decoder
395        .decode_structure(value_raw, &desc)
396        .map_err(|_| PvGetError::Decode("RPC decode failed".to_string()))?;
397
398    let mut strings = Vec::new();
399    collect_strings_from_decoded(&decoded, &mut strings);
400    if strings.is_empty() {
401        return Err(PvGetError::Decode(
402            "RPC list returned no strings".to_string(),
403        ));
404    }
405    Ok(normalize_pv_names(strings))
406}
407
408/// List PVs via server RPC, trying `"server"` then `"__server"`.
409pub async fn list_pvs_via_server_rpc(
410    opts: &PvOptions,
411    server_addr: SocketAddr,
412) -> Result<Vec<String>, PvGetError> {
413    let mut errs = Vec::new();
414    for channel in ["server", "__server"] {
415        match list_pvs_via_server_rpc_channel(opts, server_addr, channel).await {
416            Ok(names) => return Ok(names),
417            Err(err) => errs.push(format!("{}: {}", channel, err)),
418        }
419    }
420    Err(PvGetError::Protocol(format!(
421        "server RPC unavailable: {}",
422        errs.join(" | ")
423    )))
424}
425
426/// List PVs via server GET, trying `"server"` then `"__server"`.
427pub async fn list_pvs_via_server_get(
428    opts: &PvOptions,
429    server_addr: SocketAddr,
430) -> Result<Vec<String>, PvGetError> {
431    async fn get_channel(
432        opts: &PvOptions,
433        server_addr: SocketAddr,
434        channel: &str,
435    ) -> Result<Vec<String>, PvGetError> {
436        let mut get_opts = opts.clone();
437        get_opts.pv_name = channel.to_string();
438        let ChannelConn {
439            mut stream,
440            sid,
441            version,
442            is_be,
443            mut reassembler,
444            ..
445        } = establish_channel(server_addr, &get_opts).await?;
446
447        let ioid = 1u32;
448        let init_req = encode_get_request(sid, ioid, 0x08, &PV_REQUEST_EMPTY, version, is_be);
449        stream.write_all(&init_req).await?;
450        let init_resp = read_until(
451            &mut stream,
452            opts.timeout,
453            &mut reassembler,
454            |cmd| match cmd {
455                PvaPacketCommand::Op(op) => {
456                    op.command == 10 && op.ioid == ioid && (op.subcmd & 0x08) != 0
457                }
458                _ => false,
459            },
460        )
461        .await?;
462        let mut pkt = PvaPacket::new(&init_resp);
463        let init_cmd = pkt.decode_payload().ok_or(PvGetError::Protocol(
464            "server get init decode failed".to_string(),
465        ))?;
466        let init_desc = match init_cmd {
467            PvaPacketCommand::Op(op) => {
468                if op.status.as_ref().is_some_and(|s| s.is_error()) {
469                    let detail = op
470                        .status
471                        .as_ref()
472                        .map(ToString::to_string)
473                        .unwrap_or_default();
474                    return Err(PvGetError::Protocol(format!(
475                        "server GET init failed: {}",
476                        detail
477                    )));
478                }
479                op.introspection
480            }
481            _ => None,
482        };
483
484        let data_req = encode_get_request(sid, ioid, 0x00, &[], version, is_be);
485        stream.write_all(&data_req).await?;
486        let data_resp = read_until(
487            &mut stream,
488            opts.timeout,
489            &mut reassembler,
490            |cmd| match cmd {
491                PvaPacketCommand::Op(op) => {
492                    op.command == 10 && op.ioid == ioid && op.subcmd == 0x00
493                }
494                _ => false,
495            },
496        )
497        .await?;
498        let mut pkt = PvaPacket::new(&data_resp);
499        let data_cmd = pkt.decode_payload().ok_or(PvGetError::Protocol(
500            "server get data decode failed".to_string(),
501        ))?;
502
503        let PvaPacketCommand::Op(mut op) = data_cmd else {
504            return Err(PvGetError::Protocol(
505                "unexpected GET data response".to_string(),
506            ));
507        };
508        if op.status.as_ref().is_some_and(|s| s.is_error()) {
509            let detail = op
510                .status
511                .as_ref()
512                .map(ToString::to_string)
513                .unwrap_or_default();
514            return Err(PvGetError::Protocol(format!(
515                "server GET data failed: {}",
516                detail
517            )));
518        }
519
520        let mut names = Vec::new();
521        names.extend(op.pv_names.clone());
522        extract_ascii_candidates(&op.body, &mut names);
523        if let Some(desc) = &init_desc {
524            for field in &desc.fields {
525                names.push(field.name.clone());
526            }
527            // A decode failure leaves decoded_value None, handled below.
528            let _ = op.decode_with_field_desc(desc, is_be, DecodeMode::Strict);
529            if let Some(decoded) = &op.decoded_value {
530                collect_strings_from_decoded(decoded, &mut names);
531            }
532        }
533
534        let mut names = normalize_pv_names(names);
535        names.retain(|n| looks_like_pv_name(n));
536        if names.is_empty() {
537            return Err(PvGetError::Decode(
538                "server GET returned no PV-like names".to_string(),
539            ));
540        }
541        Ok(names)
542    }
543
544    let mut errs = Vec::new();
545    for channel in ["server", "__server"] {
546        match get_channel(opts, server_addr, channel).await {
547            Ok(names) => return Ok(names),
548            Err(err) => errs.push(format!("{}: {}", channel, err)),
549        }
550    }
551    Err(PvGetError::Protocol(format!(
552        "server GET unavailable: {}",
553        errs.join(" | ")
554    )))
555}
556
557// ─── Public API ──────────────────────────────────────────────────────────────
558
559/// List PV names from a server using `__pvlist` GET (preferred method).
560pub async fn pvlist(opts: &PvOptions, server_addr: SocketAddr) -> Result<Vec<String>, PvGetError> {
561    list_pvs_via_pvlist(opts, server_addr).await
562}
563
564/// List PV names with automatic fallback through all strategies.
565///
566/// Tries (in order): `__pvlist` → GET_FIELD (opt-in) → Server RPC → Server GET.
567pub async fn pvlist_with_fallback(
568    opts: &PvOptions,
569    server_addr: SocketAddr,
570) -> Result<(Vec<String>, PvListSource), PvGetError> {
571    pvlist_with_fallback_progress(opts, server_addr, |_| {}).await
572}
573
574/// List PV names with fallback and progress callback.
575pub async fn pvlist_with_fallback_progress<F>(
576    opts: &PvOptions,
577    server_addr: SocketAddr,
578    mut on_progress: F,
579) -> Result<(Vec<String>, PvListSource), PvGetError>
580where
581    F: FnMut(&str),
582{
583    let addrs = candidate_server_addrs(opts, server_addr);
584    let mut attempts = Vec::new();
585    let get_field_fallback = is_get_field_fallback_enabled();
586
587    if addrs.len() > 1 {
588        on_progress(&format!(
589            "Trying {} candidate server endpoints...",
590            addrs.len()
591        ));
592    }
593    if !get_field_fallback {
594        on_progress(
595            "GET_FIELD fallback is disabled by default (set EPICS_PVA_ENABLE_GET_FIELD_FALLBACK=YES to enable)",
596        );
597    }
598
599    for addr in addrs {
600        on_progress(&format!("Trying __pvlist on {}", addr));
601        let primary = list_pvs_via_pvlist(opts, addr).await;
602        match primary {
603            Ok(names) => return Ok((normalize_pv_names(names), PvListSource::PvList)),
604            Err(primary_err) => {
605                let get_field_result = if get_field_fallback {
606                    on_progress(&format!(
607                        "__pvlist unavailable on {}; trying GET_FIELD(*)",
608                        addr
609                    ));
610                    let fallback_star = list_pvs_via_get_field(opts, addr, Some("*")).await;
611                    match fallback_star {
612                        Ok(names) => {
613                            return Ok((normalize_pv_names(names), PvListSource::GetField));
614                        }
615                        Err(star_err) => {
616                            on_progress(&format!(
617                                "GET_FIELD(*) unavailable on {}; trying GET_FIELD(<empty>)",
618                                addr
619                            ));
620                            let fallback_empty = list_pvs_via_get_field(opts, addr, None).await;
621                            match fallback_empty {
622                                Ok(names) => {
623                                    return Ok((normalize_pv_names(names), PvListSource::GetField));
624                                }
625                                Err(empty_err) => Some(format!(
626                                    "GET_FIELD(*): {}; GET_FIELD(<empty>): {}",
627                                    star_err, empty_err
628                                )),
629                            }
630                        }
631                    }
632                } else {
633                    None
634                };
635
636                on_progress(&format!(
637                    "__pvlist unavailable on {}; trying RPC(server)",
638                    addr
639                ));
640                match list_pvs_via_server_rpc(opts, addr).await {
641                    Ok(names) => return Ok((normalize_pv_names(names), PvListSource::ServerRpc)),
642                    Err(rpc_err) => {
643                        on_progress(&format!(
644                            "RPC(server) unavailable on {}; trying GET(server)",
645                            addr
646                        ));
647                        match list_pvs_via_server_get(opts, addr).await {
648                            Ok(names) => {
649                                return Ok((normalize_pv_names(names), PvListSource::ServerGet));
650                            }
651                            Err(get_err) => {
652                                let get_field_msg = get_field_result
653                                    .unwrap_or_else(|| "GET_FIELD: disabled".to_string());
654                                attempts.push(format!(
655                                    "{} => __pvlist: {}; {}; RPC(server): {}; GET(server): {}",
656                                    addr, primary_err, get_field_msg, rpc_err, get_err
657                                ));
658                            }
659                        }
660                    }
661                }
662            }
663        }
664    }
665
666    Err(PvGetError::Protocol(format!(
667        "failed to list PVs from {}: {}",
668        server_addr,
669        attempts.join(" | ")
670    )))
671}
672
673// ─── Tests ───────────────────────────────────────────────────────────────────
674
675#[cfg(test)]
676mod tests {
677    use super::*;
678
679    #[test]
680    fn parse_pvlist_value_extracts_ntscalararray_strings() {
681        let value = DecodedValue::Structure(vec![
682            (
683                "value".to_string(),
684                DecodedValue::Array(vec![
685                    DecodedValue::String("SIM:AI".to_string()),
686                    DecodedValue::String("SIM:AO".to_string()),
687                ]),
688            ),
689            ("alarm".to_string(), DecodedValue::Structure(vec![])),
690        ]);
691
692        let parsed = parse_pvlist_value(&value).expect("parsed");
693        assert_eq!(parsed, vec!["SIM:AI".to_string(), "SIM:AO".to_string()]);
694    }
695
696    #[test]
697    fn normalize_pv_names_sorts_and_deduplicates() {
698        let names = vec!["B".into(), "A".into(), "B".into(), " ".into()];
699        let result = normalize_pv_names(names);
700        assert_eq!(result, vec!["A".to_string(), "B".to_string()]);
701    }
702
703    #[test]
704    fn candidate_server_addrs_adds_default_tcp_port_fallback() {
705        let mut opts = PvOptions::new(String::new());
706        opts.tcp_port = 5075;
707        let addr: SocketAddr = "10.0.0.2:6000".parse().unwrap();
708        let addrs = candidate_server_addrs(&opts, addr);
709        assert_eq!(addrs.len(), 2);
710        assert_eq!(addrs[0], addr);
711        assert_eq!(addrs[1], "10.0.0.2:5075".parse::<SocketAddr>().unwrap());
712    }
713
714    #[test]
715    fn candidate_server_addrs_no_dup_when_same_port() {
716        let mut opts = PvOptions::new(String::new());
717        opts.tcp_port = 6000;
718        let addr: SocketAddr = "10.0.0.2:6000".parse().unwrap();
719        let addrs = candidate_server_addrs(&opts, addr);
720        assert_eq!(addrs.len(), 1);
721    }
722
723    #[test]
724    fn collect_strings_from_decoded_extracts_nested_strings() {
725        let value = DecodedValue::Structure(vec![
726            ("a".to_string(), DecodedValue::String("ONE".to_string())),
727            (
728                "b".to_string(),
729                DecodedValue::Array(vec![DecodedValue::String("TWO".to_string())]),
730            ),
731        ]);
732        let mut out = Vec::new();
733        collect_strings_from_decoded(&value, &mut out);
734        assert_eq!(out, vec!["ONE".to_string(), "TWO".to_string()]);
735    }
736
737    #[test]
738    fn extract_ascii_candidates_finds_pv_like_tokens() {
739        let raw = b"\x00SIM:AI\x00junk\x00IOC-01:PV1\x00";
740        let mut out = Vec::new();
741        extract_ascii_candidates(raw, &mut out);
742        assert!(out.iter().any(|s| s == "SIM:AI"));
743        assert!(out.iter().any(|s| s == "IOC-01:PV1"));
744    }
745
746    #[test]
747    fn looks_like_pv_name_filters_metadata() {
748        assert!(looks_like_pv_name("SIM:AI"));
749        assert!(looks_like_pv_name("IOC-01:PV1"));
750        assert!(!looks_like_pv_name("value"));
751        assert!(!looks_like_pv_name("alarm"));
752        assert!(!looks_like_pv_name("epics:nt/NTScalar:1.0"));
753        assert!(!looks_like_pv_name(""));
754        assert!(!looks_like_pv_name("has space"));
755    }
756
757    #[test]
758    fn encode_server_rpc_channels_request_uses_nturi_channels() {
759        let payload = encode_server_rpc_channels_request(false);
760        assert_eq!(payload.first(), Some(&0x80));
761
762        let decoder = PvdDecoder::new(false);
763        let (desc, consumed) = decoder
764            .parse_introspection_with_len(&payload)
765            .expect("introspection");
766        assert_eq!(desc.struct_id.as_deref(), Some("epics:nt/NTURI:1.0"));
767
768        let (decoded, _) = decoder
769            .decode_structure(&payload[consumed..], &desc)
770            .expect("decode payload");
771        let DecodedValue::Structure(fields) = decoded else {
772            panic!("expected structure");
773        };
774
775        let mut scheme = None;
776        let mut path = None;
777        let mut op = None;
778        for (name, value) in fields {
779            match (name.as_str(), value) {
780                ("scheme", DecodedValue::String(v)) => scheme = Some(v),
781                ("path", DecodedValue::String(v)) => path = Some(v),
782                ("query", DecodedValue::Structure(query_fields)) => {
783                    for (qname, qvalue) in query_fields {
784                        if qname == "op" {
785                            if let DecodedValue::String(v) = qvalue {
786                                op = Some(v);
787                            }
788                        }
789                    }
790                }
791                _ => {}
792            }
793        }
794        assert_eq!(scheme.as_deref(), Some("pva"));
795        assert_eq!(path.as_deref(), Some("server"));
796        assert_eq!(op.as_deref(), Some("channels"));
797    }
798}