1use bhtune_driver::{
6 BrowseNode, BrowseNodeKind, BrowsePage, BrowsePageRequest, Driver, OpcDaDriver, Quality,
7 SearchEvent, SearchIndexControlAction, SearchIndexRequest, SearchIndexResponse,
8 SearchIndexStatus, SearchMatch, SearchRequest, TagWrite, close_opcda_browse_session,
9 list_opcda_servers,
10};
11
12use crate::args::{OpcCommand, OpcSearchMatchModeArg, SearchIndexCommand};
13use crate::output::OutputFormat;
14
15pub async fn run(command: OpcCommand, config: &crate::config::BhtuneConfig) -> anyhow::Result<()> {
16 run_with_output(command, config, OutputFormat::Table).await
17}
18
19pub async fn run_with_output(
20 command: OpcCommand,
21 config: &crate::config::BhtuneConfig,
22 output: OutputFormat,
23) -> anyhow::Result<()> {
24 match command {
25 OpcCommand::Servers { bridge_host } => {
26 let bridge_host = crate::config::resolve_bridge_host(bridge_host, config);
27 servers_with_output(&bridge_host, output).await
28 }
29 OpcCommand::Read {
30 bridge_host,
31 server,
32 tags,
33 } => {
34 let bridge_host = crate::config::resolve_bridge_host(bridge_host, config);
35 let server = crate::config::resolve_server(server, config)?;
36 read_with_output(&bridge_host, &server, &tags, output).await
37 }
38 OpcCommand::Write {
39 bridge_host,
40 server,
41 tag,
42 value,
43 } => {
44 let bridge_host = crate::config::resolve_bridge_host(bridge_host, config);
45 let server = crate::config::resolve_server(server, config)?;
46 write_with_output(&bridge_host, &server, &tag, &value, output).await
47 }
48 OpcCommand::Browse {
49 bridge_host,
50 server,
51 session_id,
52 parent_node_key,
53 page_token,
54 page_size,
55 all,
56 refresh,
57 } => {
58 let bridge_host = crate::config::resolve_bridge_host(bridge_host, config);
59 let server = crate::config::resolve_server(server, config)?;
60 browse_with_output(
61 &bridge_host,
62 &server,
63 BrowseOptions {
64 session_id,
65 parent_node_key,
66 page_token,
67 page_size,
68 all,
69 refresh,
70 },
71 output,
72 )
73 .await
74 }
75 OpcCommand::Close {
76 bridge_host,
77 session_id,
78 } => {
79 let bridge_host = crate::config::resolve_bridge_host(bridge_host, config);
80 close_with_output(&bridge_host, &session_id, output).await
81 }
82 OpcCommand::Search {
83 bridge_host,
84 server,
85 query,
86 match_mode,
87 max_results,
88 session_id,
89 scope_node_key,
90 include_branches,
91 refresh,
92 } => {
93 let bridge_host = crate::config::resolve_bridge_host(bridge_host, config);
94 let server = crate::config::resolve_server(server, config)?;
95 search_with_output(
96 &bridge_host,
97 &server,
98 SearchOptions {
99 query,
100 match_mode,
101 max_results,
102 session_id,
103 scope_node_key,
104 include_branches,
105 refresh,
106 },
107 output,
108 )
109 .await
110 }
111 OpcCommand::SearchIndex { command } => {
112 run_search_index_command(command, config, output).await
113 }
114 }
115}
116
117async fn run_search_index_command(
118 command: SearchIndexCommand,
119 config: &crate::config::BhtuneConfig,
120 output: OutputFormat,
121) -> anyhow::Result<()> {
122 match command {
123 SearchIndexCommand::Status {
124 bridge_host,
125 server,
126 } => {
127 let bridge_host = crate::config::resolve_bridge_host(bridge_host, config);
128 let server = crate::config::resolve_server(server, config)?;
129 let driver = OpcDaDriver::connect(&bridge_host, &server).await?;
130 let status = driver.search_index_status().await?;
131 print_search_index_status(&status, output)
132 }
133 SearchIndexCommand::Search {
134 bridge_host,
135 server,
136 query,
137 match_mode,
138 max_results,
139 } => {
140 let bridge_host = crate::config::resolve_bridge_host(bridge_host, config);
141 let server = crate::config::resolve_server(server, config)?;
142 let driver = OpcDaDriver::connect(&bridge_host, &server).await?;
143 let response = driver
144 .search_index(SearchIndexRequest::new(
145 query,
146 match_mode.into(),
147 max_results,
148 ))
149 .await?;
150 print_search_index_results(&response, output)
151 }
152 SearchIndexCommand::Refresh {
153 bridge_host,
154 server,
155 force,
156 } => {
157 let bridge_host = crate::config::resolve_bridge_host(bridge_host, config);
158 let server = crate::config::resolve_server(server, config)?;
159 let driver = OpcDaDriver::connect(&bridge_host, &server).await?;
160 let status = driver.refresh_search_index(force).await?;
161 print_search_index_status(&status, output)
162 }
163 SearchIndexCommand::Control {
164 bridge_host,
165 server,
166 action,
167 } => {
168 let bridge_host = crate::config::resolve_bridge_host(bridge_host, config);
169 let server = crate::config::resolve_server(server, config)?;
170 let driver = OpcDaDriver::connect(&bridge_host, &server).await?;
171 let status = driver
172 .control_search_index(SearchIndexControlAction::from(action))
173 .await?;
174 print_search_index_status(&status, output)
175 }
176 }
177}
178
179#[cfg(test)]
180async fn servers(bridge_host: &str) -> anyhow::Result<()> {
181 servers_with_output(bridge_host, OutputFormat::Table).await
182}
183
184async fn servers_with_output(bridge_host: &str, output: OutputFormat) -> anyhow::Result<()> {
185 let servers = list_opcda_servers(bridge_host).await?;
186 if output == OutputFormat::Json {
187 println!(
188 "{}",
189 serde_json::to_string_pretty(&serde_json::json!({ "servers": servers }))?
190 );
191 return Ok(());
192 }
193 if servers.is_empty() {
194 println!("No OPC DA servers registered on the gateway's host.");
195 return Ok(());
196 }
197 for server in servers {
198 println!("{server}");
199 }
200 Ok(())
201}
202
203#[cfg(test)]
204async fn read(bridge_host: &str, server: &str, tags: &[String]) -> anyhow::Result<()> {
205 read_with_output(bridge_host, server, tags, OutputFormat::Table).await
206}
207
208async fn read_with_output(
209 bridge_host: &str,
210 server: &str,
211 tags: &[String],
212 output: OutputFormat,
213) -> anyhow::Result<()> {
214 if tags.is_empty() {
215 anyhow::bail!("at least one tag is required");
216 }
217 let driver = OpcDaDriver::connect(bridge_host, server).await?;
218 let values = driver.read(tags).await?;
219 if output == OutputFormat::Json {
220 let values = values
221 .into_iter()
222 .map(|value| {
223 serde_json::json!({
224 "tag": value.tag,
225 "value": value.value,
226 "quality": quality_name(value.quality),
227 "timestamp": value.timestamp.map(|time| time.to_rfc3339()),
228 })
229 })
230 .collect::<Vec<_>>();
231 println!(
232 "{}",
233 serde_json::to_string_pretty(&serde_json::json!({ "values": values }))?
234 );
235 return Ok(());
236 }
237 println!(
238 "{:<40} {:<15} {:<10} {:<20}",
239 "TAG", "VALUE", "QUALITY", "TIMESTAMP"
240 );
241 for v in values {
242 println!(
243 "{:<40} {:<15} {:<10} {:<20}",
244 v.tag,
245 v.value,
246 format!("{:?}", v.quality),
247 format_timestamp(v.timestamp),
248 );
249 }
250 Ok(())
251}
252
253#[cfg(test)]
254async fn write(bridge_host: &str, server: &str, tag: &str, value: &str) -> anyhow::Result<()> {
255 write_with_output(bridge_host, server, tag, value, OutputFormat::Table).await
256}
257
258async fn write_with_output(
259 bridge_host: &str,
260 server: &str,
261 tag: &str,
262 value: &str,
263 output: OutputFormat,
264) -> anyhow::Result<()> {
265 let driver = OpcDaDriver::connect(bridge_host, server).await?;
266 let write_value = match value.parse::<f32>() {
269 Ok(f) => TagWrite::Float(f),
270 Err(_) => TagWrite::Raw(value.to_string()),
271 };
272 let outcome = driver.write(&tag.to_string(), write_value).await?;
273 if output == OutputFormat::Json {
274 println!(
275 "{}",
276 serde_json::to_string_pretty(&serde_json::json!({
277 "tag": tag,
278 "value": value,
279 "success": outcome.success,
280 "error": outcome.error_message,
281 }))?
282 );
283 if !outcome.success {
284 anyhow::bail!(
285 "driver rejected the write: {}",
286 outcome
287 .error_message
288 .unwrap_or_else(|| "unknown reason".to_string())
289 );
290 }
291 return Ok(());
292 }
293 if outcome.success {
294 println!("Wrote '{value}' to '{tag}'.");
295 Ok(())
296 } else {
297 anyhow::bail!(
298 "driver rejected the write: {}",
299 write_rejection_reason(outcome.error_message)
300 )
301 }
302}
303
304fn format_timestamp(timestamp: Option<chrono::DateTime<chrono::Utc>>) -> String {
305 timestamp
306 .map(|timestamp| timestamp.to_rfc3339())
307 .unwrap_or_else(|| "-".to_string())
308}
309
310fn write_rejection_reason(error_message: Option<String>) -> String {
311 error_message.unwrap_or_else(|| "unknown reason".to_string())
312}
313
314#[cfg(test)]
315async fn browse(bridge_host: &str, server: &str, _path: &str) -> anyhow::Result<()> {
316 browse_with_output(
317 bridge_host,
318 server,
319 BrowseOptions {
320 page_size: bhtune_driver::DEFAULT_PAGE_SIZE,
321 ..BrowseOptions::default()
322 },
323 OutputFormat::Table,
324 )
325 .await
326}
327
328#[derive(Debug, Default)]
329struct BrowseOptions {
330 session_id: Option<String>,
331 parent_node_key: Option<String>,
332 page_token: Option<String>,
333 page_size: u32,
334 all: bool,
335 refresh: bool,
336}
337
338#[derive(serde::Serialize)]
339struct BrowsePagesOutput {
340 session_id: String,
341 pages: Vec<serde_json::Value>,
342 complete: bool,
343}
344
345async fn browse_with_output(
346 bridge_host: &str,
347 server: &str,
348 options: BrowseOptions,
349 output: OutputFormat,
350) -> anyhow::Result<()> {
351 validate_browse_options(&options)?;
352
353 let driver = OpcDaDriver::connect(bridge_host, server).await?;
354 let first_request = BrowsePageRequest {
355 session_id: options.session_id,
356 parent_node_key: options.parent_node_key.clone(),
357 page_token: options.page_token,
358 page_size: options.page_size,
359 refresh: options.refresh,
360 };
361 let first = driver.browse(first_request).await?;
362 let session = first.session_id.clone();
363 let pages = collect_browse_pages(
364 &driver,
365 first,
366 options.parent_node_key,
367 options.page_size,
368 options.all,
369 )
370 .await?;
371 render_browse_pages(&pages, &session, output)
372}
373
374fn validate_browse_options(options: &BrowseOptions) -> anyhow::Result<()> {
375 if (options.parent_node_key.is_some() || options.page_token.is_some())
376 && options.session_id.is_none()
377 {
378 anyhow::bail!("--session-id is required with --parent-node-key or --page-token");
379 }
380 Ok(())
381}
382
383fn render_browse_pages(
384 pages: &[BrowsePage],
385 session_id: &str,
386 output: OutputFormat,
387) -> anyhow::Result<()> {
388 if output == OutputFormat::Json {
389 return render_browse_json(pages, session_id);
390 }
391 render_browse_table(pages);
392 Ok(())
393}
394
395fn render_browse_json(pages: &[BrowsePage], session_id: &str) -> anyhow::Result<()> {
396 if pages.len() == 1 {
397 println!(
398 "{}",
399 serde_json::to_string_pretty(&json_browse_page(&pages[0]))?
400 );
401 } else {
402 let response = BrowsePagesOutput {
403 session_id: session_id.to_owned(),
404 pages: pages.iter().map(json_browse_page).collect(),
405 complete: pages.last().is_some_and(|page| page.complete),
406 };
407 println!("{}", serde_json::to_string_pretty(&response)?);
408 }
409 Ok(())
410}
411
412fn render_browse_table(pages: &[BrowsePage]) {
413 for (index, page) in pages.iter().enumerate() {
414 if pages.len() > 1 {
415 println!("Page {}:", index + 1);
416 }
417 if page.nodes.is_empty() {
418 println!("No tags found at this level.");
419 }
420 for node in &page.nodes {
421 let marker = match node.kind {
422 BrowseNodeKind::Branch => "[+]",
423 BrowseNodeKind::Item => " ",
424 BrowseNodeKind::BranchAndItem => "[*]",
425 BrowseNodeKind::Unspecified => "[?]",
426 };
427 let item = node
428 .item_id
429 .as_deref()
430 .map(|id| format!(" -> {id}"))
431 .unwrap_or_default();
432 println!("{marker} {}{item}", node.display_name);
433 }
434 if let Some(token) = &page.next_page_token {
435 println!("More pages available (next token: {token}).");
436 }
437 }
438}
439
440async fn close_with_output(
441 bridge_host: &str,
442 session_id: &str,
443 output: OutputFormat,
444) -> anyhow::Result<()> {
445 if session_id.trim().is_empty() {
446 anyhow::bail!("a browse session ID is required");
447 }
448 close_opcda_browse_session(bridge_host, session_id).await?;
449 if output == OutputFormat::Json {
450 println!(
451 "{}",
452 serde_json::to_string_pretty(&serde_json::json!({
453 "session_id": session_id,
454 "closed": true,
455 }))?
456 );
457 } else {
458 println!("Browse session closed successfully.");
459 }
460 Ok(())
461}
462
463async fn collect_browse_pages(
464 driver: &OpcDaDriver,
465 first: BrowsePage,
466 parent_node_key: Option<String>,
467 page_size: u32,
468 all: bool,
469) -> anyhow::Result<Vec<BrowsePage>> {
470 const MAX_PAGES: usize = 10_000;
471 let mut pages = vec![first];
472 if !all {
473 return Ok(pages);
474 }
475 while let Some(token) = pages.last().and_then(|page| page.next_page_token.clone()) {
476 if pages.len() >= MAX_PAGES {
477 anyhow::bail!("browse continuation exceeded the safety limit of {MAX_PAGES} pages");
478 }
479 let last = pages.last().expect("pages always contains the first page");
480 let page = driver
481 .browse(BrowsePageRequest::next(
482 last.session_id.clone(),
483 parent_node_key.clone(),
484 token,
485 page_size,
486 ))
487 .await?;
488 pages.push(page);
489 }
490 Ok(pages)
491}
492
493fn json_browse_page(page: &BrowsePage) -> serde_json::Value {
494 serde_json::json!({
495 "session_id": page.session_id,
496 "nodes": page.nodes.iter().map(json_browse_node).collect::<Vec<_>>(),
497 "next_page_token": page.next_page_token,
498 "complete": page.complete,
499 "organization": format!("{:?}", page.organization).to_lowercase(),
500 "source": format!("{:?}", page.source).to_lowercase(),
501 "warning": page.warning,
502 })
503}
504
505fn json_browse_node(node: &BrowseNode) -> serde_json::Value {
506 serde_json::json!({
507 "node_key": node.node_key,
508 "display_name": node.display_name,
509 "kind": browse_node_kind_name(node.kind),
510 "item_id": node.item_id,
511 })
512}
513
514fn browse_node_kind_name(kind: BrowseNodeKind) -> &'static str {
515 match kind {
516 BrowseNodeKind::Unspecified => "unspecified",
517 BrowseNodeKind::Branch => "branch",
518 BrowseNodeKind::Item => "item",
519 BrowseNodeKind::BranchAndItem => "branch_and_item",
520 }
521}
522
523#[derive(Debug)]
524struct SearchOptions {
525 query: String,
526 match_mode: OpcSearchMatchModeArg,
527 max_results: u32,
528 session_id: Option<String>,
529 scope_node_key: Option<String>,
530 include_branches: bool,
531 refresh: bool,
532}
533
534async fn search_with_output(
535 bridge_host: &str,
536 server: &str,
537 options: SearchOptions,
538 output: OutputFormat,
539) -> anyhow::Result<()> {
540 let driver = OpcDaDriver::connect(bridge_host, server).await?;
541 let request = SearchRequest {
542 query: options.query,
543 match_mode: options.match_mode.into(),
544 session_id: options.session_id,
545 scope_node_key: options.scope_node_key,
546 max_results: options.max_results,
547 include_branches: options.include_branches,
548 refresh: options.refresh,
549 };
550 let mut stream = driver.search_stream(request).await?;
551 let mut matches = Vec::new();
552 let mut progress = Vec::new();
553 let mut completed = None;
554 while let Some(event) = stream.next().await? {
555 match event {
556 SearchEvent::Match(found) => matches.push(found),
557 SearchEvent::Progress(current) => {
558 eprintln!(
559 "search progress: visited={} matches={}{}",
560 current.visited_nodes,
561 current.matches,
562 if current.partial { " (partial)" } else { "" }
563 );
564 progress.push(current);
565 }
566 SearchEvent::Completed(current) => completed = Some(current),
567 }
568 }
569 if output == OutputFormat::Json {
570 let matches = matches.iter().map(json_search_match).collect::<Vec<_>>();
571 let progress = progress
572 .iter()
573 .map(|current| {
574 serde_json::json!({
575 "visited_nodes": current.visited_nodes,
576 "matches": current.matches,
577 "partial": current.partial,
578 })
579 })
580 .collect::<Vec<_>>();
581 let completed = completed.map(|current| {
582 serde_json::json!({
583 "complete": current.complete,
584 "cancelled": current.cancelled,
585 "truncated": current.truncated,
586 "warning": current.warning,
587 })
588 });
589 println!(
590 "{}",
591 serde_json::to_string_pretty(&serde_json::json!({
592 "matches": matches,
593 "progress": progress,
594 "completed": completed,
595 }))?
596 );
597 return Ok(());
598 }
599 if matches.is_empty() {
600 println!("No matching tags found.");
601 } else {
602 for found in matches {
603 let item = found.node.item_id.as_deref().unwrap_or("-");
604 let breadcrumb = found
605 .breadcrumbs
606 .iter()
607 .map(|part| part.display_name.as_str())
608 .collect::<Vec<_>>()
609 .join("/");
610 println!("{item}\t{breadcrumb}");
611 }
612 }
613 let _ = completed.map(|done| {
614 if let Some(warning) = done.warning {
615 eprintln!("search warning: {warning}");
616 }
617 if done.truncated {
618 eprintln!("search results were truncated at the requested maximum");
619 }
620 });
621 Ok(())
622}
623
624fn json_search_match(found: &SearchMatch) -> serde_json::Value {
625 serde_json::json!({
626 "node": json_browse_node(&found.node),
627 "breadcrumbs": found.breadcrumbs.iter().map(|part| {
628 serde_json::json!({
629 "node_key": part.node_key,
630 "display_name": part.display_name,
631 })
632 }).collect::<Vec<_>>(),
633 })
634}
635
636fn print_search_index_status(
637 status: &SearchIndexStatus,
638 output: OutputFormat,
639) -> anyhow::Result<()> {
640 if output == OutputFormat::Json {
641 println!(
642 "{}",
643 serde_json::to_string_pretty(&json_search_index_status(status))?
644 );
645 return Ok(());
646 }
647 println!("Server: {}", status.server);
648 println!("State: {}", status.state);
649 println!("Auto-refresh enabled: {}", status.auto_refresh_enabled);
650 println!("Active generation: {}", status.active_generation);
651 println!("Entries: {}", status.entry_count);
652 println!("Unique items: {}", status.unique_item_count);
653 if let Some(progress) = &status.progress {
654 println!(
655 "Progress: {} entries, {:.1} items/s",
656 progress.entries_seen, progress.items_per_second
657 );
658 }
659 status.last_error.iter().for_each(|diagnostic| {
660 let label = if status.state != bhtune_driver::SearchIndexState::Failed {
661 "Last warning"
662 } else {
663 "Last error"
664 };
665 println!("{label}: {diagnostic}");
666 });
667 Ok(())
668}
669
670fn print_search_index_results(
671 response: &SearchIndexResponse,
672 output: OutputFormat,
673) -> anyhow::Result<()> {
674 if output == OutputFormat::Json {
675 println!(
676 "{}",
677 serde_json::to_string_pretty(&serde_json::json!({
678 "matches": response.matches.iter().map(json_indexed_search_match).collect::<Vec<_>>(),
679 "has_more": response.has_more,
680 "status": json_search_index_status(&response.status),
681 }))?
682 );
683 return Ok(());
684 }
685 if response.matches.is_empty() {
686 println!(
687 "No matching tags found (index state: {}).",
688 response.status.state
689 );
690 } else {
691 for found in &response.matches {
692 println!("{}\t{}", found.item_id, found.breadcrumbs.join("/"));
693 }
694 }
695 if response.has_more {
696 eprintln!("More than the requested number of matches exist; refine the query.");
697 }
698 if response.status.state != bhtune_driver::SearchIndexState::Ready {
699 eprintln!(
700 "Search results come from an {} namespace index.",
701 response.status.state
702 );
703 }
704 Ok(())
705}
706
707fn json_search_index_status(status: &SearchIndexStatus) -> serde_json::Value {
708 serde_json::json!({
709 "server": status.server,
710 "state": status.state.to_string(),
711 "auto_refresh_enabled": status.auto_refresh_enabled,
712 "active_generation": status.active_generation,
713 "entry_count": status.entry_count,
714 "unique_item_count": status.unique_item_count,
715 "started_at": status.started_at,
716 "completed_at": status.completed_at,
717 "last_error": status.last_error,
718 "database_bytes": status.database_bytes,
719 "organization": format!("{:?}", status.organization).to_lowercase(),
720 "source": format!("{:?}", status.source).to_lowercase(),
721 "scheduler": {
722 "next_refresh_at": status.scheduler.next_refresh_at,
723 "last_attempt_at": status.scheduler.last_attempt_at,
724 "last_success_at": status.scheduler.last_success_at,
725 "last_success_duration_ms": status.scheduler.last_success_duration_ms,
726 "retry_after": status.scheduler.retry_after,
727 "consecutive_failures": status.scheduler.consecutive_failures,
728 "circuit_open": status.scheduler.circuit_open,
729 },
730 "progress": status.progress.as_ref().map(|progress| serde_json::json!({
731 "branches_visited": progress.branches_visited,
732 "entries_seen": progress.entries_seen,
733 "unique_items": progress.unique_items,
734 "active_time_ms": progress.active_time_ms,
735 "paused_time_ms": progress.paused_time_ms,
736 "items_per_second": progress.items_per_second,
737 "estimated_remaining_ms": progress.estimated_remaining_ms,
738 })),
739 })
740}
741
742fn json_indexed_search_match(found: &bhtune_driver::IndexedSearchMatch) -> serde_json::Value {
743 serde_json::json!({
744 "item_id": found.item_id,
745 "display_name": found.display_name,
746 "kind": browse_node_kind_name(found.kind),
747 "breadcrumbs": found.breadcrumbs,
748 })
749}
750
751fn quality_name(quality: Quality) -> &'static str {
752 match quality {
753 Quality::Good => "good",
754 Quality::Uncertain => "uncertain",
755 Quality::Bad => "bad",
756 }
757}
758
759#[cfg(test)]
760mod tests {
761 use super::*;
762 use crate::args::SearchIndexControlActionArg;
763 use crate::test_support::{MockBridgeService, start_mock_server};
764 use opcda_bridge_proto::bridge::{
765 BrowseNode, BrowseNodeKind, BrowsePage, ListServersResponse, ReadResponse, SearchCompleted,
766 SearchEvent as ProtoSearchEvent, SearchIndexResponse as ProtoSearchIndexResponse,
767 SearchIndexState as ProtoSearchIndexState, SearchIndexStatus as ProtoSearchIndexStatus,
768 SearchProgress, TagValue as ProtoTagValue, WriteResponse, search_event,
769 };
770
771 #[tokio::test]
772 async fn servers_prints_every_registered_server_from_a_mock_gateway() {
773 let (host, server) = start_mock_server(MockBridgeService {
774 list_servers_response: ListServersResponse {
775 servers: vec![
776 "Matrikon.OPC.Simulation.1".to_string(),
777 "Kepware.KEPServerEX.V6".to_string(),
778 ],
779 },
780 ..Default::default()
781 })
782 .await;
783
784 servers(&host).await.unwrap();
785
786 server.shutdown().await;
787 }
788
789 #[tokio::test]
790 async fn servers_handles_an_empty_result() {
791 let (host, server) = start_mock_server(MockBridgeService::default()).await;
792 servers(&host).await.unwrap();
793 server.shutdown().await;
794 }
795
796 #[tokio::test]
797 async fn servers_connect_failure_surfaces_as_an_error() {
798 let err = servers("127.0.0.1:1").await.unwrap_err();
799 assert!(!err.to_string().is_empty());
800 }
801
802 #[tokio::test]
803 async fn read_prints_values_from_a_mock_gateway() {
804 let (host, server) = start_mock_server(MockBridgeService {
805 read_response: ReadResponse {
806 values: vec![ProtoTagValue {
807 tag_id: "Unit1.LIC101.PV".to_string(),
808 value: "42.5".to_string(),
809 quality: "Good".to_string(),
810 timestamp: "2024-01-15 10:23:45".to_string(),
811 }],
812 },
813 ..Default::default()
814 })
815 .await;
816
817 read(&host, "Sim.Server", &["Unit1.LIC101.PV".to_string()])
818 .await
819 .unwrap();
820
821 server.shutdown().await;
822 }
823
824 #[tokio::test]
825 async fn read_requires_at_least_one_tag() {
826 let err = read("127.0.0.1:1", "Sim.Server", &[]).await.unwrap_err();
827 assert!(err.to_string().contains("at least one tag"));
828 }
829
830 #[tokio::test]
831 async fn write_reports_success_from_a_mock_gateway() {
832 let (host, server) = start_mock_server(MockBridgeService {
833 write_response: WriteResponse {
834 tag_id: "Unit1.LIC101.OP".to_string(),
835 success: true,
836 error: None,
837 },
838 ..Default::default()
839 })
840 .await;
841
842 write(&host, "Sim.Server", "Unit1.LIC101.OP", "55.0")
843 .await
844 .unwrap();
845
846 server.shutdown().await;
847 }
848
849 #[tokio::test]
850 async fn write_surfaces_a_rejected_write_as_an_error() {
851 let (host, server) = start_mock_server(MockBridgeService {
852 write_response: WriteResponse {
853 tag_id: "Unit1.LIC101.OP".to_string(),
854 success: false,
855 error: Some("tag is read-only".to_string()),
856 },
857 ..Default::default()
858 })
859 .await;
860
861 let err = write(&host, "Sim.Server", "Unit1.LIC101.OP", "55.0")
862 .await
863 .unwrap_err();
864 assert!(err.to_string().contains("read-only"));
865
866 server.shutdown().await;
867 }
868
869 #[tokio::test]
870 async fn write_reports_a_gateway_rejection_when_it_omits_a_reason() {
871 let (host, server) = start_mock_server(MockBridgeService {
872 write_response: WriteResponse {
873 tag_id: "Unit1.LIC101.OP".to_string(),
874 success: false,
875 error: None,
876 },
877 ..Default::default()
878 })
879 .await;
880
881 let err = write(&host, "Sim.Server", "Unit1.LIC101.OP", "55.0")
882 .await
883 .unwrap_err();
884 assert!(err.to_string().contains("driver rejected the write"));
885
886 server.shutdown().await;
887 }
888
889 #[test]
890 fn diagnostic_formatters_cover_timestamp_and_reason_fallbacks() {
891 assert_eq!(format_timestamp(None), "-");
892 assert_eq!(
893 format_timestamp(Some(
894 chrono::DateTime::parse_from_rfc3339("2024-01-15T10:23:45Z")
895 .unwrap()
896 .into()
897 )),
898 "2024-01-15T10:23:45+00:00"
899 );
900 assert_eq!(write_rejection_reason(None), "unknown reason");
901 assert_eq!(
902 write_rejection_reason(Some("read-only".to_string())),
903 "read-only"
904 );
905 }
906
907 #[tokio::test]
908 async fn write_accepts_a_raw_non_numeric_value() {
909 let (host, server) = start_mock_server(MockBridgeService {
910 write_response: WriteResponse {
911 tag_id: "Unit1.LIC101.OP".to_string(),
912 success: true,
913 error: None,
914 },
915 ..Default::default()
916 })
917 .await;
918
919 write(&host, "Sim.Server", "Unit1.LIC101.MODE", "MAN")
920 .await
921 .unwrap();
922
923 server.shutdown().await;
924 }
925
926 #[tokio::test]
927 async fn browse_prints_nodes_from_a_mock_gateway() {
928 let (host, server) = start_mock_server(MockBridgeService {
929 browse_response: BrowsePage {
930 session_id: "session".to_string(),
931 nodes: vec![BrowseNode {
932 node_key: "unit1".to_string(),
933 display_name: "Unit1".to_string(),
934 kind: BrowseNodeKind::Branch as i32,
935 item_id: None,
936 }],
937 complete: true,
938 ..Default::default()
939 },
940 ..Default::default()
941 })
942 .await;
943
944 browse(&host, "Sim.Server", "").await.unwrap();
945
946 server.shutdown().await;
947 }
948
949 #[tokio::test]
950 async fn browse_handles_an_empty_result() {
951 let (host, server) = start_mock_server(MockBridgeService::default()).await;
952 browse(&host, "Sim.Server", "Unit1").await.unwrap();
953 server.shutdown().await;
954 }
955
956 #[tokio::test]
957 async fn browse_connect_failure_surfaces_as_an_error() {
958 let err = browse("127.0.0.1:1", "Sim.Server", "").await.unwrap_err();
959 assert!(!err.to_string().is_empty());
960 }
961
962 #[tokio::test]
963 async fn connect_failure_surfaces_as_an_error() {
964 let err = read(
966 "127.0.0.1:1",
967 "Sim.Server",
968 &["Unit1.LIC101.PV".to_string()],
969 )
970 .await
971 .unwrap_err();
972 assert!(!err.to_string().is_empty());
973 }
974
975 #[tokio::test]
976 async fn run_dispatches_servers_read_write_and_browse() {
977 let close_browse_session_calls = std::sync::Arc::new(std::sync::atomic::AtomicU32::new(0));
978 let (host, server) = start_mock_server(MockBridgeService {
979 list_servers_response: ListServersResponse {
980 servers: vec!["Matrikon.OPC.Simulation.1".to_string()],
981 },
982 read_response: ReadResponse {
983 values: vec![ProtoTagValue {
984 tag_id: "Unit1.LIC101.PV".to_string(),
985 value: "42.5".to_string(),
986 quality: "Good".to_string(),
987 timestamp: "2024-01-15 10:23:45".to_string(),
988 }],
989 },
990 write_response: WriteResponse {
991 tag_id: "Unit1.LIC101.OP".to_string(),
992 success: true,
993 error: None,
994 },
995 browse_response: BrowsePage {
996 session_id: "session".to_string(),
997 complete: true,
998 ..Default::default()
999 },
1000 close_browse_session_calls: close_browse_session_calls.clone(),
1001 ..Default::default()
1002 })
1003 .await;
1004 let config = crate::config::BhtuneConfig::default();
1005
1006 run(
1007 OpcCommand::Servers {
1008 bridge_host: Some(host.clone()),
1009 },
1010 &config,
1011 )
1012 .await
1013 .unwrap();
1014
1015 run(
1016 OpcCommand::Read {
1017 bridge_host: Some(host.clone()),
1018 server: Some("Sim.Server".to_string()),
1019 tags: vec!["Unit1.LIC101.PV".to_string()],
1020 },
1021 &config,
1022 )
1023 .await
1024 .unwrap();
1025
1026 run(
1027 OpcCommand::Write {
1028 bridge_host: Some(host.clone()),
1029 server: Some("Sim.Server".to_string()),
1030 tag: "Unit1.LIC101.OP".to_string(),
1031 value: "55.0".to_string(),
1032 },
1033 &config,
1034 )
1035 .await
1036 .unwrap();
1037
1038 run(
1039 OpcCommand::Browse {
1040 bridge_host: Some(host.clone()),
1041 server: Some("Sim.Server".to_string()),
1042 session_id: None,
1043 parent_node_key: None,
1044 page_token: None,
1045 page_size: bhtune_driver::DEFAULT_PAGE_SIZE,
1046 all: false,
1047 refresh: false,
1048 },
1049 &config,
1050 )
1051 .await
1052 .unwrap();
1053 assert_eq!(
1054 close_browse_session_calls.load(std::sync::atomic::Ordering::SeqCst),
1055 0
1056 );
1057
1058 run(
1059 OpcCommand::Close {
1060 bridge_host: Some(host),
1061 session_id: "session".to_string(),
1062 },
1063 &config,
1064 )
1065 .await
1066 .unwrap();
1067 assert_eq!(
1068 close_browse_session_calls.load(std::sync::atomic::Ordering::SeqCst),
1069 1
1070 );
1071
1072 server.shutdown().await;
1073 }
1074
1075 #[tokio::test]
1076 async fn run_resolves_bridge_host_and_server_from_config_when_cli_flags_are_unset() {
1077 let (host, server) = start_mock_server(MockBridgeService {
1078 read_response: ReadResponse {
1079 values: vec![ProtoTagValue {
1080 tag_id: "Unit1.LIC101.PV".to_string(),
1081 value: "42.5".to_string(),
1082 quality: "Good".to_string(),
1083 timestamp: "2024-01-15 10:23:45".to_string(),
1084 }],
1085 },
1086 ..Default::default()
1087 })
1088 .await;
1089 let config = crate::config::BhtuneConfig {
1090 bridge_host: Some(host),
1091 server: Some("Sim.Server".to_string()),
1092 ..Default::default()
1093 };
1094
1095 run(
1096 OpcCommand::Read {
1097 bridge_host: None,
1098 server: None,
1099 tags: vec!["Unit1.LIC101.PV".to_string()],
1100 },
1101 &config,
1102 )
1103 .await
1104 .unwrap();
1105
1106 server.shutdown().await;
1107 }
1108
1109 #[tokio::test]
1110 async fn run_resolves_bridge_host_from_config_for_servers_when_cli_flag_is_unset() {
1111 let (host, server) = start_mock_server(MockBridgeService::default()).await;
1112 let config = crate::config::BhtuneConfig {
1113 bridge_host: Some(host),
1114 ..Default::default()
1115 };
1116
1117 run(OpcCommand::Servers { bridge_host: None }, &config)
1118 .await
1119 .unwrap();
1120
1121 server.shutdown().await;
1122 }
1123
1124 #[tokio::test]
1125 async fn run_errors_when_server_is_unset_in_both_cli_and_config() {
1126 let err = run(
1127 OpcCommand::Read {
1128 bridge_host: None,
1129 server: None,
1130 tags: vec!["Unit1.LIC101.PV".to_string()],
1131 },
1132 &crate::config::BhtuneConfig::default(),
1133 )
1134 .await
1135 .unwrap_err();
1136 assert!(err.to_string().contains("no OPC server specified"));
1137 }
1138
1139 #[tokio::test]
1140 async fn search_and_index_commands_cover_table_json_and_control_paths() {
1141 let status = ProtoSearchIndexStatus {
1142 server: "Sim.Server".into(),
1143 state: ProtoSearchIndexState::Partial as i32,
1144 configured: true,
1145 active_generation: 2,
1146 entry_count: 3,
1147 unique_item_count: 2,
1148 started_at: Some("2026-01-01T00:00:00Z".into()),
1149 completed_at: None,
1150 last_error: Some("inventory incomplete".into()),
1151 database_bytes: 42,
1152 progress: Some(opcda_bridge_proto::bridge::IndexedSearchProgress {
1153 branches_visited: 1,
1154 entries_seen: 3,
1155 unique_items: 2,
1156 active_time_ms: 10,
1157 paused_time_ms: 2,
1158 items_per_second: 1.5,
1159 estimated_remaining_ms: Some(5),
1160 }),
1161 ..Default::default()
1162 };
1163 let (host, server) = start_mock_server(MockBridgeService {
1164 search_events: vec![
1165 ProtoSearchEvent {
1166 event: Some(search_event::Event::Progress(SearchProgress {
1167 visited_nodes: 4,
1168 matches: 1,
1169 partial: true,
1170 })),
1171 },
1172 ProtoSearchEvent {
1173 event: Some(search_event::Event::Match(
1174 opcda_bridge_proto::bridge::SearchMatch {
1175 node: Some(BrowseNode {
1176 node_key: "pv".into(),
1177 display_name: "PV".into(),
1178 kind: BrowseNodeKind::Item as i32,
1179 item_id: Some("Area.PV".into()),
1180 }),
1181 breadcrumbs: vec![opcda_bridge_proto::bridge::BrowseBreadcrumb {
1182 node_key: "area".into(),
1183 display_name: "Area".into(),
1184 }],
1185 },
1186 )),
1187 },
1188 ProtoSearchEvent {
1189 event: Some(search_event::Event::Completed(SearchCompleted {
1190 complete: false,
1191 cancelled: false,
1192 truncated: true,
1193 warning: Some("partial result".into()),
1194 })),
1195 },
1196 ],
1197 search_index_status_response: status.clone(),
1198 search_index_response: ProtoSearchIndexResponse {
1199 matches: vec![],
1200 has_more: true,
1201 status: Some(status),
1202 },
1203 ..Default::default()
1204 })
1205 .await;
1206 let config = crate::config::BhtuneConfig::default();
1207
1208 run_with_output(
1209 OpcCommand::Search {
1210 bridge_host: Some(host.clone()),
1211 server: Some("Sim.Server".into()),
1212 query: "PV".into(),
1213 match_mode: OpcSearchMatchModeArg::Contains,
1214 max_results: 10,
1215 session_id: Some("session".into()),
1216 scope_node_key: Some("area".into()),
1217 include_branches: true,
1218 refresh: true,
1219 },
1220 &config,
1221 OutputFormat::Table,
1222 )
1223 .await
1224 .unwrap();
1225 run_with_output(
1226 OpcCommand::Search {
1227 bridge_host: Some(host.clone()),
1228 server: Some("Sim.Server".into()),
1229 query: "PV".into(),
1230 match_mode: OpcSearchMatchModeArg::Exact,
1231 max_results: 10,
1232 session_id: None,
1233 scope_node_key: None,
1234 include_branches: false,
1235 refresh: false,
1236 },
1237 &config,
1238 OutputFormat::Json,
1239 )
1240 .await
1241 .unwrap();
1242
1243 for command in [
1244 SearchIndexCommand::Status {
1245 bridge_host: Some(host.clone()),
1246 server: Some("Sim.Server".into()),
1247 },
1248 SearchIndexCommand::Refresh {
1249 bridge_host: Some(host.clone()),
1250 server: Some("Sim.Server".into()),
1251 force: true,
1252 },
1253 SearchIndexCommand::Control {
1254 bridge_host: Some(host.clone()),
1255 server: Some("Sim.Server".into()),
1256 action: SearchIndexControlActionArg::Pause,
1257 },
1258 ] {
1259 run_with_output(
1260 OpcCommand::SearchIndex { command },
1261 &config,
1262 OutputFormat::Table,
1263 )
1264 .await
1265 .unwrap();
1266 }
1267 run_with_output(
1268 OpcCommand::SearchIndex {
1269 command: SearchIndexCommand::Search {
1270 bridge_host: Some(host.clone()),
1271 server: Some("Sim.Server".into()),
1272 query: "PV".into(),
1273 match_mode: OpcSearchMatchModeArg::Prefix,
1274 max_results: 5,
1275 },
1276 },
1277 &config,
1278 OutputFormat::Json,
1279 )
1280 .await
1281 .unwrap();
1282
1283 server.shutdown().await;
1284 }
1285
1286 #[tokio::test]
1287 async fn diagnostic_commands_cover_json_and_validation_branches() {
1288 let status = bhtune_driver::SearchIndexStatus {
1289 server: "Sim.Server".into(),
1290 state: bhtune_driver::SearchIndexState::Failed,
1291 auto_refresh_enabled: true,
1292 active_generation: 1,
1293 entry_count: 2,
1294 unique_item_count: 1,
1295 started_at: None,
1296 completed_at: None,
1297 last_error: Some("failed".into()),
1298 database_bytes: 10,
1299 organization: bhtune_driver::NamespaceOrganization::Flat,
1300 source: bhtune_driver::BrowseSource::Flat,
1301 progress: Some(bhtune_driver::IndexedSearchProgress {
1302 branches_visited: 1,
1303 entries_seen: 2,
1304 unique_items: 1,
1305 active_time_ms: 1,
1306 paused_time_ms: 0,
1307 items_per_second: 2.0,
1308 estimated_remaining_ms: None,
1309 }),
1310 scheduler: bhtune_driver::IndexSchedulerDiagnostics::default(),
1311 };
1312 let indexed = bhtune_driver::SearchIndexResponse {
1313 matches: vec![bhtune_driver::IndexedSearchMatch {
1314 item_id: "Area.PV".into(),
1315 display_name: "PV".into(),
1316 kind: bhtune_driver::BrowseNodeKind::BranchAndItem,
1317 breadcrumbs: vec!["Area".into()],
1318 }],
1319 has_more: true,
1320 status: status.clone(),
1321 };
1322 assert_eq!(quality_name(Quality::Good), "good");
1323 assert_eq!(quality_name(Quality::Uncertain), "uncertain");
1324 assert_eq!(quality_name(Quality::Bad), "bad");
1325 print_search_index_status(&status, OutputFormat::Table).unwrap();
1326 print_search_index_status(&status, OutputFormat::Json).unwrap();
1327 print_search_index_results(&indexed, OutputFormat::Table).unwrap();
1328 print_search_index_results(&indexed, OutputFormat::Json).unwrap();
1329 print_search_index_results(
1330 &bhtune_driver::SearchIndexResponse {
1331 matches: vec![],
1332 has_more: false,
1333 status,
1334 },
1335 OutputFormat::Table,
1336 )
1337 .unwrap();
1338
1339 let host_service = MockBridgeService {
1340 read_response: ReadResponse {
1341 values: vec![ProtoTagValue {
1342 tag_id: "Area.PV".into(),
1343 value: "1.5".into(),
1344 quality: "Uncertain".into(),
1345 timestamp: "N/A".into(),
1346 }],
1347 },
1348 write_response: WriteResponse {
1349 tag_id: "Area.MV".into(),
1350 success: false,
1351 error: None,
1352 },
1353 browse_response: BrowsePage {
1354 session_id: "s".into(),
1355 nodes: vec![
1356 BrowseNode {
1357 node_key: "u".into(),
1358 display_name: "Unspecified".into(),
1359 kind: BrowseNodeKind::Unspecified as i32,
1360 item_id: None,
1361 },
1362 BrowseNode {
1363 node_key: "i".into(),
1364 display_name: "Item".into(),
1365 kind: BrowseNodeKind::Item as i32,
1366 item_id: Some("Area.PV".into()),
1367 },
1368 BrowseNode {
1369 node_key: "b".into(),
1370 display_name: "Branch".into(),
1371 kind: BrowseNodeKind::Branch as i32,
1372 item_id: None,
1373 },
1374 BrowseNode {
1375 node_key: "both".into(),
1376 display_name: "Both".into(),
1377 kind: BrowseNodeKind::BranchAndItem as i32,
1378 item_id: Some("Area.Both".into()),
1379 },
1380 ],
1381 next_page_token: Some("next".into()),
1382 complete: false,
1383 ..Default::default()
1384 },
1385 browse_continuation_response: Some(BrowsePage {
1386 session_id: "s".into(),
1387 complete: true,
1388 ..Default::default()
1389 }),
1390 list_servers_response: ListServersResponse {
1391 servers: vec!["Sim.Server".into()],
1392 },
1393 ..Default::default()
1394 };
1395 let (host, server) = start_mock_server(host_service).await;
1396 servers_with_output(&host, OutputFormat::Json)
1397 .await
1398 .unwrap();
1399 read_with_output(&host, "Sim.Server", &["Area.PV".into()], OutputFormat::Json)
1400 .await
1401 .unwrap();
1402 write_with_output(&host, "Sim.Server", "Area.MV", "1.0", OutputFormat::Json)
1403 .await
1404 .unwrap_err();
1405 let err = write_with_output(
1406 &host,
1407 "Sim.Server",
1408 "Area.MV",
1409 "not-numeric",
1410 OutputFormat::Json,
1411 )
1412 .await
1413 .unwrap_err();
1414 assert!(err.to_string().contains("driver rejected"));
1415 browse_with_output(
1416 &host,
1417 "Sim.Server",
1418 BrowseOptions {
1419 page_size: 2,
1420 all: true,
1421 ..Default::default()
1422 },
1423 OutputFormat::Json,
1424 )
1425 .await
1426 .unwrap();
1427 browse_with_output(
1428 &host,
1429 "Sim.Server",
1430 BrowseOptions {
1431 page_size: 2,
1432 all: false,
1433 ..Default::default()
1434 },
1435 OutputFormat::Json,
1436 )
1437 .await
1438 .unwrap();
1439 browse_with_output(
1440 &host,
1441 "Sim.Server",
1442 BrowseOptions {
1443 page_size: 2,
1444 all: true,
1445 ..Default::default()
1446 },
1447 OutputFormat::Table,
1448 )
1449 .await
1450 .unwrap();
1451 browse_with_output(
1452 &host,
1453 "Sim.Server",
1454 BrowseOptions {
1455 page_size: 2,
1456 all: false,
1457 ..Default::default()
1458 },
1459 OutputFormat::Table,
1460 )
1461 .await
1462 .unwrap();
1463 let err = browse_with_output(
1464 &host,
1465 "Sim.Server",
1466 BrowseOptions {
1467 parent_node_key: Some("parent".into()),
1468 ..Default::default()
1469 },
1470 OutputFormat::Table,
1471 )
1472 .await
1473 .unwrap_err();
1474 assert!(err.to_string().contains("--session-id"));
1475 let err = close_with_output(&host, " ", OutputFormat::Table)
1476 .await
1477 .unwrap_err();
1478 assert!(err.to_string().contains("session ID"));
1479 close_with_output(&host, "s", OutputFormat::Json)
1480 .await
1481 .unwrap();
1482 server.shutdown().await;
1483
1484 let (host, server) = start_mock_server(MockBridgeService {
1485 write_response: WriteResponse {
1486 tag_id: "Area.MV".into(),
1487 success: true,
1488 error: None,
1489 },
1490 ..Default::default()
1491 })
1492 .await;
1493 write_with_output(&host, "Sim.Server", "Area.MV", "1.0", OutputFormat::Json)
1494 .await
1495 .unwrap();
1496 server.shutdown().await;
1497 }
1498
1499 #[tokio::test]
1500 async fn table_search_reports_empty_results_and_completion_diagnostics() {
1501 let (host, server) = start_mock_server(MockBridgeService {
1502 search_events: vec![ProtoSearchEvent {
1503 event: Some(search_event::Event::Completed(SearchCompleted {
1504 complete: false,
1505 cancelled: false,
1506 truncated: true,
1507 warning: Some("partial".into()),
1508 })),
1509 }],
1510 ..Default::default()
1511 })
1512 .await;
1513 search_with_output(
1514 &host,
1515 "Sim.Server",
1516 SearchOptions {
1517 query: "PV".into(),
1518 match_mode: OpcSearchMatchModeArg::Contains,
1519 max_results: 1,
1520 session_id: None,
1521 scope_node_key: None,
1522 include_branches: false,
1523 refresh: false,
1524 },
1525 OutputFormat::Table,
1526 )
1527 .await
1528 .unwrap();
1529 server.shutdown().await;
1530 }
1531
1532 #[tokio::test]
1533 async fn browse_all_stops_at_the_safety_page_limit() {
1534 let (host, server) = start_mock_server(MockBridgeService {
1535 browse_response: BrowsePage {
1536 session_id: "s".into(),
1537 next_page_token: Some("next".into()),
1538 complete: false,
1539 ..Default::default()
1540 },
1541 browse_continuation_response: Some(BrowsePage {
1542 session_id: "s".into(),
1543 next_page_token: Some("next".into()),
1544 complete: false,
1545 ..Default::default()
1546 }),
1547 ..Default::default()
1548 })
1549 .await;
1550 let err = browse_with_output(
1551 &host,
1552 "Sim.Server",
1553 BrowseOptions {
1554 all: true,
1555 ..Default::default()
1556 },
1557 OutputFormat::Table,
1558 )
1559 .await
1560 .unwrap_err();
1561 assert!(err.to_string().contains("safety limit"));
1562 server.shutdown().await;
1563 }
1564}