unb-runtime 2.0.3

unb session runtime: transport codec, session/writer engine, Wire handle, cancellation
Documentation
use std::collections::BTreeSet;
use std::pin::Pin;

use futures_util::stream::{unfold, Stream};
use unb_core::{Detail, DiscoverEvent, DiscoverPlan, Mode, Scope, DEFAULT_HOPS};

use crate::error::WsError;
use crate::wire::Wire;

impl Wire {
    /// Open a `Discover` walk and stream its typed catalog events.
    ///
    /// Yields one [`DiscoverEvent`] per graph observation — `NodeCatalog`, `Edge`,
    /// `Warning` — and finally the `Done` marker (carrying the `discover_id`),
    /// after which the stream ends. An error terminal or a closed session ends the
    /// stream without a `Done`.
    pub async fn discover_catalog(
        &self,
        target_path: &str,
        detail: Detail,
        scope: Scope,
    ) -> Result<Pin<Box<dyn Stream<Item = DiscoverEvent> + '_>>, WsError> {
        let plan = DiscoverPlan {
            discover_id: String::new(),
            detail,
            scope,
            hops: DEFAULT_HOPS,
            visited: BTreeSet::new(),
            timeout_ms: None,
            mode: Mode::PartialOk,
        };
        let stream = self.client_session().discover(target_path, plan).await?;

        Ok(Box::pin(unfold(
            (stream, false),
            |(mut stream, done)| async move {
                if done {
                    return None;
                }
                while let Ok(Some(envelope)) = stream.next().await {
                    match serde_json::from_slice::<DiscoverEvent>(&envelope.payload) {
                        Ok(marker @ DiscoverEvent::Done { .. }) => {
                            return Some((marker, (stream, true)))
                        }
                        Ok(event) => return Some((event, (stream, false))),
                        Err(_) => continue,
                    }
                }
                None
            },
        )))
    }
}