use crate::backend::connector::{ComConnector, ServerConnector};
use crate::com_worker::{ComRequest, ComWorker};
use crate::opc_da::errors::OpcResult;
use crate::provider::{
BrowseCapabilities, BrowsePage, BrowsePageRequest, BrowseSessionToken, OpcProvider, OpcValue,
TagValue, WriteResult,
};
use async_trait::async_trait;
use std::sync::Arc;
use std::sync::atomic::AtomicUsize;
pub struct OpcDaClient<C: ServerConnector + 'static = ComConnector> {
pub worker: ComWorker<C>,
}
impl Default for OpcDaClient<ComConnector> {
fn default() -> Self {
match Self::new(ComConnector) {
Ok(client) => client,
Err(err) => {
tracing::error!(error = ?err, "Failed to initialize default OpcDaClient");
Self {
worker: ComWorker::closed(),
}
}
}
}
}
impl<C: ServerConnector + 'static> OpcDaClient<C> {
pub fn new(connector: C) -> OpcResult<Self> {
tracing::info!("Initializing OpcDaClient...");
let worker = ComWorker::start(Arc::new(connector))?;
tracing::info!("OpcDaClient initialized successfully");
Ok(Self { worker })
}
}
#[allow(clippy::too_many_lines)]
#[async_trait]
impl<C: ServerConnector + 'static> OpcProvider for OpcDaClient<C> {
async fn list_servers(&self, host: &str) -> OpcResult<Vec<String>> {
let host_owned = host.to_string();
self.worker
.send_request(|reply| ComRequest::ListServers {
host: host_owned,
reply,
})
.await
}
async fn browse_tags(
&self,
server: &str,
max_tags: usize,
progress: Arc<AtomicUsize>,
tags_sink: Arc<std::sync::Mutex<Vec<String>>>,
) -> OpcResult<Vec<String>> {
let server_owned = server.to_string();
self.worker
.send_request(|reply| ComRequest::BrowseTags {
server: server_owned,
max_tags,
progress,
tags_sink,
reply,
})
.await
}
async fn browse_capabilities(&self, server: &str) -> OpcResult<BrowseCapabilities> {
let server_owned = server.to_string();
self.worker
.send_request(|reply| ComRequest::BrowseCapabilities {
server: server_owned,
reply,
})
.await
}
async fn open_browse_session(&self, server: &str) -> OpcResult<BrowseSessionToken> {
let server_owned = server.to_string();
self.worker
.send_request(|reply| ComRequest::OpenBrowseSession {
server: server_owned,
reply,
})
.await
}
async fn browse_page(
&self,
session: &BrowseSessionToken,
request: BrowsePageRequest,
) -> OpcResult<BrowsePage> {
let session = *session;
self.worker
.send_request(|reply| ComRequest::BrowsePage {
session,
request,
reply,
})
.await
}
async fn close_browse_session(&self, session: &BrowseSessionToken) -> OpcResult<()> {
let session = *session;
self.worker
.send_request(|reply| ComRequest::CloseBrowseSession { session, reply })
.await
}
async fn read_tag_values(
&self,
server: &str,
tag_ids: Vec<String>,
) -> OpcResult<Vec<TagValue>> {
let server_owned = server.to_string();
self.worker
.send_request(|reply| ComRequest::ReadTagValues {
server: server_owned,
tag_ids,
reply,
})
.await
}
async fn write_tag_value(
&self,
server: &str,
tag_id: &str,
value: OpcValue,
) -> OpcResult<WriteResult> {
let server_owned = server.to_string();
let tag_id_owned = tag_id.to_string();
self.worker
.send_request(|reply| ComRequest::WriteTagValue {
server: server_owned,
tag_id: tag_id_owned,
value,
reply,
})
.await
}
}