Skip to main content

zenkey_fleet/bus/
admin.rs

1//! Zenoh admin-space access (issue #14): browse `@/**` — the middleware's
2//! own introspection — from the same un-namespaced session the convention
3//! tooling already holds (a namespaced session's admin selector would be
4//! rewritten and match nothing, RFC 09 §5).
5//!
6//! Admin key layouts vary between zenoh versions (report §3.1's caveat), so
7//! this module stays a thin, honest transport: keys + JSON values, no
8//! hardcoded schema. `routers` extracts the few fields every 1.x layout
9//! carries, and leaves the rest visible in `raw`.
10//!
11//! **Base-less by design.** Everything here but [`origin_attachments`] takes a
12//! bare `&Session`, not a [`crate::Fleet`]: `@/**` is the middleware's own
13//! space and sits outside every deployment namespace, so there is no base for
14//! these calls to run *against* — a `Fleet` would offer one they must ignore.
15//! [`origin_attachments`] is the exception because it joins admin tokens back
16//! onto convention keys, which only parse under a base.
17
18use std::time::Duration;
19
20use crate::{Error, Result};
21use zenoh::Session;
22
23use crate::bus::query::GetOpts;
24use crate::report::{
25    Coverage, CoverageRow, DeclaredEntities, DeclaredEntity, EntityKind, MeshLink,
26    OriginAttachment, RouterInfo, StorageInfo, TopologyEdge, TopologyNode, TopologyReport,
27};
28
29/// One admin-space entry.
30#[derive(Debug, Clone)]
31pub struct AdminEntry {
32    pub key: String,
33    pub value: serde_json::Value,
34}
35
36/// GET an admin selector (default `@/**`). Fans to every node (target All,
37/// consolidation None — several routers may answer).
38///
39/// Bounded at [`crate::DEFAULT_MAX_REPLIES`] (#339): `@/**` against a router
40/// with a large storage is a lot of replies, each holding a payload. To see
41/// what the bound cost — or to raise it — use [`admin_get_within`], which
42/// takes the options the count rides on.
43pub async fn admin_get(
44    session: &Session,
45    selector: &str,
46    timeout: Duration,
47) -> Result<Vec<AdminEntry>> {
48    admin_get_within(session, selector, &GetOpts::new(timeout)).await
49}
50
51/// [`admin_get`] under the caller's own options — the reply bound and, after
52/// the call, what it cost ([`GetOpts::elided`], RFC 13 §3 O6).
53pub async fn admin_get_within(
54    session: &Session,
55    selector: &str,
56    opts: &GetOpts,
57) -> Result<Vec<AdminEntry>> {
58    let replies = crate::bus::query::disciplined_get(session, selector, opts)
59        .await
60        .map_err(|e| Error::bus("admin get", selector, e))?;
61    let mut out = Vec::new();
62    let mut elided = 0u64;
63    while let Ok(reply) = replies.recv_async().await {
64        // Past the bound the replies are drained but not kept: the count
65        // stays exact, the memory stays bounded.
66        if out.len() >= opts.reply_bound() {
67            elided += 1;
68            continue;
69        }
70        let Ok(sample) = reply.result() else { continue };
71        let bytes = sample.payload().to_bytes();
72        let value = serde_json::from_slice(&bytes).unwrap_or_else(|_| {
73            serde_json::Value::String(String::from_utf8_lossy(&bytes).to_string())
74        });
75        out.push(AdminEntry {
76            key: sample.key_expr().as_str().to_string(),
77            value,
78        });
79    }
80    opts.note_elided(elided);
81    out.sort_by(|a, b| a.key.cmp(&b.key));
82    Ok(out)
83}
84
85/// Enumerate routers/peers from `@/*/router` (and the fields every layout
86/// carries).
87pub async fn routers(session: &Session, timeout: Duration) -> Result<Vec<RouterInfo>> {
88    let entries = admin_get(session, "@/*/router", timeout).await?;
89    Ok(entries
90        .into_iter()
91        .map(|e| {
92            let zid = e
93                .value
94                .get("zid")
95                .and_then(|v| v.as_str())
96                .map(str::to_string)
97                .unwrap_or_else(|| {
98                    // Fall back to the key's zid chunk: @/<zid>/router.
99                    e.key.split('/').nth(1).unwrap_or("?").to_string()
100                });
101            let version = e
102                .value
103                .get("version")
104                .and_then(|v| v.as_str())
105                .map(str::to_string);
106            let locators = e
107                .value
108                .get("locators")
109                .and_then(|v| v.as_array())
110                .map(|a| {
111                    a.iter()
112                        .filter_map(|l| l.as_str().map(str::to_string))
113                        .collect()
114                })
115                .unwrap_or_default();
116            RouterInfo {
117                zid,
118                version,
119                locators,
120                raw: e.value,
121            }
122        })
123        .collect())
124}
125
126/// Extract a storage from one admin entry, tolerantly: the key shape is
127/// `@/<zid>/router/…/storage_manager/storages/<name>[…]`, the value a config
128/// document whose `key_expr` field names what it captures. Pure — the
129/// version-variance lives here, unit-tested.
130pub fn storage_from_admin_entry(key: &str, value: &serde_json::Value) -> Option<StorageInfo> {
131    let chunks: Vec<&str> = key.split('/').collect();
132    let storages_pos = chunks.iter().position(|c| *c == "storages")?;
133    // Only storage_manager subtrees qualify (volumes etc. share the plugin).
134    if chunks.get(storages_pos.checked_sub(1)?) != Some(&"storage_manager") {
135        return None;
136    }
137    let name = chunks.get(storages_pos + 1)?;
138    let zid = chunks.get(1).unwrap_or(&"?");
139    let text = |field: &str| {
140        value
141            .get(field)
142            .and_then(|v| v.as_str())
143            .map(str::to_string)
144    };
145    // The volume is a bare string in some layouts and an object with an `id`
146    // in others. Absorbing both here is what this pure parser is for.
147    let volume = text("volume").or_else(|| {
148        value
149            .get("volume")
150            .and_then(|v| v.get("id"))
151            .and_then(|v| v.as_str())
152            .map(str::to_string)
153    });
154    Some(StorageInfo {
155        zid: (*zid).to_string(),
156        name: (*name).to_string(),
157        key_expr: text("key_expr"),
158        strip_prefix: text("strip_prefix"),
159        volume,
160        raw: value.clone(),
161    })
162}
163
164/// One row per `(zid, name)`, merging field by field.
165///
166/// The config and status subtrees both answer a storage sweep, and neither is
167/// reliably the richer one — so a row that names a `key_expr` and a row that
168/// names a `volume` must combine rather than one winning outright. Pure, and
169/// separated from [`storages`] because with three optional fields a hand-rolled
170/// `dedup_by` is where a quietly-dropped field would hide.
171pub fn merge_storage_rows(mut rows: Vec<StorageInfo>) -> Vec<StorageInfo> {
172    rows.sort_by(|a, b| (&a.zid, &a.name).cmp(&(&b.zid, &b.name)));
173    let mut out: Vec<StorageInfo> = Vec::with_capacity(rows.len());
174    for row in rows {
175        match out.last_mut() {
176            Some(prev) if prev.zid == row.zid && prev.name == row.name => {
177                prev.key_expr = prev.key_expr.take().or(row.key_expr);
178                prev.strip_prefix = prev.strip_prefix.take().or(row.strip_prefix);
179                prev.volume = prev.volume.take().or(row.volume);
180                // Keep the document that said more, so the raw disclosure is
181                // the useful one.
182                if prev.raw.as_object().map(|o| o.len()).unwrap_or(0)
183                    < row.raw.as_object().map(|o| o.len()).unwrap_or(0)
184                {
185                    prev.raw = row.raw;
186                }
187            }
188            _ => out.push(row),
189        }
190    }
191    out
192}
193
194/// Enumerate configured storages across the mesh (issue #14). Zero routers
195/// (peer mesh, admin disabled) is an empty vec, never an error.
196pub async fn storages(session: &Session, timeout: Duration) -> Result<Vec<StorageInfo>> {
197    let entries = admin_get(
198        session,
199        "@/*/router/**/storage_manager/storages/**",
200        timeout,
201    )
202    .await?;
203    let rows: Vec<StorageInfo> = entries
204        .iter()
205        .filter_map(|e| storage_from_admin_entry(&e.key, &e.value))
206        .collect();
207    Ok(merge_storage_rows(rows))
208}
209
210/// Judge every declared **state** family against the configured storages
211/// (issue #14): the family's wire selector vs each storage's key expression,
212/// by key algebra (`includes` ⇒ covered, `intersects` ⇒ partial). Pure.
213pub fn state_coverage(
214    slices: &crate::model::registry::SliceSet,
215    base: &str,
216    storages: &[StorageInfo],
217) -> Vec<CoverageRow> {
218    use zenoh::key_expr::keyexpr;
219
220    let storage_kes: Vec<(&StorageInfo, &keyexpr)> = storages
221        .iter()
222        .filter_map(|s| {
223            let ke = s.key_expr.as_deref()?;
224            keyexpr::new(ke).ok().map(|ke| (s, ke))
225        })
226        .collect();
227    let mut rows = Vec::new();
228    for slice in slices.slices() {
229        for subject in &slice.subjects {
230            if !subject.class.is(&zenkey::Class::State) {
231                continue;
232            }
233            let Ok(pattern) = zenkey::pattern::SubjectPattern::parse(&subject.path) else {
234                continue;
235            };
236            // Composed via `with_base` so the empty base stays a valid
237            // keyexpr (`format!("{base}/…")` would grow a leading slash and
238            // silently drop every family below).
239            let selector = match &slice.service_origin {
240                Some(origin) => zenkey::grammar::with_base(
241                    base,
242                    format!("v1/{origin}/state/{}", pattern.selector_tail()),
243                ),
244                None => zenkey::grammar::with_base(
245                    base,
246                    format!("v1/*/state/{}/{}", slice.name, pattern.selector_tail()),
247                ),
248            };
249            let Ok(family) = keyexpr::new(selector.as_str()) else {
250                continue;
251            };
252            let mut coverage = Coverage::Uncovered;
253            for (info, ke) in &storage_kes {
254                if ke.includes(family) {
255                    coverage = Coverage::Covered(format!("{}@{}", info.name, info.zid));
256                    break;
257                }
258                if ke.intersects(family) && coverage == Coverage::Uncovered {
259                    coverage = Coverage::Partial(format!("{}@{}", info.name, info.zid));
260                }
261            }
262            rows.push(CoverageRow {
263                producer: slice.name.clone(),
264                path: subject.path.clone(),
265                ttl_s: subject.ttl_s,
266                coverage,
267            });
268        }
269    }
270    rows
271}
272
273impl EntityKind {
274    fn from_chunk(chunk: &str) -> Option<EntityKind> {
275        Some(match chunk {
276            "subscriber" => EntityKind::Subscriber,
277            "publisher" => EntityKind::Publisher,
278            "queryable" => EntityKind::Queryable,
279            "querier" => EntityKind::Querier,
280            "token" => EntityKind::Token,
281            _ => return None,
282        })
283    }
284
285    fn chunk(self) -> &'static str {
286        match self {
287            EntityKind::Subscriber => "subscriber",
288            EntityKind::Publisher => "publisher",
289            EntityKind::Queryable => "queryable",
290            EntityKind::Querier => "querier",
291            EntityKind::Token => "token",
292        }
293    }
294}
295
296/// Parse one admin entry (`@/<zid>/<whatami>/<kind>/<keyexpr...>`) into a
297/// declared entity. Pure; unknown shapes yield `None` (the admin space is
298/// version-dependent surface — tolerate, never fail).
299pub fn declared_from_admin_entry(key: &str, value: &serde_json::Value) -> Option<DeclaredEntity> {
300    let mut chunks = key.split('/');
301    if chunks.next()? != "@" {
302        return None;
303    }
304    let zid = chunks.next()?;
305    let _whatami = chunks.next()?;
306    let kind = EntityKind::from_chunk(chunks.next()?)?;
307    let keyexpr: Vec<&str> = chunks.collect();
308    if keyexpr.is_empty() {
309        return None;
310    }
311    Some(DeclaredEntity {
312        kind,
313        keyexpr: keyexpr.join("/"),
314        node_zid: zid.to_string(),
315        sources: value.clone(),
316    })
317}
318
319/// Enumerate declared subscribers/publishers/queryables/tokens from every
320/// reachable admin space.
321///
322/// `Ok(None)` when **nothing answered at all**: zenoh's `adminspace.enabled`
323/// defaults to *false* (routers ship with it on; a pure peer mesh has none),
324/// so an empty sweep means "not available", never "nothing declared" —
325/// callers MUST render the difference (RFC 09 §5.1 O4). A *publisher* is
326/// visible only if it was declared (P7's rule for the data planes); a bare
327/// `session.put()` never appears here.
328pub async fn declared_entities(
329    session: &Session,
330    timeout: Duration,
331) -> Result<Option<DeclaredEntities>> {
332    let mut entities = Vec::new();
333
334    let mut any_reply = false;
335
336    for kind in [
337        EntityKind::Subscriber,
338        EntityKind::Publisher,
339        EntityKind::Queryable,
340        EntityKind::Querier,
341        EntityKind::Token,
342    ] {
343        let selector = format!("@/*/*/{}/**", kind.chunk());
344        let entries = admin_get(session, &selector, timeout).await?;
345        any_reply |= !entries.is_empty();
346        entities.extend(
347            entries
348                .iter()
349                .filter_map(|e| declared_from_admin_entry(&e.key, &e.value)),
350        );
351    }
352    if !any_reply {
353        return Ok(None);
354    }
355    Ok(Some(DeclaredEntities { entities }))
356}
357
358/// Collapse the per-reporter edges into undirected links.
359///
360/// [`TopologyEdge`]'s own doc has anticipated this since it was written — "a
361/// renderer that wants an undirected mesh dedups by unordered zid pair" — and
362/// two renderers now want it, which is why it stopped being a private helper
363/// in the CLI (issue #207).
364pub fn mesh_links(report: &TopologyReport) -> Vec<MeshLink> {
365    let mut out: Vec<MeshLink> = Vec::new();
366    for e in &report.edges {
367        let (a, b) = if e.reporter <= e.peer {
368            (e.reporter.clone(), e.peer.clone())
369        } else {
370            (e.peer.clone(), e.reporter.clone())
371        };
372        match out.iter_mut().find(|l| l.a == a && l.b == b) {
373            Some(link) => link.corroborated = true,
374            None => out.push(MeshLink {
375                a,
376                b,
377                corroborated: false,
378                links: e.links.clone(),
379            }),
380        }
381    }
382    out
383}
384
385/// Graphviz, self-contained: routers as boxes, peers/clients as ellipses,
386/// heard-of nodes dashed, our own session bold.
387pub fn render_dot(report: &TopologyReport, attachments: &[OriginAttachment]) -> String {
388    use std::fmt::Write as _;
389    let mut out = String::from("graph zenoh_mesh {\n");
390    for n in &report.nodes {
391        let shape = if n.whatami == "router" {
392            "box"
393        } else {
394            "ellipse"
395        };
396        let mut style = Vec::new();
397        if !n.answered {
398            style.push("dashed");
399        }
400        if n.zid == report.self_zid {
401            style.push("bold");
402        }
403        let label = match (&n.answered, &n.version) {
404            (false, _) => format!("{}\\n{} (heard of)", n.zid, n.whatami),
405            (true, Some(v)) => format!("{}\\n{} {v}", n.zid, n.whatami),
406            (true, None) => format!("{}\\n{}", n.zid, n.whatami),
407        };
408        let _ = writeln!(
409            out,
410            "  \"{}\" [shape={shape}, label=\"{label}\"{}];",
411            n.zid,
412            if style.is_empty() {
413                String::new()
414            } else {
415                format!(", style=\"{}\"", style.join(","))
416            }
417        );
418    }
419    for link in mesh_links(report) {
420        let proto = link
421            .links
422            .first()
423            .and_then(|l| l.split('/').next())
424            .unwrap_or("");
425        let _ = writeln!(
426            out,
427            "  \"{}\" -- \"{}\"{};",
428            link.a,
429            link.b,
430            if proto.is_empty() {
431                String::new()
432            } else {
433                format!(" [label=\"{proto}\"]")
434            }
435        );
436    }
437    // Origins as their own small nodes (#131): attached by a solid edge to
438    // the session the admin sources named, or by a dotted one to the mere
439    // reporter — the picture keeps the evidence distinction the join made.
440    for (i, a) in attachments.iter().enumerate() {
441        let id = format!("origin_{i}");
442        let _ = writeln!(
443            out,
444            "  \"{id}\" [shape=hexagon, label=\"{}\", fontsize=10];",
445            a.origin
446        );
447        match &a.session_zid {
448            Some(z) => {
449                let _ = writeln!(out, "  \"{id}\" -- \"{z}\";");
450            }
451            None => {
452                let _ = writeln!(
453                    out,
454                    "  \"{id}\" -- \"{}\" [style=dotted, label=\"reported\"];",
455                    a.reporter_zid
456                );
457            }
458        }
459    }
460    out.push('}');
461    out
462}
463
464/// Collect every zid string under the zenoh 1.9 `Sources` shape
465/// (`{ routers: [...], peers: [...], clients: [...] }`) — tolerant of the
466/// layout varying by version: unknown shapes yield nothing, never an error.
467fn source_zids(sources: &serde_json::Value) -> Vec<String> {
468    let mut out = Vec::new();
469    for kind in ["routers", "peers", "clients"] {
470        if let Some(list) = sources.get(kind).and_then(|v| v.as_array()) {
471            out.extend(list.iter().filter_map(|z| z.as_str().map(str::to_string)));
472        }
473    }
474    out.sort();
475    out.dedup();
476    out
477}
478
479/// Join the admin space's declared liveliness tokens against the keyspace:
480/// which origin hangs off which session (#131).
481///
482/// One `@/*/*/token/**` sweep; each token whose keyexpr parses under `base`
483/// as an `alive` leaf yields an attachment. The session zid is taken from
484/// the token's `sources` **only when they name exactly one** — several
485/// candidates or none degrade to reporter-only, stated rather than guessed.
486/// An empty result means the admin space served no tokens (or none parse
487/// under this base) — an observation, not an empty fleet (O4).
488pub async fn origin_attachments(
489    fleet: &crate::Fleet<'_>,
490    timeout: Duration,
491) -> Result<Vec<OriginAttachment>> {
492    let base = fleet.base();
493
494    let entries = admin_get(fleet.session(), "@/*/*/token/**", timeout).await?;
495
496    let mut out: Vec<OriginAttachment> = Vec::new();
497
498    for e in &entries {
499        let Some(decl) = declared_from_admin_entry(&e.key, &e.value) else {
500            continue;
501        };
502        if decl.kind != EntityKind::Token {
503            continue;
504        }
505        let Some(parsed) = zenkey::grammar::parse_full(base, &decl.keyexpr) else {
506            continue;
507        };
508        // The framework liveliness shape: an `alive` leaf on the state
509        // class (RFC 04 §5). Anything else declared as a token is not an
510        // origin claim and is left alone.
511        if parsed.subject.last().copied() != Some("alive") {
512            continue;
513        }
514        let origin = parsed.origin.chunk().to_string();
515        let zids = source_zids(&decl.sources);
516        let session_zid = match zids.as_slice() {
517            [only] => Some(only.clone()),
518            _ => None,
519        };
520        let attachment = OriginAttachment {
521            origin,
522            session_zid,
523            reporter_zid: decl.node_zid.clone(),
524            token_key: decl.keyexpr.clone(),
525        };
526        // One origin can hold several sessions (one per producer process);
527        // dedup only exact repeats.
528        if !out.iter().any(|a| {
529            a.origin == attachment.origin
530                && a.session_zid == attachment.session_zid
531                && a.reporter_zid == attachment.reporter_zid
532        }) {
533            out.push(attachment);
534        }
535    }
536    Ok(out)
537}
538
539/// Whether a node's admin root doc filters loopback endpoints out of its
540/// `locators` — true from zenoh 1.10.0 (eclipse-zenoh/zenoh#2671, the
541/// loopback scouting fix: the root doc switched to
542/// `get_locators_noloopback()`). Judged from the leading `major.minor`
543/// of the version string the doc itself declares; a version that does
544/// not parse answers `false` — "cannot say", never a claim (O4).
545///
546/// One definition, used by both renderers, so the two tools explain an
547/// empty locator column with one voice (#155).
548pub fn admin_doc_omits_loopback(version: &str) -> bool {
549    let nums: Vec<u64> = version
550        .trim_start_matches(|c: char| !c.is_ascii_digit())
551        .split(|c: char| !c.is_ascii_digit())
552        .take(2)
553        .map_while(|p| p.parse().ok())
554        .collect();
555    matches!(nums.as_slice(), [maj, min] if (*maj, *min) >= (1, 10))
556}
557
558/// Join the admin root docs (`@/<zid>/<whatami>`) into a topology: every
559/// answering node with its locators and version, every session it reports
560/// as an edge, and every zid that is *only* mentioned as a
561/// heard-of-not-queryable node.
562///
563/// Where a root doc declares no locators — since zenoh 1.10.0 that is the
564/// normal answer for a loopback-only node (eclipse-zenoh/zenoh#2671
565/// filters loopback endpoints from the root doc) — the join corroborates
566/// from session links instead: the node-side endpoint of each reported
567/// link lands in [`TopologyNode::locators_via_links`], kept apart from
568/// `locators` because it is link evidence, not a listen-endpoint claim.
569/// Nothing is invented: a node no link names stays honestly empty.
570pub async fn topology(session: &Session, timeout: Duration) -> Result<TopologyReport> {
571    const ASKED: &str = "@/*/*";
572    let entries = admin_get(session, ASKED, timeout).await?;
573    let mut nodes: Vec<TopologyNode> = Vec::new();
574    let mut edges: Vec<TopologyEdge> = Vec::new();
575    for e in &entries {
576        // Root docs only: @/<zid>/<whatami>. Anything deeper is a
577        // different handler and not a node document.
578        let mut chunks = e.key.split('/');
579        let (Some("@"), Some(zid), Some(whatami), None) =
580            (chunks.next(), chunks.next(), chunks.next(), chunks.next())
581        else {
582            continue;
583        };
584        let doc = &e.value;
585        nodes.push(TopologyNode {
586            zid: doc
587                .get("zid")
588                .and_then(|v| v.as_str())
589                .unwrap_or(zid)
590                .to_string(),
591            whatami: whatami.to_string(),
592            version: doc
593                .get("version")
594                .and_then(|v| v.as_str())
595                .map(str::to_string),
596            locators: doc
597                .get("locators")
598                .and_then(|v| v.as_array())
599                .map(|a| {
600                    a.iter()
601                        .filter_map(|l| l.as_str().map(str::to_string))
602                        .collect()
603                })
604                .unwrap_or_default(),
605            locators_via_links: Vec::new(),
606            answered: true,
607        });
608        for s in doc
609            .get("sessions")
610            .and_then(|v| v.as_array())
611            .map(|a| a.as_slice())
612            .unwrap_or_default()
613        {
614            let Some(peer) = s.get("peer").and_then(|v| v.as_str()) else {
615                continue;
616            };
617            edges.push(TopologyEdge {
618                reporter: zid.to_string(),
619                peer: peer.to_string(),
620                whatami: s
621                    .get("whatami")
622                    .and_then(|v| v.as_str())
623                    .unwrap_or("unknown")
624                    .to_string(),
625                region: s.get("region").and_then(|v| v.as_str()).map(str::to_string),
626                links: s
627                    .get("links")
628                    .and_then(|v| v.as_array())
629                    .map(|a| {
630                        a.iter()
631                            .filter_map(|l| {
632                                Some(format!(
633                                    "{} -> {}",
634                                    l.get("src")?.as_str()?,
635                                    l.get("dst")?.as_str()?
636                                ))
637                            })
638                            .collect()
639                    })
640                    .unwrap_or_default(),
641            });
642        }
643    }
644    let answered = nodes.len();
645    // Heard-of nodes: mentioned as a session peer, but no root doc answered
646    // for them (admin space off, or out of reach). Shown, never omitted.
647    for e in &edges {
648        if !nodes.iter().any(|n| n.zid == e.peer) {
649            nodes.push(TopologyNode {
650                zid: e.peer.clone(),
651                whatami: e.whatami.clone(),
652                version: None,
653                locators: Vec::new(),
654                locators_via_links: Vec::new(),
655                answered: false,
656            });
657        }
658    }
659    nodes.sort_by(|a, b| a.zid.cmp(&b.zid));
660    nodes.dedup_by(|a, b| a.zid == b.zid);
661    // Corroborate where the root doc declared nothing (see the fn doc):
662    // for each link `src -> dst`, `src` is an address on the reporter's
663    // side and `dst` one on the peer's. That is what a link *used*, no
664    // more — kept out of `locators` and labelled by the renderers.
665    for n in nodes.iter_mut().filter(|n| n.locators.is_empty()) {
666        for e in &edges {
667            let reporter_side = if e.reporter == n.zid {
668                true
669            } else if e.peer == n.zid {
670                false
671            } else {
672                continue;
673            };
674            for l in &e.links {
675                let mut parts = l.splitn(2, " -> ");
676                let (Some(src), Some(dst)) = (parts.next(), parts.next()) else {
677                    continue;
678                };
679                let end = if reporter_side { src } else { dst };
680                if !n.locators_via_links.iter().any(|x| x == end) {
681                    n.locators_via_links.push(end.to_string());
682                }
683            }
684        }
685    }
686    Ok(TopologyReport {
687        nodes,
688        edges,
689        asked: ASKED.to_string(),
690        answered,
691        self_zid: session.zid().to_string(),
692    })
693}
694
695#[cfg(test)]
696mod tests {
697    use super::*;
698
699    #[test]
700    fn storage_extraction_tolerates_layouts() {
701        // 1.x config-subtree shape.
702        let v = serde_json::json!({"key_expr": "zs/v1/*/state/**", "volume": "fs"});
703        let s = storage_from_admin_entry(
704            "@/abc123/router/config/plugins/storage_manager/storages/latest",
705            &v,
706        )
707        .unwrap();
708        assert_eq!(s.zid, "abc123");
709        assert_eq!(s.name, "latest");
710        assert_eq!(s.key_expr.as_deref(), Some("zs/v1/*/state/**"));
711        assert_eq!(s.volume.as_deref(), Some("fs"), "a bare-string volume");
712        // The other spelling of the same field.
713        let v = serde_json::json!({
714            "key_expr": "zs/v1/*/state/**",
715            "strip_prefix": "zs/v1",
716            "volume": {"id": "rocksdb", "dir": "latest"},
717        });
718        let s = storage_from_admin_entry(
719            "@/abc123/router/config/plugins/storage_manager/storages/durable",
720            &v,
721        )
722        .unwrap();
723        assert_eq!(s.volume.as_deref(), Some("rocksdb"), "an object volume");
724        assert_eq!(s.strip_prefix.as_deref(), Some("zs/v1"));
725        // Status-subtree shape without key_expr still names the storage.
726        let s = storage_from_admin_entry(
727            "@/abc123/router/status/plugins/storage_manager/storages/latest/info",
728            &serde_json::json!("ok"),
729        )
730        .unwrap();
731        assert_eq!(s.name, "latest");
732        assert!(s.key_expr.is_none());
733        // Non-storage subtrees do not match.
734        assert!(
735            storage_from_admin_entry(
736                "@/abc123/router/config/plugins/storage_manager/volumes/fs",
737                &serde_json::json!({}),
738            )
739            .is_none()
740        );
741    }
742
743    /// Config and status subtrees both answer, neither is reliably richer, and
744    /// a field named by only one of them must survive either arrival order.
745    #[test]
746    fn merging_a_storage_keeps_every_field_either_side_named() {
747        let config = StorageInfo {
748            zid: "z1".into(),
749            name: "latest".into(),
750            key_expr: Some("zs/**".into()),
751            strip_prefix: Some("zs".into()),
752            volume: None,
753            raw: serde_json::json!({"key_expr": "zs/**", "strip_prefix": "zs"}),
754        };
755        let status = StorageInfo {
756            zid: "z1".into(),
757            name: "latest".into(),
758            key_expr: None,
759            strip_prefix: None,
760            volume: Some("fs".into()),
761            raw: serde_json::json!({"volume": "fs"}),
762        };
763        for rows in [
764            vec![config.clone(), status.clone()],
765            vec![status, config.clone()],
766        ] {
767            let merged = merge_storage_rows(rows);
768            assert_eq!(merged.len(), 1, "one row per (zid, name)");
769            assert_eq!(merged[0].key_expr.as_deref(), Some("zs/**"));
770            assert_eq!(merged[0].strip_prefix.as_deref(), Some("zs"));
771            assert_eq!(merged[0].volume.as_deref(), Some("fs"));
772        }
773
774        // Two names under one zid stay two rows.
775        let other = StorageInfo {
776            zid: "z1".into(),
777            name: "history".into(),
778            key_expr: None,
779            strip_prefix: None,
780            volume: None,
781            raw: serde_json::Value::Null,
782        };
783        assert_eq!(merge_storage_rows(vec![config, other]).len(), 2);
784    }
785
786    fn slices_with_state() -> crate::model::registry::SliceSet {
787        let toml = r#"
788            [registry]
789            version = "1.0"
790            app = "t"
791            convention = 1
792            [producer]
793            name = "tc"
794            [[subject]]
795            path = "health"
796            class = "state"
797            type = "Health"
798            ttl_s = 60
799            [[subject]]
800            path = "config/{iface}"
801            class = "state"
802            type = "Config"
803            ttl_s = 120
804            [[subject]]
805            path = "bandwidth"
806            class = "telemetry"
807            type = "Point"
808        "#;
809        crate::model::registry::SliceSet::from_toml_for_tests(toml)
810    }
811
812    fn storage(name: &str, key_expr: &str) -> StorageInfo {
813        StorageInfo {
814            zid: "z1".into(),
815            name: name.into(),
816            key_expr: Some(key_expr.into()),
817            strip_prefix: None,
818            volume: None,
819            raw: serde_json::Value::Null,
820        }
821    }
822
823    #[test]
824    fn coverage_judges_covered_partial_uncovered() {
825        let slices = slices_with_state();
826        // Full state storage: everything covered; telemetry not judged.
827        let rows = state_coverage(&slices, "zs", &[storage("latest", "zs/v1/*/state/**")]);
828        assert_eq!(rows.len(), 2);
829        assert!(
830            rows.iter()
831                .all(|r| matches!(r.coverage, Coverage::Covered(_)))
832        );
833
834        // A one-interface storage: config/{iface} is partial, health uncovered.
835        let rows = state_coverage(
836            &slices,
837            "zs",
838            &[storage("one", "zs/v1/*/state/tc/config/eth0")],
839        );
840        let health = rows.iter().find(|r| r.path == "health").unwrap();
841        assert_eq!(health.coverage, Coverage::Uncovered);
842        let config = rows.iter().find(|r| r.path == "config/{iface}").unwrap();
843        assert!(matches!(config.coverage, Coverage::Partial(_)));
844
845        // No storages at all.
846        let rows = state_coverage(&slices, "zs", &[]);
847        assert!(rows.iter().all(|r| r.coverage == Coverage::Uncovered));
848
849        // The empty base composes a valid selector (`v1/…`, no leading
850        // slash) instead of silently dropping every family.
851        let rows = state_coverage(&slices, "", &[storage("latest", "v1/*/state/**")]);
852        assert_eq!(rows.len(), 2);
853        assert!(
854            rows.iter()
855                .all(|r| matches!(r.coverage, Coverage::Covered(_)))
856        );
857    }
858
859    /// The zenoh 1.9 admin key shape (`@/<zid>/<whatami>/<kind>/<keyexpr...>`)
860    /// parses into a declared entity; foreign shapes are tolerated as None.
861    #[test]
862    fn declared_entities_parse_the_admin_key_shape() {
863        let v = serde_json::json!({"routers": [], "peers": ["p1"], "clients": []});
864        let e = declared_from_admin_entry("@/a1b2c3/router/subscriber/zensight/v1/*/state/**", &v)
865            .unwrap();
866        assert_eq!(e.kind, EntityKind::Subscriber);
867        assert_eq!(e.keyexpr, "zensight/v1/*/state/**");
868        assert_eq!(e.node_zid, "a1b2c3");
869
870        let e = declared_from_admin_entry("@/z/peer/publisher/v1/h-a/telemetry/x/m", &v).unwrap();
871        assert_eq!(e.kind, EntityKind::Publisher);
872        assert_eq!(e.keyexpr, "v1/h-a/telemetry/x/m");
873
874        // Tolerated, never fatal:
875        assert!(declared_from_admin_entry("@/z/router/config/x", &v).is_none());
876        assert!(declared_from_admin_entry("@/z/router/subscriber", &v).is_none());
877        assert!(declared_from_admin_entry("not/admin/at/all", &v).is_none());
878    }
879
880    fn report() -> TopologyReport {
881        TopologyReport {
882            nodes: vec![
883                TopologyNode {
884                    zid: "aaa".into(),
885                    whatami: "router".into(),
886                    version: Some("1.9.0".into()),
887                    locators: vec!["tcp/10.0.0.1:7447".into()],
888                    locators_via_links: vec![],
889                    answered: true,
890                },
891                TopologyNode {
892                    zid: "bbb".into(),
893                    whatami: "peer".into(),
894                    version: None,
895                    locators: vec![],
896                    locators_via_links: vec![],
897                    answered: false,
898                },
899            ],
900            edges: vec![
901                TopologyEdge {
902                    reporter: "aaa".into(),
903                    peer: "bbb".into(),
904                    whatami: "peer".into(),
905                    region: None,
906                    links: vec!["tcp/10.0.0.1:7447 -> tcp/10.0.0.2:53210".into()],
907                },
908                TopologyEdge {
909                    reporter: "bbb".into(),
910                    peer: "aaa".into(),
911                    whatami: "router".into(),
912                    region: None,
913                    links: vec![],
914                },
915            ],
916            asked: "@/*/*".into(),
917            answered: 1,
918            self_zid: "bbb".into(),
919        }
920    }
921
922    /// The version gate for the 1.10 loopback filter: judged from the
923    /// doc's own version string, and an unparseable one answers "cannot
924    /// say" — false — never a claim either way.
925    #[test]
926    fn the_loopback_filter_is_judged_from_the_docs_own_version() {
927        assert!(admin_doc_omits_loopback("1.10.0"));
928        assert!(admin_doc_omits_loopback(
929            "v1.10.0-12-gabcdef built with rustc"
930        ));
931        assert!(admin_doc_omits_loopback("1.11.2"));
932        assert!(admin_doc_omits_loopback("2.0.0"));
933        assert!(!admin_doc_omits_loopback("1.9.0"));
934        assert!(!admin_doc_omits_loopback("0.11.0-dev"));
935        assert!(!admin_doc_omits_loopback("unknown"));
936        assert!(!admin_doc_omits_loopback(""));
937    }
938
939    /// Reciprocal reports collapse to one undirected edge, marked as
940    /// corroborated — not drawn twice, not silently merged.
941    #[test]
942    fn reciprocal_reports_corroborate_one_edge() {
943        let edges = mesh_links(&report());
944        assert_eq!(edges.len(), 1);
945        let link = &edges[0];
946        assert_eq!((link.a.as_str(), link.b.as_str()), ("aaa", "bbb"));
947        assert!(link.corroborated, "both ends reported it");
948    }
949
950    /// The DOT form: routers boxed, heard-of nodes dashed, our session
951    /// bold, edges labeled by protocol — pipeable to `dot -Tsvg` as-is.
952    #[test]
953    fn the_dot_form_marks_what_the_join_knows() {
954        let dot = render_dot(&report(), &[]);
955        assert!(dot.starts_with("graph zenoh_mesh {"), "{dot}");
956        assert!(dot.contains("\"aaa\" [shape=box"), "{dot}");
957        assert!(dot.contains("heard of"), "{dot}");
958        assert!(dot.contains("style=\"dashed,bold\""), "{dot}");
959        assert!(dot.contains("\"aaa\" -- \"bbb\" [label=\"tcp\"]"), "{dot}");
960        assert!(dot.ends_with('}'), "{dot}");
961    }
962
963    /// The origin overlay (#131): a sources-named attachment is a solid
964    /// edge; a reporter-only one is dotted and says "reported" — the DOT
965    /// keeps the evidence distinction the join made.
966    #[test]
967    fn the_dot_form_keeps_the_attachment_evidence_distinction() {
968        let attachments = vec![
969            OriginAttachment {
970                origin: "h-cccccccccccc".into(),
971                session_zid: Some("bbb".into()),
972                reporter_zid: "aaa".into(),
973                token_key: "v1/h-cccccccccccc/state/demo/alive".into(),
974            },
975            OriginAttachment {
976                origin: "h-dddddddddddd".into(),
977                session_zid: None,
978                reporter_zid: "aaa".into(),
979                token_key: "v1/h-dddddddddddd/state/demo/alive".into(),
980            },
981        ];
982        let dot = render_dot(&report(), &attachments);
983        assert!(dot.contains("label=\"h-cccccccccccc\""), "{dot}");
984        assert!(dot.contains("\"origin_0\" -- \"bbb\";"), "{dot}");
985        assert!(
986            dot.contains("\"origin_1\" -- \"aaa\" [style=dotted, label=\"reported\"]"),
987            "{dot}"
988        );
989    }
990}