1use 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#[derive(Clone, Copy, Debug, PartialEq, Eq)]
34pub enum PvListSource {
35 PvList,
36 GetField,
37 ServerRpc,
38 ServerGet,
39}
40
41const PV_REQUEST_EMPTY: [u8; 6] = [0xfd, 0x02, 0x00, 0x80, 0x00, 0x00];
44
45pub 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
53pub 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
201async 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
217pub 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 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
301async 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
408pub 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
426pub 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 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
557pub async fn pvlist(opts: &PvOptions, server_addr: SocketAddr) -> Result<Vec<String>, PvGetError> {
561 list_pvs_via_pvlist(opts, server_addr).await
562}
563
564pub 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
574pub 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#[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}