1use std::{collections::VecDeque, future::Future};
2
3use futures::{stream::FuturesUnordered, Stream, StreamExt};
4use hashbrown::{Equivalent, HashSet};
5use opcua_types::{BrowseDescription, BrowseDirection, ByteString, Error, NodeId, StatusCode};
6
7use crate::{
8 session::{Browse, BrowseNext},
9 RequestRetryPolicy, Session,
10};
11
12use super::{result::BrowserResult, BrowseResultItem, Browser, BrowserPolicy, RequestWithRetries};
13
14impl<'a, T: BrowserPolicy + 'a, R: RequestRetryPolicy + Clone + 'a> Browser<'a, T, R> {
15 pub fn run(
19 self,
20 initial: Vec<BrowseDescription>,
21 ) -> impl Stream<Item = Result<BrowseResultItem, Error>> + 'a {
22 let initial_exec = BrowserExecution {
29 browser: self,
30 running_browses: FuturesUnordered::new(),
31 pending: initial
32 .into_iter()
33 .map(|r| RequestWithRetries {
34 request: r,
35 num_outer_retries: 0,
36 depth: 0,
37 })
38 .collect(),
39 pending_out: VecDeque::new(),
40 browsed_nodes: HashSet::new(),
41 pending_continuation_points: HashSet::new(),
42 };
43 futures::stream::try_unfold(initial_exec, |mut s| async move {
44 loop {
45 if s.browser.token.is_cancelled() {
47 s.cleanup().await;
48 return Ok(None);
49 }
50
51 if let Some(n) = s.pending_out.pop_front() {
53 return Ok(Some((n, s)));
54 }
55
56 while s.running_browses.len() < s.browser.config.max_concurrent_requests
59 || s.browser.config.max_concurrent_requests == 0
60 {
61 let mut chunk = Vec::new();
62 while chunk.len() < s.browser.config.max_nodes_per_request
63 || s.browser.config.max_nodes_per_request == 0
64 {
65 let Some(it) = s.pending.pop_front() else {
66 break;
67 };
68 chunk.push(it);
69 }
70 if !chunk.is_empty() {
71 s.running_browses.push(run_browse(
72 s.browser.session,
73 BrowseBatch::Browse(chunk),
74 s.browser.retry_policy.clone(),
75 s.browser.config.max_references_per_node,
76 ));
77 } else {
78 break;
79 }
80 }
81
82 let Some(next) = s.running_browses.next().await else {
84 return Ok(None);
85 };
86
87 let browse_next = match next.and_then(|m| s.process_result(m)) {
89 Ok(next) => next,
90 Err(e) => {
91 s.cleanup().await;
92 return Err(e);
93 }
94 };
95
96 if !browse_next.is_empty() {
100 s.running_browses.push(run_browse(
101 s.browser.session,
102 BrowseBatch::Next(browse_next),
103 s.browser.retry_policy.clone(),
104 s.browser.config.max_references_per_node,
105 ));
106 }
107 }
108 })
109 }
110
111 pub async fn run_into_result(
113 self,
114 initial: Vec<BrowseDescription>,
115 ) -> Result<BrowserResult, Error> {
116 BrowserResult::build_from_browser(self.run(initial)).await
117 }
118}
119
120struct BrowserExecution<'a, T, R, TFut> {
125 browser: Browser<'a, T, R>,
126 running_browses: FuturesUnordered<TFut>,
127 pending: VecDeque<RequestWithRetries>,
128 pending_out: VecDeque<BrowseResultItem>,
129 browsed_nodes: HashSet<BrowsedNode>,
130 pending_continuation_points: HashSet<ByteString>,
131}
132
133#[derive(PartialEq, Eq, Hash)]
134enum Direction {
135 Forward,
136 Inverse,
137 Both,
138}
139
140#[derive(PartialEq, Eq, Hash)]
141struct BrowsedNode {
142 id: NodeId,
143 direction: Direction,
144}
145
146impl From<&BrowseDescription> for BrowsedNode {
147 fn from(value: &BrowseDescription) -> Self {
148 Self {
149 id: value.node_id.clone(),
150 direction: match value.browse_direction {
151 BrowseDirection::Forward => Direction::Forward,
152 BrowseDirection::Inverse => Direction::Inverse,
153 BrowseDirection::Both => Direction::Both,
154 BrowseDirection::Invalid => Direction::Both,
155 },
156 }
157 }
158}
159
160#[derive(PartialEq, Eq, Hash)]
161struct BrowsedNodeRef<'a> {
162 id: &'a NodeId,
163 direction: Direction,
164}
165
166impl Equivalent<BrowsedNode> for BrowsedNodeRef<'_> {
167 fn equivalent(&self, key: &BrowsedNode) -> bool {
168 self.id == &key.id && self.direction == key.direction
169 }
170}
171
172struct InnerResultItem {
173 it: BrowseResultItem,
174 cp: Option<ByteString>,
175}
176
177struct BrowseNextItem {
178 original: RequestWithRetries,
179 cp: ByteString,
180}
181
182enum BrowseBatch {
183 Browse(Vec<RequestWithRetries>),
184 Next(Vec<BrowseNextItem>),
185}
186
187async fn run_browse<R: RequestRetryPolicy>(
188 session: &Session,
189 batch: BrowseBatch,
190 policy: R,
191 max_references_per_node: u32,
192) -> Result<Vec<InnerResultItem>, Error> {
193 match batch {
194 BrowseBatch::Browse(items) => {
195 let r = session
196 .send_with_retry(
197 Browse::new(session)
198 .max_references_per_node(max_references_per_node)
199 .nodes_to_browse(items.iter().map(|r| r.request.clone()).collect()),
200 policy,
201 )
202 .await?;
203
204 let res = r.results.unwrap_or_default();
205 if res.len() != items.len() {
206 return Err(Error::new(
207 StatusCode::BadUnexpectedError,
208 format!(
209 "Incorrect number of results returned from Browse, expected {}, got {}",
210 items.len(),
211 res.len()
212 ),
213 ));
214 }
215 Ok(res
216 .into_iter()
217 .zip(items)
218 .map(|(res, it)| InnerResultItem {
219 it: BrowseResultItem {
220 status: res.status_code,
221 request: it,
222 references: res.references.unwrap_or_default(),
223 request_continuation_point: None,
224 },
225 cp: if res.continuation_point.is_null() {
226 None
227 } else {
228 Some(res.continuation_point)
229 },
230 })
231 .collect())
232 }
233 BrowseBatch::Next(items) => {
234 let r = session
235 .send_with_retry(
236 BrowseNext::new(session)
237 .continuation_points(items.iter().map(|r| r.cp.clone()).collect()),
238 policy,
239 )
240 .await?;
241
242 let res = r.results.unwrap_or_default();
243 if res.len() != items.len() {
244 return Err(Error::new(
245 StatusCode::BadUnexpectedError,
246 format!(
247 "Incorrect number of results returned from BrowseNext, expected {}, got {}",
248 items.len(),
249 res.len()
250 ),
251 ));
252 }
253 Ok(res
254 .into_iter()
255 .zip(items)
256 .map(|(res, it)| InnerResultItem {
257 it: BrowseResultItem {
258 status: res.status_code,
259 request: it.original,
260 references: res.references.unwrap_or_default(),
261 request_continuation_point: Some(it.cp),
262 },
263 cp: if res.continuation_point.is_null() {
264 None
265 } else {
266 Some(res.continuation_point)
267 },
268 })
269 .collect())
270 }
271 }
272}
273
274impl<
275 'a,
276 T: BrowserPolicy,
277 R: RequestRetryPolicy + Clone + 'a,
278 TFut: Future<Output = Result<Vec<InnerResultItem>, Error>> + 'a,
279 > BrowserExecution<'a, T, R, TFut>
280{
281 fn visited(&self, dir: Direction, node_id: &NodeId) -> bool {
282 self.browsed_nodes.contains(&BrowsedNodeRef {
283 id: node_id,
284 direction: dir,
285 })
286 }
287
288 async fn cleanup(&mut self) {
289 while let Some(r) = self.running_browses.next().await {
291 let Ok(r) = r else {
292 continue;
293 };
294 for res in r {
295 if let Some(old_cp) = res.it.request_continuation_point.as_ref() {
296 self.pending_continuation_points.remove(old_cp);
297 }
298 if let Some(cp) = res.cp {
299 self.pending_continuation_points.insert(cp.clone());
300 }
301 }
302 }
303
304 let to_consume: Vec<_> = self.pending_continuation_points.drain().collect();
306 let mut futures = Vec::new();
307 for chunk in to_consume.chunks(self.browser.config.max_nodes_per_request) {
308 let session = self.browser.session;
309 futures.push(async move {
310 let _ = session.browse_next(true, chunk).await;
312 });
313 }
314
315 let mut it = FuturesUnordered::new();
317 loop {
318 while it.len() < self.browser.config.max_concurrent_requests
319 || self.browser.config.max_concurrent_requests == 0
320 {
321 let Some(fut) = futures.pop() else {
322 break;
323 };
324 it.push(fut);
325 }
326
327 if it.next().await.is_none() {
328 break;
329 }
330 }
331 }
332
333 pub(crate) fn process_result(
334 &mut self,
335 next: Vec<InnerResultItem>,
336 ) -> Result<Vec<BrowseNextItem>, Error> {
337 let mut browse_next = Vec::new();
338 for res in next {
339 let to_enqueue = self.browser.handler.get_next(&res.it);
341 for mut it in to_enqueue {
342 match it.browse_direction {
344 BrowseDirection::Forward => {
345 if self.visited(Direction::Forward, &it.node_id)
346 || self.visited(Direction::Both, &it.node_id)
347 {
348 continue;
349 }
350 }
351 BrowseDirection::Inverse => {
352 if self.visited(Direction::Inverse, &it.node_id)
353 || self.visited(Direction::Both, &it.node_id)
354 {
355 continue;
356 }
357 }
358 BrowseDirection::Both => {
359 if self.visited(Direction::Both, &it.node_id) {
360 continue;
361 }
362
363 let visited_inv = self.visited(Direction::Inverse, &it.node_id);
364 let visited_for = self.visited(Direction::Forward, &it.node_id);
365 if visited_for && visited_inv {
366 continue;
367 } else if visited_for {
368 it.browse_direction = BrowseDirection::Inverse;
369 } else if visited_inv {
370 it.browse_direction = BrowseDirection::Forward;
371 }
372 }
373 BrowseDirection::Invalid => {
374 return Err(Error::new(
375 StatusCode::BadBrowseDirectionInvalid,
376 "Produced an invalid browse direction",
377 ))
378 }
379 }
380
381 self.browsed_nodes.insert(BrowsedNode::from(&it));
383 self.pending.push_back(RequestWithRetries {
384 request: it,
385 num_outer_retries: 0,
386 depth: res.it.request.depth + 1,
387 });
388 }
389 if let Some(old_cp) = res.it.request_continuation_point.as_ref() {
391 self.pending_continuation_points.remove(old_cp);
392 }
393 if let Some(cp) = res.cp {
394 self.pending_continuation_points.insert(cp.clone());
397 browse_next.push(BrowseNextItem {
398 original: res.it.request.clone(),
399 cp,
400 });
401 } else if matches!(res.it.status, StatusCode::BadContinuationPointInvalid)
402 && self.browser.config.max_continuation_point_retries
403 > res.it.request.num_outer_retries
404 {
405 self.pending.push_back(RequestWithRetries {
408 request: res.it.request.request.clone(),
409 num_outer_retries: res.it.request.num_outer_retries + 1,
410 depth: res.it.request.depth,
411 });
412 }
413
414 self.pending_out.push_back(res.it);
415 }
416 Ok(browse_next)
417 }
418}