Skip to main content

agentbridge/
member_source.rs

1// Copyright AGNTCY Contributors (https://github.com/agntcy)
2// SPDX-License-Identifier: Apache-2.0
3
4//! Pluggable resolution of SLIM group member candidates.
5//!
6//! A group moderator can build up who's trusted/invitable in more than one
7//! way: by searching the Agent Directory for a skill, by looking up an
8//! already-known DID, or by naming exact `{name, did, endpoint}` triples by
9//! hand. [`MemberSource`] is the common abstraction — each technique is a
10//! small, independent implementation, and a group's candidate pool can come
11//! from any combination of them.
12
13use std::process::Command;
14
15use serde_json::Value;
16use shadi_a2a::{A2ABinding, A2ALocator};
17
18use crate::dir_registry::dirctl_binary;
19
20/// A directory-resolved (or manually-named) agent, ready to be admitted into
21/// a group's DID trust set and/or invited into a live session.
22///
23/// `did` is the portable name. `slim_endpoint` / `a2a_url` are locators that
24/// can change when the agent moves; look the DID up again to refresh them.
25/// Local aliases for the DID: the tool name and the SLIM channel
26/// (`agntcy/shadi/<name>-a2a`).
27#[derive(Debug, Clone, PartialEq, Eq)]
28pub struct CandidateMember {
29    pub name: String,
30    pub did: String,
31    pub slim_endpoint: Option<String>,
32    pub a2a_url: Option<String>,
33    /// Official unicast binding for [`Self::a2a_url`] (`GRPC`, `JSONRPC`,
34    /// `HTTP+JSON`). `None` when there is no HTTP locator.
35    pub a2a_binding: Option<A2ABinding>,
36}
37
38impl CandidateMember {
39    /// SLIM locator: the local node this agent attached to.
40    pub fn slim_locator(&self) -> Option<A2ALocator> {
41        let node = self.slim_endpoint.as_deref()?.trim();
42        if node.is_empty() {
43            return None;
44        }
45        Some(A2ALocator::slim(node))
46    }
47
48    /// Current official A2A unicast locator, if this candidate advertises one.
49    pub fn unicast_locator(&self) -> Option<A2ALocator> {
50        let url = self.a2a_url.as_deref()?.trim();
51        if url.is_empty() {
52            return None;
53        }
54        let binding = self.a2a_binding.unwrap_or(A2ABinding::Grpc);
55        if !binding.is_unicast() {
56            return None;
57        }
58        Some(A2ALocator::new(binding, url))
59    }
60
61    /// Every current locator (SLIM node first, then unicast).
62    pub fn locators(&self) -> Vec<A2ALocator> {
63        let mut out = Vec::new();
64        if let Some(locator) = self.slim_locator() {
65            out.push(locator);
66        }
67        if let Some(locator) = self.unicast_locator() {
68            out.push(locator);
69        }
70        out
71    }
72
73    /// Dispatch locator: prefer unicast when both exist (same as NEXT).
74    pub fn preferred_locator(&self) -> Option<A2ALocator> {
75        self.unicast_locator().or_else(|| self.slim_locator())
76    }
77}
78
79/// A technique for resolving a set of candidate group members.
80pub trait MemberSource {
81    fn resolve(&self) -> Result<Vec<CandidateMember>, String>;
82}
83
84/// Server/auth used by every Agent Directory-backed `MemberSource`.
85#[derive(Debug, Clone)]
86pub struct DirLookupOptions {
87    pub server_addr: String,
88    pub gh_token: Option<String>,
89    pub limit: usize,
90}
91
92/// Discover candidates by capability: `dirctl search --skill <skill>`, then
93/// pull each matching record for its real A2A card and DID.
94pub struct SkillSearchSource {
95    pub skill: String,
96    pub dir: DirLookupOptions,
97}
98
99impl MemberSource for SkillSearchSource {
100    fn resolve(&self) -> Result<Vec<CandidateMember>, String> {
101        resolve_via_dirctl_query(&["--skill", &self.skill], &self.dir)
102    }
103}
104
105/// Discover a candidate by an already-known DID: `dirctl search --author
106/// <did>` resolves its current name/skills/SLIM endpoint from the Directory
107/// without the moderator needing to know anything but the DID up front.
108pub struct DidLookupSource {
109    pub did: String,
110    pub dir: DirLookupOptions,
111}
112
113impl MemberSource for DidLookupSource {
114    fn resolve(&self) -> Result<Vec<CandidateMember>, String> {
115        resolve_via_dirctl_query(&["--author", &self.did], &self.dir)
116    }
117}
118
119/// The fully-manual technique: candidates named directly, no Directory
120/// round-trip at all. Formalizes what a moderator could already do by hand
121/// into the same abstraction as the discovery-based sources.
122pub struct ExplicitListSource {
123    pub entries: Vec<CandidateMember>,
124}
125
126impl MemberSource for ExplicitListSource {
127    fn resolve(&self) -> Result<Vec<CandidateMember>, String> {
128        Ok(self.entries.clone())
129    }
130}
131
132// ---------------------------------------------------------------------------
133// `--members <spec>` parsing — the CLI-facing surface shared by every call
134// site that lets a moderator name member sources on the command line
135// (`shadictl slim create-group --members ...`, `/slim invite-from <spec>`).
136// ---------------------------------------------------------------------------
137
138/// Parse one `--members`/`invite-from` spec into the `MemberSource` it names:
139/// `skill:<skill>` → [`SkillSearchSource`], `did:<did>` → [`DidLookupSource`],
140/// `explicit:<name>=<did>[@<endpoint>]` → [`ExplicitListSource`].
141pub fn parse_member_spec(
142    spec: &str,
143    dir: &DirLookupOptions,
144) -> Result<Box<dyn MemberSource>, String> {
145    if let Some(skill) = spec.strip_prefix("skill:") {
146        if skill.is_empty() {
147            return Err(format!(
148                "invalid member spec '{spec}': skill: needs a skill name"
149            ));
150        }
151        return Ok(Box::new(SkillSearchSource {
152            skill: skill.to_string(),
153            dir: dir.clone(),
154        }));
155    }
156
157    if let Some(did) = spec.strip_prefix("did:") {
158        if did.is_empty() {
159            return Err(format!("invalid member spec '{spec}': did: needs a DID"));
160        }
161        return Ok(Box::new(DidLookupSource {
162            did: did.to_string(),
163            dir: dir.clone(),
164        }));
165    }
166
167    if let Some(rest) = spec.strip_prefix("explicit:") {
168        let (name, did_and_endpoint) = rest.split_once('=').ok_or_else(|| {
169            format!("invalid member spec '{spec}': expected explicit:<name>=<did>[@<endpoint>]")
170        })?;
171        if name.is_empty() {
172            return Err(format!(
173                "invalid member spec '{spec}': explicit: needs a name"
174            ));
175        }
176        let (did, endpoint) = match did_and_endpoint.split_once('@') {
177            Some((did, endpoint)) => (did, Some(endpoint.to_string())),
178            None => (did_and_endpoint, None),
179        };
180        if did.is_empty() {
181            return Err(format!(
182                "invalid member spec '{spec}': explicit: needs a DID"
183            ));
184        }
185        return Ok(Box::new(ExplicitListSource {
186            entries: vec![CandidateMember {
187                name: name.to_string(),
188                did: did.to_string(),
189                slim_endpoint: endpoint,
190                a2a_url: None,
191                a2a_binding: None,
192            }],
193        }));
194    }
195
196    Err(format!(
197        "invalid member spec '{spec}': expected skill:<skill>, did:<did>, or explicit:<name>=<did>[@<endpoint>]"
198    ))
199}
200
201/// Resolve every `--members` spec and concatenate their candidates. Specs are
202/// resolved independently, in order — a group's candidate pool can come from
203/// any combination of techniques at once.
204pub fn resolve_members(
205    specs: &[String],
206    dir: &DirLookupOptions,
207) -> Result<Vec<CandidateMember>, String> {
208    let mut resolved = Vec::new();
209    for spec in specs {
210        let source = parse_member_spec(spec, dir)?;
211        resolved.extend(source.resolve()?);
212    }
213    Ok(resolved)
214}
215
216/// SLIM channel advertised for a registered adapter (`agntcy/shadi/<id>-a2a`).
217pub fn slim_channel_name(agent_id: &str) -> String {
218    format!("agntcy/shadi/{agent_id}-a2a")
219}
220
221/// True when `query` is the tool name or its SLIM channel (aliases for the DID).
222pub fn matches_local_alias(name: &str, query: &str) -> bool {
223    let query = query.trim();
224    query == name || query == slim_channel_name(name) || query == format!("{name}-a2a")
225}
226
227/// True when `query` is a DID (`did:key:…`, `did:web:…`, …), including the
228/// member-spec form `did:did:key:…`.
229pub fn parse_peer_did(query: &str) -> Option<String> {
230    let trimmed = query.trim();
231    if let Some(rest) = trimmed.strip_prefix("did:did:") {
232        if rest.is_empty() {
233            return None;
234        }
235        return Some(format!("did:{rest}"));
236    }
237    if trimmed.starts_with("did:") {
238        let method_and_id = &trimmed["did:".len()..];
239        if method_and_id.contains(':') && !method_and_id.starts_with(':') {
240            return Some(trimmed.to_string());
241        }
242    }
243    None
244}
245
246/// Resolve an agent by DID (portable name) or local adapter name.
247///
248/// Order: on-host leases, then Agent Directory by author DID. A URL is never
249/// the name — pass it separately as a locator override after this returns.
250pub fn resolve_adapter_peer(
251    query: &str,
252    registry: &crate::local_registry::LocalAdapterRegistry,
253    dir: Option<&DirLookupOptions>,
254) -> Result<CandidateMember, String> {
255    match registry.resolve_live(query) {
256        Ok(record) => return Ok(record.to_candidate()),
257        Err(err) if err.contains("ambiguous") => return Err(err),
258        Err(_) => {}
259    }
260    let Some(did) = parse_peer_did(query) else {
261        return Err(format!(
262            "no live adapter named '{query}'. Use the DID from `list --local` \
263             (`--to did:key:…`) so a URL change still finds the same agent"
264        ));
265    };
266    let Some(dir) = dir else {
267        return Err(format!(
268            "DID {did} is not listening on this host. Publish it to DIR or \
269             pass --a2a-url as a locator override (the DID remains the name)"
270        ));
271    };
272    let found = DidLookupSource {
273        did: did.clone(),
274        dir: dir.clone(),
275    }
276    .resolve()?;
277    found
278        .into_iter()
279        .find(|c| c.did == did)
280        .ok_or_else(|| {
281            format!(
282                "DID {did} was not found locally or in the Agent Directory. \
283                 Re-register the adapter (new URL, same DID) or pass --a2a-url"
284            )
285        })
286}
287
288/// Resolve a portable name (DID, tool name, or SLIM channel) to locators.
289///
290/// The name stays the DID. Locators are the current SLIM node and/or unicast
291/// `{binding, url}` from the live lease or DIR card.
292pub fn resolve_locator(
293    query: &str,
294    registry: &crate::local_registry::LocalAdapterRegistry,
295    dir: Option<&DirLookupOptions>,
296) -> Result<(CandidateMember, Vec<A2ALocator>), String> {
297    let peer = resolve_adapter_peer(query, registry, dir)?;
298    let locators = peer.locators();
299    Ok((peer, locators))
300}
301
302/// Resolve a portable name to the current unicast locator only.
303pub fn resolve_unicast_locator(
304    query: &str,
305    registry: &crate::local_registry::LocalAdapterRegistry,
306    dir: Option<&DirLookupOptions>,
307) -> Result<(CandidateMember, Option<A2ALocator>), String> {
308    let peer = resolve_adapter_peer(query, registry, dir)?;
309    let locator = peer.unicast_locator();
310    Ok((peer, locator))
311}
312
313// ---------------------------------------------------------------------------
314// dirctl-backed resolution
315// ---------------------------------------------------------------------------
316
317fn apply_dir_auth(cmd: &mut Command, dir: &DirLookupOptions) {
318    if let Some(token) = &dir.gh_token {
319        cmd.env("DIRECTORY_CLIENT_AUTH_MODE", "github")
320            .env("DIRECTORY_CLIENT_GITHUB_TOKEN", token);
321    }
322}
323
324/// Run `dirctl search <query_args> --output jsonl`, then pull and parse each
325/// matching CID for a real `{name, did, slim_endpoint}` candidate.
326fn resolve_via_dirctl_query(
327    query_args: &[&str],
328    dir: &DirLookupOptions,
329) -> Result<Vec<CandidateMember>, String> {
330    let cids = search_cids(query_args, dir)?;
331    let mut members = Vec::with_capacity(cids.len());
332    for cid in cids {
333        match pull_record_json(&cid, dir) {
334            Ok(record) => {
335                if let Some(candidate) = extract_candidate(&record) {
336                    members.push(candidate);
337                } else {
338                    eprintln!(
339                        "[member_source] {cid}: no DID (authors) or integration/a2a card — skipped"
340                    );
341                }
342            }
343            Err(e) => eprintln!("[member_source] failed to pull {cid}: {e}"),
344        }
345    }
346    Ok(members)
347}
348
349fn search_cids(query_args: &[&str], dir: &DirLookupOptions) -> Result<Vec<String>, String> {
350    let mut cmd = Command::new(dirctl_binary());
351    cmd.arg("search");
352    for arg in query_args {
353        cmd.arg(arg);
354    }
355    cmd.arg("--limit")
356        .arg(dir.limit.to_string())
357        .arg("--server-addr")
358        .arg(&dir.server_addr)
359        .arg("--output")
360        .arg("jsonl");
361    apply_dir_auth(&mut cmd, dir);
362
363    let output = match cmd.output() {
364        Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
365            return Err("dirctl not found in PATH".to_string());
366        }
367        Err(e) => return Err(format!("dirctl search failed: {e}")),
368        Ok(o) => o,
369    };
370    if !output.status.success() {
371        return Err(String::from_utf8_lossy(&output.stderr).into_owned());
372    }
373
374    let stdout = String::from_utf8_lossy(&output.stdout);
375    Ok(stdout
376        .lines()
377        .filter(|l| !l.trim().is_empty())
378        .filter_map(|l| serde_json::from_str::<String>(l.trim()).ok())
379        .collect())
380}
381
382fn pull_record_json(cid: &str, dir: &DirLookupOptions) -> Result<Value, String> {
383    let mut cmd = Command::new(dirctl_binary());
384    cmd.arg("pull")
385        .arg(cid)
386        .arg("--server-addr")
387        .arg(&dir.server_addr)
388        .arg("--output")
389        .arg("json");
390    apply_dir_auth(&mut cmd, dir);
391
392    let output = match cmd.output() {
393        Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
394            return Err("dirctl not found in PATH".to_string());
395        }
396        Err(e) => return Err(format!("dirctl pull failed: {e}")),
397        Ok(o) => o,
398    };
399    if !output.status.success() {
400        return Err(String::from_utf8_lossy(&output.stderr).into_owned());
401    }
402
403    serde_json::from_slice(&output.stdout).map_err(|e| format!("parse pulled record: {e}"))
404}
405
406/// Extract `{name, did, slim_endpoint}` from a pulled OASF record: the DID is
407/// the record's own `authors[0]`, the name/endpoint come from the
408/// `integration/a2a` module's `card_data` (see [`crate::dir_registry::wrap_agent_card`]
409/// for the shape this mirrors). Returns `None` if either is missing — a
410/// record without a DID or without a real A2A card isn't a usable candidate.
411fn extract_candidate(record: &Value) -> Option<CandidateMember> {
412    let did = record
413        .get("authors")
414        .and_then(Value::as_array)
415        .and_then(|a| a.first())
416        .and_then(Value::as_str)?
417        .to_string();
418
419    let card = record
420        .get("modules")
421        .and_then(Value::as_array)?
422        .iter()
423        .find(|m| m.get("name").and_then(Value::as_str) == Some("integration/a2a"))
424        .and_then(|m| m.get("data"))
425        .and_then(|d| d.get("card_data"))?;
426
427    let name = card
428        .get("name")
429        .and_then(Value::as_str)
430        .unwrap_or("(unnamed)")
431        .to_string();
432
433    let supported = card
434        .get("supportedInterfaces")
435        .and_then(Value::as_array)
436        .cloned()
437        .unwrap_or_default();
438
439    let slim_endpoint = supported.iter().find_map(|iface| {
440        let url = iface.get("url")?.as_str()?;
441        let rest = url.strip_prefix("slim://")?;
442        Some(rest.split('/').next().unwrap_or(rest).to_string())
443    });
444    let (a2a_binding, a2a_url) = supported
445        .iter()
446        .find_map(unicast_locator_from_interface)
447        .map(|(binding, url)| (Some(binding), Some(url)))
448        .unwrap_or((None, None));
449
450    Some(CandidateMember {
451        name,
452        did,
453        slim_endpoint,
454        a2a_url,
455        a2a_binding,
456    })
457}
458
459fn unicast_locator_from_interface(iface: &Value) -> Option<(A2ABinding, String)> {
460    let url = iface.get("url")?.as_str()?;
461    if url.starts_with("slim://") {
462        return None;
463    }
464    let binding_raw = iface
465        .get("protocolBinding")
466        .and_then(Value::as_str)
467        .unwrap_or("");
468    let binding = if binding_raw.is_empty() {
469        if url.starts_with("http://") || url.starts_with("https://") {
470            A2ABinding::Grpc
471        } else {
472            return None;
473        }
474    } else {
475        A2ABinding::parse(binding_raw).ok()?
476    };
477    if !binding.is_unicast() {
478        return None;
479    }
480    let url = if url.starts_with("http://") || url.starts_with("https://") {
481        url.to_string()
482    } else if url.contains("://") {
483        url.to_string()
484    } else {
485        // a2a-lf strips `http://` on some agent-card URLs.
486        format!("http://{url}")
487    };
488    Some((binding, url))
489}
490
491#[cfg(test)]
492mod tests {
493    use super::*;
494
495    fn a2a_record(name: &str, did: &str, endpoint: &str) -> Value {
496        serde_json::json!({
497            "authors": [did],
498            "modules": [{
499                "name": "integration/a2a",
500                "data": {
501                    "card_data": {
502                        "name": name,
503                        "supportedInterfaces": [{
504                            "url": format!("slim://{endpoint}/agntcy/shadi/{name}-a2a"),
505                            "protocolBinding": "SLIMRPC",
506                            "protocolVersion": "0.3.0",
507                        }],
508                    },
509                    "card_schema_version": "v1.0.0",
510                },
511            }],
512        })
513    }
514
515    #[test]
516    fn extract_candidate_reads_did_and_slim_endpoint_from_module() {
517        let record = a2a_record("copilot", "did:key:z6Mk...", "127.0.0.1:47357");
518        let candidate = extract_candidate(&record).expect("candidate");
519        assert_eq!(candidate.name, "copilot");
520        assert_eq!(candidate.did, "did:key:z6Mk...");
521        assert_eq!(candidate.slim_endpoint.as_deref(), Some("127.0.0.1:47357"));
522    }
523
524    #[test]
525    fn extract_candidate_reads_grpc_url_even_when_http_prefix_was_stripped() {
526        let record = serde_json::json!({
527            "authors": ["did:key:zGrpc"],
528            "modules": [{
529                "name": "integration/a2a",
530                "data": {
531                    "card_data": {
532                        "name": "copilot",
533                        "supportedInterfaces": [{
534                            "url": "127.0.0.1:50051",
535                            "protocolBinding": "GRPC",
536                            "protocolVersion": "1.0",
537                        }],
538                    },
539                    "card_schema_version": "v1.0.0",
540                },
541            }],
542        });
543        let candidate = extract_candidate(&record).expect("candidate");
544        assert_eq!(candidate.did, "did:key:zGrpc");
545        assert_eq!(candidate.a2a_url.as_deref(), Some("http://127.0.0.1:50051"));
546        assert_eq!(candidate.a2a_binding, Some(A2ABinding::Grpc));
547        assert_eq!(candidate.slim_endpoint, None);
548        assert_eq!(
549            candidate.unicast_locator().unwrap().display_uri(),
550            "grpc://127.0.0.1:50051"
551        );
552    }
553
554    #[test]
555    fn extract_candidate_reads_jsonrpc_and_http_json_bindings() {
556        let jsonrpc = serde_json::json!({
557            "authors": ["did:key:zJson"],
558            "modules": [{
559                "name": "integration/a2a",
560                "data": {
561                    "card_data": {
562                        "name": "copilot",
563                        "supportedInterfaces": [{
564                            "url": "http://127.0.0.1:8080",
565                            "protocolBinding": "JSONRPC",
566                            "protocolVersion": "1.0",
567                        }],
568                    },
569                    "card_schema_version": "v1.0.0",
570                },
571            }],
572        });
573        let candidate = extract_candidate(&jsonrpc).expect("jsonrpc");
574        assert_eq!(candidate.a2a_binding, Some(A2ABinding::Jsonrpc));
575        assert_eq!(candidate.a2a_url.as_deref(), Some("http://127.0.0.1:8080"));
576        assert_eq!(
577            candidate.unicast_locator().unwrap().display_uri(),
578            "jsonrpc://127.0.0.1:8080"
579        );
580
581        let rest = serde_json::json!({
582            "authors": ["did:key:zRest"],
583            "modules": [{
584                "name": "integration/a2a",
585                "data": {
586                    "card_data": {
587                        "name": "codex",
588                        "supportedInterfaces": [{
589                            "url": "127.0.0.1:8081",
590                            "protocolBinding": "HTTP+JSON",
591                            "protocolVersion": "1.0",
592                        }],
593                    },
594                    "card_schema_version": "v1.0.0",
595                },
596            }],
597        });
598        let candidate = extract_candidate(&rest).expect("http+json");
599        assert_eq!(candidate.a2a_binding, Some(A2ABinding::HttpJson));
600        assert_eq!(candidate.a2a_url.as_deref(), Some("http://127.0.0.1:8081"));
601    }
602
603    #[test]
604    fn parse_peer_did_accepts_did_key_and_member_spec_form() {
605        assert_eq!(
606            parse_peer_did("did:key:z6Mkabc").as_deref(),
607            Some("did:key:z6Mkabc")
608        );
609        assert_eq!(
610            parse_peer_did("did:did:key:z6Mkabc").as_deref(),
611            Some("did:key:z6Mkabc")
612        );
613        assert_eq!(parse_peer_did("copilot"), None);
614        assert_eq!(parse_peer_did("did:"), None);
615        assert_eq!(parse_peer_did("did:did:"), None);
616        assert_eq!(
617            parse_peer_did("  did:key:z6Mkabc  ").as_deref(),
618            Some("did:key:z6Mkabc")
619        );
620    }
621
622    #[test]
623    fn locators_skip_empty_urls_and_slim_unicast_binding() {
624        let empty = CandidateMember {
625            name: "copilot".to_string(),
626            did: "did:key:zEmpty".to_string(),
627            slim_endpoint: Some("   ".to_string()),
628            a2a_url: Some(String::new()),
629            a2a_binding: None,
630        };
631        assert!(empty.slim_locator().is_none());
632        assert!(empty.unicast_locator().is_none());
633
634        let slim_as_unicast = CandidateMember {
635            name: "copilot".to_string(),
636            did: "did:key:zSlim".to_string(),
637            slim_endpoint: None,
638            a2a_url: Some("http://127.0.0.1:9".to_string()),
639            a2a_binding: Some(A2ABinding::Slim),
640        };
641        assert!(slim_as_unicast.unicast_locator().is_none());
642    }
643
644    #[test]
645    fn unicast_locator_from_interface_infers_http_and_skips_non_unicast() {
646        let http = serde_json::json!({"url": "http://127.0.0.1:9"});
647        assert_eq!(
648            unicast_locator_from_interface(&http),
649            Some((A2ABinding::Grpc, "http://127.0.0.1:9".to_string()))
650        );
651        let https = serde_json::json!({"url": "https://example.test"});
652        assert_eq!(
653            unicast_locator_from_interface(&https),
654            Some((A2ABinding::Grpc, "https://example.test".to_string()))
655        );
656        let bare = serde_json::json!({"url": "127.0.0.1:9"});
657        assert!(unicast_locator_from_interface(&bare).is_none());
658        let slim = serde_json::json!({
659            "url": "http://127.0.0.1:9",
660            "protocolBinding": "SLIMRPC"
661        });
662        assert!(unicast_locator_from_interface(&slim).is_none());
663        let grpc_scheme = serde_json::json!({
664            "url": "grpc://127.0.0.1:9",
665            "protocolBinding": "GRPC"
666        });
667        assert_eq!(
668            unicast_locator_from_interface(&grpc_scheme),
669            Some((A2ABinding::Grpc, "grpc://127.0.0.1:9".to_string()))
670        );
671    }
672
673    #[test]
674    fn resolve_adapter_peer_requires_dir_for_unknown_did() {
675        let (_dir, registry) = temp_registry();
676        let err = resolve_adapter_peer("did:key:zMissing", &registry, None).unwrap_err();
677        assert!(err.contains("not listening on this host"), "{err}");
678        assert!(err.contains("did:key:zMissing"), "{err}");
679    }
680
681    fn temp_registry() -> (tempfile::TempDir, crate::local_registry::LocalAdapterRegistry) {
682        let dir = tempfile::tempdir().unwrap();
683        let registry =
684            crate::local_registry::LocalAdapterRegistry::with_dir(dir.path().to_path_buf());
685        (dir, registry)
686    }
687
688    fn grpc_lease(name: &str, did: &str, url: &str) -> crate::local_registry::LocalAdapterRecord {
689        crate::local_registry::LocalAdapterRecord {
690            name: name.to_string(),
691            did: did.to_string(),
692            slim_endpoint: String::new(),
693            a2a_url: url.to_string(),
694            a2a_binding: A2ABinding::Grpc,
695            pid: std::process::id(),
696        }
697    }
698
699    #[test]
700    fn resolve_adapter_peer_uses_local_lease_by_name_or_did() {
701        let (_dir, registry) = temp_registry();
702        let _lease = registry
703            .publish(&grpc_lease(
704                "copilot",
705                "did:key:zCopilot",
706                "http://127.0.0.1:50151",
707            ))
708            .unwrap();
709        let by_name = resolve_adapter_peer("copilot", &registry, None).expect("name");
710        assert_eq!(by_name.did, "did:key:zCopilot");
711        assert_eq!(
712            by_name.a2a_url.as_deref(),
713            Some("http://127.0.0.1:50151")
714        );
715        let by_did = resolve_adapter_peer("did:key:zCopilot", &registry, None).expect("did");
716        assert_eq!(by_did.name, "copilot");
717        assert_eq!(by_did.a2a_url, by_name.a2a_url);
718    }
719
720    #[test]
721    fn resolve_adapter_peer_follows_did_after_url_change() {
722        let (_dir, registry) = temp_registry();
723        let mut record = grpc_lease("copilot", "did:key:zSame", "http://127.0.0.1:50151");
724        let _first = registry.publish(&record).unwrap();
725        record.a2a_url = "http://127.0.0.1:50152".to_string();
726        let _second = registry.publish(&record).unwrap();
727        let found = resolve_adapter_peer("did:key:zSame", &registry, None).expect("moved");
728        assert_eq!(found.a2a_url.as_deref(), Some("http://127.0.0.1:50152"));
729        assert_eq!(found.a2a_binding, Some(A2ABinding::Grpc));
730        assert_eq!(found.name, "copilot");
731    }
732
733    #[test]
734    fn resolve_unicast_locator_follows_did_when_binding_changes() {
735        let (_dir, registry) = temp_registry();
736        let mut record = grpc_lease("copilot", "did:key:zSame", "http://127.0.0.1:50151");
737        let _first = registry.publish(&record).unwrap();
738        record.a2a_url = "http://127.0.0.1:8080".to_string();
739        record.a2a_binding = A2ABinding::Jsonrpc;
740        let _second = registry.publish(&record).unwrap();
741        let (peer, locator) =
742            resolve_unicast_locator("did:key:zSame", &registry, None).expect("moved");
743        assert_eq!(peer.name, "copilot");
744        let locator = locator.expect("unicast locator");
745        assert_eq!(locator.binding, A2ABinding::Jsonrpc);
746        assert_eq!(locator.url, "http://127.0.0.1:8080");
747        assert_eq!(locator.display_uri(), "jsonrpc://127.0.0.1:8080");
748    }
749
750    #[test]
751    fn resolve_locator_returns_slim_node_for_slim_only_lease() {
752        let (_dir, registry) = temp_registry();
753        let mut record = grpc_lease("copilot", "did:key:zSlim", "");
754        record.a2a_url.clear();
755        record.slim_endpoint = "127.0.0.1:47357".to_string();
756        let _lease = registry.publish(&record).unwrap();
757        let (peer, locators) =
758            resolve_locator("agntcy/shadi/copilot-a2a", &registry, None).expect("channel alias");
759        assert_eq!(peer.did, "did:key:zSlim");
760        assert_eq!(locators.len(), 1);
761        assert_eq!(locators[0].binding, A2ABinding::Slim);
762        assert_eq!(locators[0].display_uri(), "slim://127.0.0.1:47357");
763        assert_eq!(peer.preferred_locator(), Some(locators[0].clone()));
764    }
765
766    #[test]
767    fn resolve_adapter_peer_rejects_ambiguous_did_and_missing_name() {
768        let (_dir, registry) = temp_registry();
769        let _a = registry
770            .publish(&grpc_lease(
771                "copilot",
772                "did:key:zShared",
773                "http://127.0.0.1:50151",
774            ))
775            .unwrap();
776        let _b = registry
777            .publish(&grpc_lease(
778                "codex",
779                "did:key:zShared",
780                "http://127.0.0.1:50152",
781            ))
782            .unwrap();
783        let err = resolve_adapter_peer("did:key:zShared", &registry, None).unwrap_err();
784        assert!(err.contains("ambiguous"), "{err}");
785        let missing = resolve_adapter_peer("goose", &registry, None).unwrap_err();
786        assert!(missing.contains("did:key"), "{missing}");
787    }
788
789    #[test]
790    fn extract_candidate_returns_none_without_authors() {
791        let mut record = a2a_record("copilot", "did:key:z6Mk...", "127.0.0.1:47357");
792        record.as_object_mut().unwrap().remove("authors");
793        assert!(extract_candidate(&record).is_none());
794    }
795
796    #[test]
797    fn extract_candidate_returns_none_without_a2a_module() {
798        let record = serde_json::json!({"authors": ["did:key:z6Mk..."], "modules": []});
799        assert!(extract_candidate(&record).is_none());
800    }
801
802    fn test_dir() -> DirLookupOptions {
803        DirLookupOptions {
804            server_addr: "localhost:9999".to_string(),
805            gh_token: None,
806            limit: 10,
807        }
808    }
809
810    #[test]
811    #[cfg(unix)]
812    fn parse_member_spec_skill_builds_a_working_skill_search_source() {
813        let _guard = crate::dir_registry::dirctl_env_lock().lock().expect("lock");
814        let record = a2a_record("copilot", "did:key:z6Mk...", "127.0.0.1:47357");
815        let (script, _dir) = fake_dirctl_script("bafkreitest", &record.to_string());
816        std::env::set_var("SHADI_DIRCTL_BINARY", &script);
817
818        let source =
819            parse_member_spec("skill:code_generation/implementation", &test_dir()).expect("parse");
820        let candidates = source.resolve().expect("resolve");
821
822        std::env::remove_var("SHADI_DIRCTL_BINARY");
823
824        assert_eq!(candidates.len(), 1);
825        assert_eq!(candidates[0].name, "copilot");
826    }
827
828    #[test]
829    #[cfg(unix)]
830    fn parse_member_spec_did_builds_a_working_did_lookup_source() {
831        let _guard = crate::dir_registry::dirctl_env_lock().lock().expect("lock");
832        let record = a2a_record("claude-code", "did:key:z6Mkagent", "127.0.0.1:47560");
833        let (script, _dir) = fake_dirctl_script("bafkreiagent", &record.to_string());
834        std::env::set_var("SHADI_DIRCTL_BINARY", &script);
835
836        let source = parse_member_spec("did:did:key:z6Mkagent", &test_dir()).expect("parse");
837        let candidates = source.resolve().expect("resolve");
838
839        std::env::remove_var("SHADI_DIRCTL_BINARY");
840
841        assert_eq!(candidates.len(), 1);
842        assert_eq!(candidates[0].did, "did:key:z6Mkagent");
843    }
844
845    #[test]
846    fn parse_member_spec_explicit_with_endpoint() {
847        let source =
848            parse_member_spec("explicit:avatar=did:key:human@127.0.0.1:47560", &test_dir())
849                .expect("parse");
850        let candidates = source.resolve().expect("resolve");
851        assert_eq!(candidates.len(), 1);
852        assert_eq!(candidates[0].name, "avatar");
853        assert_eq!(candidates[0].did, "did:key:human");
854        assert_eq!(
855            candidates[0].slim_endpoint.as_deref(),
856            Some("127.0.0.1:47560")
857        );
858    }
859
860    #[test]
861    fn parse_member_spec_explicit_without_endpoint() {
862        let source =
863            parse_member_spec("explicit:avatar=did:key:human", &test_dir()).expect("parse");
864        let candidates = source.resolve().expect("resolve");
865        assert_eq!(candidates[0].slim_endpoint, None);
866    }
867
868    #[test]
869    fn parse_member_spec_rejects_unknown_prefix() {
870        assert!(parse_member_spec("bogus:whatever", &test_dir()).is_err());
871    }
872
873    #[test]
874    fn parse_member_spec_rejects_explicit_without_equals() {
875        assert!(parse_member_spec("explicit:avatar", &test_dir()).is_err());
876    }
877
878    #[test]
879    fn parse_member_spec_rejects_empty_skill_name() {
880        let err = parse_member_spec("skill:", &test_dir()).err().unwrap();
881        assert!(err.contains("needs a skill name"));
882    }
883
884    #[test]
885    fn parse_member_spec_rejects_empty_did() {
886        let err = parse_member_spec("did:", &test_dir()).err().unwrap();
887        assert!(err.contains("needs a DID"));
888    }
889
890    #[test]
891    fn parse_member_spec_rejects_explicit_with_empty_name() {
892        let err = parse_member_spec("explicit:=did:key:human", &test_dir())
893            .err()
894            .unwrap();
895        assert!(err.contains("needs a name"));
896    }
897
898    #[test]
899    fn parse_member_spec_rejects_explicit_with_empty_did() {
900        let err = parse_member_spec("explicit:avatar=", &test_dir())
901            .err()
902            .unwrap();
903        assert!(err.contains("needs a DID"));
904    }
905
906    #[test]
907    fn resolve_members_concatenates_multiple_specs() {
908        let specs = vec![
909            "explicit:avatar=did:key:human".to_string(),
910            "explicit:claude-code=did:key:agent@127.0.0.1:47560".to_string(),
911        ];
912        let resolved = resolve_members(&specs, &test_dir()).expect("resolve");
913        assert_eq!(resolved.len(), 2);
914        assert_eq!(resolved[0].name, "avatar");
915        assert_eq!(resolved[1].name, "claude-code");
916    }
917
918    #[test]
919    fn explicit_list_source_returns_configured_entries_verbatim() {
920        let entries = vec![
921            CandidateMember {
922                name: "avatar".to_string(),
923                did: "did:key:human".to_string(),
924                slim_endpoint: None,
925                a2a_url: None,
926                a2a_binding: None,
927            },
928            CandidateMember {
929                name: "claude-code".to_string(),
930                did: "did:key:agent".to_string(),
931                slim_endpoint: Some("127.0.0.1:47560".to_string()),
932                a2a_url: None,
933                a2a_binding: None,
934            },
935        ];
936        let source = ExplicitListSource {
937            entries: entries.clone(),
938        };
939        assert_eq!(source.resolve().expect("resolve"), entries);
940    }
941
942    // ── dirctl-backed sources against a fake dirctl script ────────────────────
943    //
944    // These tests mutate the process-global SHADI_DIRCTL_BINARY env var, so
945    // they share a crate-wide lock with dir_registry's own dirctl-faking tests
946    // rather than a module-local one — two different locks guarding the same
947    // global would race under parallel test execution.
948
949    use crate::dir_registry::dirctl_env_lock as env_lock;
950
951    /// A fake `dirctl` that responds to `search ... --output jsonl` with one
952    /// CID and to `pull <cid> ... --output json` with a full a2a-module record,
953    /// so `SkillSearchSource`/`DidLookupSource` can be exercised end to end
954    /// without a real Directory server.
955    #[cfg(unix)]
956    fn fake_dirctl_script(cid: &str, record_json: &str) -> (std::path::PathBuf, tempfile::TempDir) {
957        let dir = tempfile::tempdir().expect("tempdir");
958        let path = dir.path().join("fake_dirctl.sh");
959        let script = format!(
960            r#"#!/bin/sh
961case "$1" in
962  search) echo '"{cid}"' ;;
963  pull) echo '{record_json}' ;;
964  *) exit 1 ;;
965esac
966"#
967        );
968        std::fs::write(&path, script).expect("write script");
969        use std::os::unix::fs::PermissionsExt;
970        std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o755)).expect("chmod");
971        (path, dir)
972    }
973
974    /// A fake `dirctl` where `search` always reports one CID, but the
975    /// caller controls what `pull` does with it — exits with `pull_exit`,
976    /// printing `pull_stdout`. Used to exercise `pull_record_json`'s and
977    /// `resolve_via_dirctl_query`'s failure/skip branches.
978    #[cfg(unix)]
979    fn fake_dirctl_script_pull_behavior(
980        cid: &str,
981        pull_exit: i32,
982        pull_stdout: &str,
983    ) -> (std::path::PathBuf, tempfile::TempDir) {
984        let dir = tempfile::tempdir().expect("tempdir");
985        let path = dir.path().join("fake_dirctl.sh");
986        let script = format!(
987            r#"#!/bin/sh
988case "$1" in
989  search) echo '"{cid}"' ;;
990  pull) echo '{pull_stdout}'; exit {pull_exit} ;;
991  *) exit 1 ;;
992esac
993"#
994        );
995        std::fs::write(&path, script).expect("write script");
996        use std::os::unix::fs::PermissionsExt;
997        std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o755)).expect("chmod");
998        (path, dir)
999    }
1000
1001    /// A fake `dirctl` where `search` itself fails (nonzero exit).
1002    #[cfg(unix)]
1003    fn fake_dirctl_script_search_fails() -> (std::path::PathBuf, tempfile::TempDir) {
1004        let dir = tempfile::tempdir().expect("tempdir");
1005        let path = dir.path().join("fake_dirctl.sh");
1006        std::fs::write(&path, "#!/bin/sh\necho 'boom' >&2\nexit 1\n").expect("write script");
1007        use std::os::unix::fs::PermissionsExt;
1008        std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o755)).expect("chmod");
1009        (path, dir)
1010    }
1011
1012    #[test]
1013    #[cfg(unix)]
1014    fn skill_search_source_resolves_candidate_via_search_then_pull() {
1015        let _guard = env_lock().lock().expect("lock");
1016        let record = a2a_record("copilot", "did:key:z6Mk...", "127.0.0.1:47357");
1017        let (script, _dir) = fake_dirctl_script("bafkreitest", &record.to_string());
1018        std::env::set_var("SHADI_DIRCTL_BINARY", &script);
1019
1020        let source = SkillSearchSource {
1021            skill: "code_generation/implementation".to_string(),
1022            dir: DirLookupOptions {
1023                server_addr: "localhost:9999".to_string(),
1024                gh_token: Some("tok".to_string()),
1025                limit: 10,
1026            },
1027        };
1028        let candidates = source.resolve().expect("resolve");
1029
1030        std::env::remove_var("SHADI_DIRCTL_BINARY");
1031
1032        assert_eq!(candidates.len(), 1);
1033        assert_eq!(candidates[0].name, "copilot");
1034        assert_eq!(candidates[0].did, "did:key:z6Mk...");
1035    }
1036
1037    #[test]
1038    #[cfg(unix)]
1039    fn search_cids_returns_err_when_search_exits_nonzero() {
1040        let _guard = env_lock().lock().expect("lock");
1041        let (script, _dir) = fake_dirctl_script_search_fails();
1042        std::env::set_var("SHADI_DIRCTL_BINARY", &script);
1043
1044        let result = search_cids(&["--skill", "x"], &test_dir());
1045
1046        std::env::remove_var("SHADI_DIRCTL_BINARY");
1047
1048        let err = result.unwrap_err();
1049        assert!(err.contains("boom"));
1050    }
1051
1052    #[test]
1053    #[cfg(unix)]
1054    fn resolve_via_dirctl_query_skips_a_cid_that_pulls_but_has_no_usable_card() {
1055        let _guard = env_lock().lock().expect("lock");
1056        // Pull succeeds, but the record has neither `authors` nor an
1057        // `integration/a2a` module — extract_candidate returns None, so
1058        // this CID is skipped rather than producing a bogus candidate.
1059        let (script, _dir) = fake_dirctl_script_pull_behavior("bafkreiempty", 0, "{}");
1060        std::env::set_var("SHADI_DIRCTL_BINARY", &script);
1061
1062        let candidates = resolve_via_dirctl_query(&["--skill", "x"], &test_dir()).expect("resolve");
1063
1064        std::env::remove_var("SHADI_DIRCTL_BINARY");
1065
1066        assert!(candidates.is_empty());
1067    }
1068
1069    #[test]
1070    #[cfg(unix)]
1071    fn resolve_via_dirctl_query_skips_a_cid_whose_pull_fails() {
1072        let _guard = env_lock().lock().expect("lock");
1073        let (script, _dir) = fake_dirctl_script_pull_behavior("bafkreifail", 1, "pull boom");
1074        std::env::set_var("SHADI_DIRCTL_BINARY", &script);
1075
1076        let candidates = resolve_via_dirctl_query(&["--skill", "x"], &test_dir()).expect("resolve");
1077
1078        std::env::remove_var("SHADI_DIRCTL_BINARY");
1079
1080        assert!(candidates.is_empty());
1081    }
1082
1083    #[test]
1084    #[cfg(unix)]
1085    fn did_lookup_source_resolves_candidate_via_author_search() {
1086        let _guard = env_lock().lock().expect("lock");
1087        let record = a2a_record("claude-code", "did:key:z6Mkagent", "127.0.0.1:47560");
1088        let (script, _dir) = fake_dirctl_script("bafkreiagent", &record.to_string());
1089        std::env::set_var("SHADI_DIRCTL_BINARY", &script);
1090
1091        let source = DidLookupSource {
1092            did: "did:key:z6Mkagent".to_string(),
1093            dir: DirLookupOptions {
1094                server_addr: "localhost:9999".to_string(),
1095                gh_token: None,
1096                limit: 10,
1097            },
1098        };
1099        let candidates = source.resolve().expect("resolve");
1100
1101        std::env::remove_var("SHADI_DIRCTL_BINARY");
1102
1103        assert_eq!(candidates.len(), 1);
1104        assert_eq!(candidates[0].name, "claude-code");
1105        assert_eq!(
1106            candidates[0].slim_endpoint.as_deref(),
1107            Some("127.0.0.1:47560")
1108        );
1109    }
1110
1111    #[test]
1112    #[cfg(unix)]
1113    fn resolve_adapter_peer_reports_dir_miss_when_card_did_differs() {
1114        let _guard = env_lock().lock().expect("lock");
1115        let record = a2a_record("copilot", "did:key:zOther", "127.0.0.1:47357");
1116        let (script, _dir) = fake_dirctl_script("bafkreiother", &record.to_string());
1117        std::env::set_var("SHADI_DIRCTL_BINARY", &script);
1118        let (_tmp, registry) = temp_registry();
1119        let err = resolve_adapter_peer("did:key:zWanted", &registry, Some(&test_dir())).unwrap_err();
1120        std::env::remove_var("SHADI_DIRCTL_BINARY");
1121        assert!(err.contains("was not found"), "{err}");
1122        assert!(err.contains("did:key:zWanted"), "{err}");
1123    }
1124
1125    #[test]
1126    fn search_cids_returns_dirctl_not_found_when_missing() {
1127        let _guard = env_lock().lock().expect("lock");
1128        std::env::set_var("SHADI_DIRCTL_BINARY", "/nonexistent/shadi_test_dirctl");
1129        let result = search_cids(
1130            &["--skill", "x"],
1131            &DirLookupOptions {
1132                server_addr: "localhost:9999".to_string(),
1133                gh_token: None,
1134                limit: 10,
1135            },
1136        );
1137        std::env::remove_var("SHADI_DIRCTL_BINARY");
1138        assert!(result.is_err());
1139    }
1140}