use std::collections::BTreeSet;
use std::time::Duration;
use affinidi_did_resolver_cache_sdk::DIDCacheClient;
use serde::Serialize;
use serde_json::Value;
use vta_sdk::protocol::matching::{Protocol, ServiceCapabilities, select_protocol};
const STEP_TIMEOUT: Duration = Duration::from_secs(10);
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize)]
#[serde(rename_all = "lowercase")]
pub enum Role {
Persona,
Vta,
Vtc,
Mediator,
}
impl Role {
#[must_use]
pub fn as_str(self) -> &'static str {
match self {
Role::Persona => "persona",
Role::Vta => "vta",
Role::Vtc => "vtc",
Role::Mediator => "mediator",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct ServiceEntry {
pub id: String,
pub types: Vec<String>,
pub endpoint: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(tag = "status", rename_all = "lowercase")]
pub enum Probe {
Reachable { url: String, http_status: u16 },
Unreachable { url: String, error: String },
Blocked { url: String, reason: String },
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum ProbePolicy {
#[default]
PublicOnly,
AllowPrivate,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ProbeGrade {
Ok,
Responding,
ServerError,
}
impl Probe {
#[must_use]
pub fn grade(&self) -> Option<ProbeGrade> {
match self {
Probe::Reachable { http_status, .. } => Some(match http_status {
200..=399 => ProbeGrade::Ok,
400..=499 => ProbeGrade::Responding,
_ => ProbeGrade::ServerError,
}),
Probe::Unreachable { .. } | Probe::Blocked { .. } => None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct Resolved {
pub services: Vec<ServiceEntry>,
pub tsp_endpoint: Option<String>,
pub didcomm_endpoint: Option<String>,
pub rest_endpoint: Option<String>,
pub probes: Vec<Probe>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct Party {
pub role: Role,
pub label: String,
pub did: String,
pub resolved: Option<Resolved>,
pub error: Option<String>,
}
impl Party {
fn caps(&self) -> ServiceCapabilities {
match &self.resolved {
Some(r) => ServiceCapabilities {
tsp: r.tsp_endpoint.clone(),
didcomm: r.didcomm_endpoint.clone(),
rest: r.rest_endpoint.clone(),
},
None => ServiceCapabilities::default(),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(tag = "outcome", rename_all = "snake_case")]
pub enum LinkOutcome {
Selected {
protocol: Protocol,
peer_endpoint: String,
},
NoCommonProtocol {
ours: Vec<Protocol>,
theirs: Vec<Protocol>,
},
Unknown { reason: String },
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct Link {
pub from: String,
pub to: String,
#[serde(flatten)]
pub outcome: LinkOutcome,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct HealthReport {
pub parties: Vec<Party>,
pub links: Vec<Link>,
pub notes: Vec<String>,
}
impl HealthReport {
#[must_use]
pub fn is_healthy(&self) -> bool {
self.parties.iter().all(|p| p.resolved.is_some())
&& !self
.links
.iter()
.any(|l| matches!(l.outcome, LinkOutcome::NoCommonProtocol { .. }))
}
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub enum Step {
ResolverStarting,
Resolving {
role: Role,
label: String,
did: String,
},
Resolved {
label: String,
services: usize,
transports: Vec<Protocol>,
elapsed: Duration,
},
ResolveFailed {
label: String,
error: String,
elapsed: Duration,
},
FollowingMediators { count: usize },
Probing { url: String },
Probed { probe: Probe, elapsed: Duration },
Negotiating { pairs: usize },
Finished { elapsed: Duration },
}
pub type ProgressFn<'a> = &'a (dyn Fn(Step) + Send + Sync);
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Subject {
pub role: Role,
pub label: String,
pub did: String,
}
impl Subject {
pub fn new(role: Role, label: impl Into<String>, did: impl Into<String>) -> Self {
Self {
role,
label: label.into(),
did: did.into(),
}
}
}
pub async fn build_report(subjects: &[Subject]) -> HealthReport {
build_report_with_progress(subjects, &|_| {}, ProbePolicy::PublicOnly).await
}
pub async fn build_report_with_progress(
subjects: &[Subject],
progress: ProgressFn<'_>,
policy: ProbePolicy,
) -> HealthReport {
let started = std::time::Instant::now();
progress(Step::ResolverStarting);
let host_policy = match policy {
ProbePolicy::PublicOnly => affinidi_did_web::HostPolicy::PublicOnly,
ProbePolicy::AllowPrivate => affinidi_did_web::HostPolicy::AllowPrivate,
};
let resolver = match DIDCacheClient::new(
affinidi_did_resolver_cache_sdk::config::DIDCacheConfigBuilder::default()
.with_host_policy(host_policy)
.build(),
)
.await
{
Ok(r) => r,
Err(e) => {
return HealthReport {
parties: Vec::new(),
links: Vec::new(),
notes: vec![format!(
"could not start a DID resolver ({e}) — nothing could be checked. \
This is a local fault, not a fault in the chain."
)],
};
}
};
let http = probe_client(policy);
let mut parties: Vec<Party> = Vec::new();
for subject in subjects {
if subject.did.is_empty() {
continue;
}
if parties.iter().any(|p| p.did == subject.did) {
continue;
}
parties.push(resolve_party(&resolver, http.as_ref(), policy, subject, progress).await);
}
let referenced = mediator_references(&parties);
let fresh: Vec<(String, String)> = referenced
.into_iter()
.filter(|(did, _)| !parties.iter().any(|p| p.did == *did))
.collect();
progress(Step::FollowingMediators { count: fresh.len() });
for (did, users) in fresh {
let subject = Subject::new(Role::Mediator, format!("mediator of {users}"), did);
parties.push(resolve_party(&resolver, http.as_ref(), policy, &subject, progress).await);
}
let links = negotiate_links(&parties);
progress(Step::Negotiating { pairs: links.len() });
let notes = collect_notes(&parties, &links);
progress(Step::Finished {
elapsed: started.elapsed(),
});
HealthReport {
parties,
links,
notes,
}
}
fn mediator_references(parties: &[Party]) -> std::collections::BTreeMap<String, String> {
let mut refs: std::collections::BTreeMap<String, Vec<(String, Vec<&'static str>)>> =
std::collections::BTreeMap::new();
for party in parties {
let Some(resolved) = &party.resolved else {
continue;
};
for (protocol, endpoint) in [
("TSP", resolved.tsp_endpoint.as_deref()),
("DIDComm", resolved.didcomm_endpoint.as_deref()),
] {
let Some(endpoint) = endpoint else { continue };
if !endpoint.starts_with("did:") {
continue;
}
let users = refs.entry(endpoint.to_string()).or_default();
match users.iter_mut().find(|(label, _)| *label == party.label) {
Some((_, protocols)) => protocols.push(protocol),
None => users.push((party.label.clone(), vec![protocol])),
}
}
}
refs.into_iter()
.map(|(did, users)| {
let described = users
.into_iter()
.map(|(label, protocols)| format!("{label} ({})", protocols.join(", ")))
.collect::<Vec<_>>()
.join(", ");
(did, described)
})
.collect()
}
async fn resolve_party(
resolver: &DIDCacheClient,
http: Option<&reqwest::Client>,
policy: ProbePolicy,
subject: &Subject,
progress: ProgressFn<'_>,
) -> Party {
let started = std::time::Instant::now();
progress(Step::Resolving {
role: subject.role,
label: subject.label.clone(),
did: subject.did.clone(),
});
let fail = |error: String| {
progress(Step::ResolveFailed {
label: subject.label.clone(),
error: error.clone(),
elapsed: started.elapsed(),
});
Party {
role: subject.role,
label: subject.label.clone(),
did: subject.did.clone(),
resolved: None,
error: Some(error),
}
};
let resolved = match tokio::time::timeout(STEP_TIMEOUT, resolver.resolve(&subject.did)).await {
Ok(Ok(resolved)) => resolved,
Ok(Err(e)) => return fail(e.to_string()),
Err(_) => {
return fail(format!(
"resolution timed out after {}s",
STEP_TIMEOUT.as_secs()
));
}
};
let doc = match serde_json::to_value(&resolved.doc) {
Ok(doc) => doc,
Err(e) => return fail(format!("document could not be re-serialised: {e}")),
};
let caps = ServiceCapabilities::from_did_document(&doc);
let services = service_entries(&doc);
progress(Step::Resolved {
label: subject.label.clone(),
services: services.len(),
transports: caps.advertised(),
elapsed: started.elapsed(),
});
let mut probes = Vec::new();
if let Some(client) = http {
let urls: BTreeSet<&str> = services
.iter()
.filter(|s| is_routable(s))
.map(|s| s.endpoint.as_str())
.filter(|e| e.starts_with("http://") || e.starts_with("https://"))
.collect();
for url in urls {
progress(Step::Probing {
url: url.to_string(),
});
let probe_started = std::time::Instant::now();
let result = probe(client, url, policy).await;
progress(Step::Probed {
probe: result.clone(),
elapsed: probe_started.elapsed(),
});
probes.push(result);
}
}
Party {
role: subject.role,
label: subject.label.clone(),
did: subject.did.clone(),
resolved: Some(Resolved {
services,
tsp_endpoint: caps.tsp,
didcomm_endpoint: caps.didcomm,
rest_endpoint: caps.rest,
probes,
}),
error: None,
}
}
const NON_ROUTABLE_SERVICE_TYPES: [&str; 2] = ["relativeRef", "LinkedVerifiablePresentation"];
fn is_routable(service: &ServiceEntry) -> bool {
!service
.types
.iter()
.any(|t| NON_ROUTABLE_SERVICE_TYPES.contains(&t.as_str()))
}
fn service_entries(doc: &Value) -> Vec<ServiceEntry> {
let Some(services) = doc.get("service").and_then(Value::as_array) else {
return Vec::new();
};
services
.iter()
.map(|svc| {
let id = svc
.get("id")
.and_then(Value::as_str)
.unwrap_or("<no id>")
.to_string();
let types = match svc.get("type") {
Some(Value::String(s)) => vec![s.clone()],
Some(Value::Array(arr)) => arr
.iter()
.filter_map(Value::as_str)
.map(str::to_string)
.collect(),
_ => Vec::new(),
};
let endpoint = svc
.get("serviceEndpoint")
.and_then(endpoint_uri)
.unwrap_or_else(|| "<unreadable>".to_string());
ServiceEntry {
id,
types,
endpoint,
}
})
.collect()
}
fn endpoint_uri(endpoint: &Value) -> Option<String> {
match endpoint {
Value::String(s) => Some(s.clone()),
Value::Object(map) => map.get("uri")?.as_str().map(str::to_string),
Value::Array(arr) => arr.iter().find_map(endpoint_uri),
_ => None,
}
}
fn probe_client(policy: ProbePolicy) -> Option<reqwest::Client> {
let mut builder = reqwest::Client::builder()
.timeout(STEP_TIMEOUT)
.connect_timeout(STEP_TIMEOUT)
.redirect(reqwest::redirect::Policy::none())
.no_proxy();
if policy == ProbePolicy::PublicOnly {
builder = builder.dns_resolver(affinidi_did_web::guarded_dns_resolver());
}
builder.build().ok()
}
async fn probe(client: &reqwest::Client, url: &str, policy: ProbePolicy) -> Probe {
let parsed = match vet_probe_url(url, policy) {
Ok(parsed) => parsed,
Err(reason) => {
return Probe::Blocked {
url: url.to_string(),
reason,
};
}
};
match client.get(parsed).send().await {
Ok(response) => Probe::Reachable {
url: url.to_string(),
http_status: response.status().as_u16(),
},
Err(e) => send_failure(url, &e),
}
}
pub(crate) fn vet_probe_url(url: &str, policy: ProbePolicy) -> Result<reqwest::Url, String> {
let parsed = reqwest::Url::parse(url).map_err(|e| format!("unparseable URL: {e}"))?;
match (parsed.scheme(), policy) {
("https", _) | ("http", ProbePolicy::AllowPrivate) => {}
("http", ProbePolicy::PublicOnly) => {
return Err("plaintext http".to_string());
}
(other, _) => return Err(format!("{other} scheme")),
}
if !parsed.username().is_empty() || parsed.password().is_some() {
return Err("URL carries userinfo".to_string());
}
if policy == ProbePolicy::PublicOnly
&& let Some(host) = parsed.host_str()
&& crate::net_guard::is_blocked_host(host)
{
return Err(format!("non-public host {host}"));
}
Ok(parsed)
}
fn send_failure(url: &str, error: &reqwest::Error) -> Probe {
match dns_refusal_in_chain(error) {
Some(reason) => Probe::Blocked {
url: url.to_string(),
reason,
},
None => Probe::Unreachable {
url: url.to_string(),
error: error.to_string(),
},
}
}
fn dns_refusal_in_chain(error: &(dyn std::error::Error + 'static)) -> Option<String> {
affinidi_net_guard::blocked_in_chain(error).map(ToString::to_string)
}
fn negotiate_links(parties: &[Party]) -> Vec<Link> {
let personas: Vec<&Party> = parties.iter().filter(|p| p.role == Role::Persona).collect();
let peers: Vec<&Party> = parties
.iter()
.filter(|p| matches!(p.role, Role::Vta | Role::Vtc))
.collect();
let mut links = Vec::new();
for persona in &personas {
for peer in &peers {
let outcome = if persona.resolved.is_none() {
LinkOutcome::Unknown {
reason: format!("{} did not resolve", persona.label),
}
} else if peer.resolved.is_none() {
LinkOutcome::Unknown {
reason: format!("{} did not resolve", peer.label),
}
} else {
match select_protocol(&persona.caps(), &peer.caps(), &peer.did) {
Ok(m) => LinkOutcome::Selected {
protocol: m.protocol,
peer_endpoint: m.peer_endpoint,
},
Err(_) => LinkOutcome::NoCommonProtocol {
ours: persona.caps().advertised(),
theirs: peer.caps().advertised(),
},
}
};
links.push(Link {
from: persona.label.clone(),
to: peer.label.clone(),
outcome,
});
}
}
links
}
fn collect_notes(parties: &[Party], links: &[Link]) -> Vec<String> {
let mut notes = Vec::new();
for party in parties {
if let Some(error) = &party.error {
notes.push(format!(
"{} ({}) did not resolve: {error}",
party.label, party.did
));
continue;
}
let Some(resolved) = &party.resolved else {
continue;
};
if resolved.services.is_empty() {
notes.push(format!(
"{} publishes no service endpoints at all — nothing can route to it.",
party.label
));
} else if resolved.tsp_endpoint.is_none()
&& resolved.didcomm_endpoint.is_none()
&& party.role != Role::Mediator
{
notes.push(format!(
"{} advertises no TSP or DIDComm transport (services present, but none of a \
recognised type) — it cannot be messaged.",
party.label
));
}
for probe in &resolved.probes {
match probe {
Probe::Unreachable { url, error } => {
notes.push(format!("{}: {url} is unreachable ({error}).", party.label));
}
Probe::Reachable { url, http_status }
if probe.grade() == Some(ProbeGrade::ServerError) =>
{
notes.push(format!(
"{}: {url} answered HTTP {http_status} — the host is up but the \
service behind it is failing.",
party.label
));
}
Probe::Reachable { .. } => {}
Probe::Blocked { url, reason } => {
notes.push(format!(
"{}: {url} not probed ({reason}) — a DID document advertising a \
plaintext or non-public endpoint is itself worth a look.",
party.label
));
}
}
}
}
for link in links {
match &link.outcome {
LinkOutcome::NoCommonProtocol { ours, theirs } => {
let mut note = format!(
"{} and {} share no transport: we offer [{}], they offer [{}]. A message \
between them cannot be sent until one side adds the other's.",
link.from,
link.to,
join_protocols(ours),
join_protocols(theirs),
);
if ours.as_slice() == [Protocol::Didcomm] && theirs.as_slice() == [Protocol::Tsp] {
note.push_str(
" This persona's document predates `#tsp` being requested at mint \
time, so it advertises DIDComm only; a persona minted by a current \
client carries both. Re-minting a persona for this community is the \
fix — the existing document will not gain the service on its own.",
);
}
notes.push(note);
}
LinkOutcome::Unknown { reason } => notes.push(format!(
"{} → {}: transport could not be determined ({reason}).",
link.from, link.to
)),
LinkOutcome::Selected { .. } => {}
}
}
let mediators: BTreeSet<&str> = parties
.iter()
.filter(|p| p.role == Role::Mediator)
.map(|p| p.did.as_str())
.collect();
if mediators.len() > 1 {
notes.push(format!(
"{} distinct mediators are in play; messages cross between them. This is \
supported, but it means a delivery problem can live in either.",
mediators.len()
));
}
notes
}
fn join_protocols(protocols: &[Protocol]) -> String {
if protocols.is_empty() {
return "none".to_string();
}
protocols
.iter()
.map(|p| p.as_str())
.collect::<Vec<_>>()
.join(", ")
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn resolved_party(role: Role, label: &str, did: &str, doc: &Value) -> Party {
let caps = ServiceCapabilities::from_did_document(doc);
Party {
role,
label: label.to_string(),
did: did.to_string(),
resolved: Some(Resolved {
services: service_entries(doc),
tsp_endpoint: caps.tsp,
didcomm_endpoint: caps.didcomm,
rest_endpoint: caps.rest,
probes: Vec::new(),
}),
error: None,
}
}
fn persona_doc(mediator: &str) -> Value {
json!({
"service": [
{
"id": "did:webvh:scid:example:persona#tsp",
"type": "TSPTransport",
"serviceEndpoint": mediator,
},
{
"id": "did:webvh:scid:example:persona#vta-didcomm",
"type": "DIDCommMessaging",
"serviceEndpoint": [{ "uri": mediator, "accept": ["didcomm/v2"] }],
},
]
})
}
#[test]
fn services_are_read_verbatim_including_unrecognised_types() {
let doc = json!({
"service": [
{ "id": "#tsp", "type": "TSPTransport", "serviceEndpoint": "did:webvh:m" },
{ "id": "#odd", "type": "SomeFutureTransport", "serviceEndpoint": "https://x/y" },
]
});
let entries = service_entries(&doc);
assert_eq!(entries.len(), 2);
assert_eq!(entries[1].types, vec!["SomeFutureTransport".to_string()]);
assert_eq!(
entries[1].endpoint, "https://x/y",
"an unrecognised transport must still be shown — a party that publishes \
one looks empty through ServiceCapabilities alone, and telling that \
apart from publishing nothing is the point of the map"
);
}
#[test]
fn endpoint_and_type_shapes_are_all_read() {
assert_eq!(endpoint_uri(&json!("https://a")), Some("https://a".into()));
assert_eq!(
endpoint_uri(&json!({"uri": "did:webvh:m"})),
Some("did:webvh:m".into())
);
assert_eq!(
endpoint_uri(&json!([{"uri": "did:webvh:m"}])),
Some("did:webvh:m".into())
);
let doc = json!({
"service": [{
"id": "#both", "type": ["DIDCommMessaging", "Other"],
"serviceEndpoint": [{ "uri": "did:webvh:m" }],
}]
});
assert_eq!(service_entries(&doc)[0].types.len(), 2);
}
#[test]
fn a_shared_transport_negotiates_and_names_the_peer_mediator() {
let parties = vec![
resolved_party(
Role::Persona,
"persona",
"did:webvh:p",
&persona_doc("did:webvh:our-mediator"),
),
resolved_party(
Role::Vtc,
"VTC",
"did:webvh:v",
&persona_doc("did:webvh:their-mediator"),
),
];
let links = negotiate_links(&parties);
assert_eq!(links.len(), 1);
match &links[0].outcome {
LinkOutcome::Selected {
protocol,
peer_endpoint,
} => {
assert_eq!(*protocol, Protocol::Tsp, "TSP outranks DIDComm");
assert_eq!(peer_endpoint, "did:webvh:their-mediator");
}
other => panic!("expected a selected protocol, got {other:?}"),
}
}
#[test]
fn a_disjoint_pair_reports_both_sides() {
let tsp_only = json!({
"service": [{ "id": "#tsp", "type": "TSPTransport", "serviceEndpoint": "did:webvh:m1" }]
});
let didcomm_only = json!({
"service": [{
"id": "#dc", "type": "DIDCommMessaging",
"serviceEndpoint": [{ "uri": "did:webvh:m2" }],
}]
});
let parties = vec![
resolved_party(Role::Persona, "persona", "did:webvh:p", &tsp_only),
resolved_party(Role::Vtc, "VTC", "did:webvh:v", &didcomm_only),
];
let links = negotiate_links(&parties);
assert!(matches!(
links[0].outcome,
LinkOutcome::NoCommonProtocol { .. }
));
let notes = collect_notes(&parties, &links);
let note = notes
.iter()
.find(|n| n.contains("share no transport"))
.expect("a disjoint pair must produce a note");
assert!(note.contains("tsp"), "our side must be named: {note}");
assert!(note.contains("didcomm"), "their side must be named: {note}");
assert!(
!HealthReport {
parties,
links,
notes,
}
.is_healthy()
);
}
#[test]
fn an_unresolvable_peer_is_named_with_its_reason() {
let parties = vec![
resolved_party(
Role::Persona,
"persona",
"did:webvh:p",
&persona_doc("did:webvh:m"),
),
Party {
role: Role::Vtc,
label: "VTC".into(),
did: "did:webvh:missing".into(),
resolved: None,
error: Some("404 fetching did.jsonl".into()),
},
];
let links = negotiate_links(&parties);
assert!(matches!(links[0].outcome, LinkOutcome::Unknown { .. }));
let notes = collect_notes(&parties, &links);
assert!(
notes
.iter()
.any(|n| n.contains("did:webvh:missing") && n.contains("404")),
"the failing DID and its reason must both appear: {notes:?}"
);
assert!(
!HealthReport {
parties,
links,
notes
}
.is_healthy()
);
}
#[test]
fn a_shared_mediator_names_every_party_and_transport() {
let parties = vec![
resolved_party(
Role::Persona,
"persona joy-ahead",
"did:webvh:p",
&persona_doc("did:webvh:shared"),
),
resolved_party(
Role::Vtc,
"VTC acme",
"did:webvh:v",
&json!({
"service": [{
"id": "#dc", "type": "DIDCommMessaging",
"serviceEndpoint": [{ "uri": "did:webvh:shared" }],
}]
}),
),
];
let refs = mediator_references(&parties);
assert_eq!(refs.len(), 1, "one host, resolved once: {refs:?}");
let label = &refs["did:webvh:shared"];
assert_eq!(
label, "persona joy-ahead (TSP, DIDComm), VTC acme (DIDComm)",
"both parties, and TSP before DIDComm (negotiation order, not \
alphabetical): {label}"
);
}
#[test]
fn document_adjacent_services_are_not_probed() {
let files = ServiceEntry {
id: "#files".into(),
types: vec!["relativeRef".into()],
endpoint: "https://webvh.storm.ws/army-provide".into(),
};
let whois = ServiceEntry {
id: "#whois".into(),
types: vec!["LinkedVerifiablePresentation".into()],
endpoint: "https://webvh.storm.ws/army-provide/whois.vp".into(),
};
assert!(
!is_routable(&files),
"#files is the document host, not a route"
);
assert!(
!is_routable(&whois),
"#whois is not somewhere a message goes"
);
}
#[test]
fn transports_and_unknown_types_are_still_probed() {
for types in [
vec!["TSPTransport".to_string()],
vec!["DIDCommMessaging".to_string()],
vec!["Authentication".to_string()],
vec!["SomeFutureTransport".to_string()],
] {
let entry = ServiceEntry {
id: "#x".into(),
types: types.clone(),
endpoint: "https://example/x".into(),
};
assert!(
is_routable(&entry),
"{types:?} names a route (or might); it must be probed"
);
}
}
#[test]
fn probe_status_is_graded_not_flattened() {
let at = |status| Probe::Reachable {
url: "https://m/x".into(),
http_status: status,
};
assert_eq!(at(200).grade(), Some(ProbeGrade::Ok));
assert_eq!(
at(405).grade(),
Some(ProbeGrade::Responding),
"a POST/websocket endpoint declining a GET is the host working"
);
assert_eq!(at(404).grade(), Some(ProbeGrade::Responding));
assert_eq!(
at(503).grade(),
Some(ProbeGrade::ServerError),
"host up, application failing — the case a flat `reachable` hid"
);
assert_eq!(
Probe::Unreachable {
url: "https://m/x".into(),
error: "dns".into()
}
.grade(),
None
);
}
#[test]
fn only_a_server_error_becomes_a_finding() {
let with_probe = |probe: Probe| {
let mut party = resolved_party(
Role::Mediator,
"mediator",
"did:webvh:m",
&json!({"service": [{
"id": "#tsp", "type": "TSPTransport",
"serviceEndpoint": "https://m/x",
}]}),
);
party.resolved.as_mut().expect("resolved").probes = vec![probe];
collect_notes(&[party], &[])
};
let five_hundred = with_probe(Probe::Reachable {
url: "https://m/x".into(),
http_status: 502,
});
assert!(
five_hundred.iter().any(|n| n.contains("502")),
"a 5xx must surface: {five_hundred:?}"
);
let four_oh_four = with_probe(Probe::Reachable {
url: "https://m/x".into(),
http_status: 404,
});
assert!(
four_oh_four.is_empty(),
"a 4xx on a non-GET endpoint is normal and must stay quiet: {four_oh_four:?}"
);
}
#[test]
fn a_didcomm_only_persona_is_told_why_it_cannot_reach_a_tsp_community() {
let didcomm_only = json!({"service": [{
"id": "#vta-didcomm", "type": "DIDCommMessaging",
"serviceEndpoint": [{"uri": "did:webvh:m"}],
}]});
let tsp_only = json!({"service": [{
"id": "#tsp", "type": "TSPTransport", "serviceEndpoint": "did:webvh:m",
}]});
let parties = vec![
resolved_party(
Role::Persona,
"persona hello-fury",
"did:webvh:p",
&didcomm_only,
),
resolved_party(Role::Vtc, "VTC first-vtc", "did:webvh:v", &tsp_only),
];
let links = negotiate_links(&parties);
let notes = collect_notes(&parties, &links);
let note = notes
.iter()
.find(|n| n.contains("share no transport"))
.expect("disjoint pair must be reported");
assert!(
note.contains("predates") && note.contains("Re-minting"),
"the note must say why this persona is stuck and what fixes it, not \
just list the two sets: {note}"
);
}
#[test]
fn the_remint_advice_is_not_given_for_the_reverse_mismatch() {
let didcomm_only = json!({"service": [{
"id": "#dc", "type": "DIDCommMessaging",
"serviceEndpoint": [{"uri": "did:webvh:m"}],
}]});
let tsp_only = json!({"service": [{
"id": "#tsp", "type": "TSPTransport", "serviceEndpoint": "did:webvh:m",
}]});
let parties = vec![
resolved_party(Role::Persona, "persona", "did:webvh:p", &tsp_only),
resolved_party(Role::Vtc, "VTC", "did:webvh:v", &didcomm_only),
];
let notes = collect_notes(&parties, &negotiate_links(&parties));
let note = notes
.iter()
.find(|n| n.contains("share no transport"))
.expect("still reported");
assert!(
!note.contains("Re-minting"),
"re-minting our persona does not fix a DIDComm-only community: {note}"
);
}
#[test]
fn a_url_endpoint_is_not_followed_as_a_mediator() {
let parties = vec![resolved_party(
Role::Mediator,
"mediator",
"did:webvh:m",
&json!({
"service": [{
"id": "#tsp", "type": "TSPTransport",
"serviceEndpoint": "https://mediator.example/mediator/v1",
}]
}),
)];
assert!(
mediator_references(&parties).is_empty(),
"following a transport URL as a DID would resolve nothing and add a \
bogus party to the map"
);
}
#[test]
fn split_mediators_are_reported_without_being_called_a_fault() {
let parties = vec![
resolved_party(
Role::Persona,
"persona",
"did:webvh:p",
&persona_doc("did:webvh:m1"),
),
resolved_party(
Role::Vtc,
"VTC",
"did:webvh:v",
&persona_doc("did:webvh:m2"),
),
resolved_party(
Role::Mediator,
"mediator of persona",
"did:webvh:m1",
&json!({}),
),
resolved_party(
Role::Mediator,
"mediator of VTC",
"did:webvh:m2",
&json!({}),
),
];
let links = negotiate_links(&parties);
let notes = collect_notes(&parties, &links);
assert!(
notes.iter().any(|n| n.contains("distinct mediators")),
"split mediators must be surfaced: {notes:?}"
);
assert!(
HealthReport {
parties,
links,
notes
}
.is_healthy(),
"two mediators is a supported topology, not a failure"
);
}
#[test]
fn a_mediator_is_not_faulted_for_advertising_no_mediator() {
let mediator_doc = json!({
"service": [{
"id": "#tsp", "type": "TSPTransport",
"serviceEndpoint": "https://mediator.example/mediator/v1",
}]
});
let parties = vec![resolved_party(
Role::Mediator,
"mediator of persona",
"did:webvh:m",
&mediator_doc,
)];
let notes = collect_notes(&parties, &[]);
assert!(
!notes.iter().any(|n| n.contains("cannot be messaged")),
"a mediator publishing a transport URL is healthy: {notes:?}"
);
}
async fn loopback_listener() -> (tokio::net::TcpListener, u16) {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind a loopback listener");
let port = listener.local_addr().expect("listener address").port();
(listener, port)
}
async fn assert_never_dialled(listener: &tokio::net::TcpListener) {
assert!(
tokio::time::timeout(Duration::from_millis(200), listener.accept())
.await
.is_err(),
"something connected to the stand-in internal service"
);
}
fn public_client() -> reqwest::Client {
probe_client(ProbePolicy::PublicOnly).expect("the probe client builds")
}
#[tokio::test]
async fn probe_refuses_loopback_literal_without_dialing() {
let (listener, port) = loopback_listener().await;
let url = format!("https://127.0.0.1:{port}/latest/meta-data/");
let result = probe(&public_client(), &url, ProbePolicy::PublicOnly).await;
assert!(
matches!(&result, Probe::Blocked { url: u, reason } if *u == url && reason.contains("127.0.0.1")),
"a loopback literal must be refused by name: {result:?}"
);
assert_eq!(result.grade(), None, "nothing was asked, so nothing grades");
assert_never_dialled(&listener).await;
}
#[tokio::test]
async fn probe_refuses_plain_http() {
let (listener, port) = loopback_listener().await;
let url = format!("http://127.0.0.1:{port}/latest/meta-data/");
let result = probe(&public_client(), &url, ProbePolicy::PublicOnly).await;
assert!(
matches!(&result, Probe::Blocked { reason, .. } if reason.contains("plaintext")),
"a plaintext endpoint is listed, not dialled: {result:?}"
);
assert_never_dialled(&listener).await;
}
#[tokio::test]
async fn probe_client_dns_guard_refuses_localhost_name() {
let (listener, port) = loopback_listener().await;
let url = format!("https://localhost:{port}/");
let error = public_client()
.get(&url)
.send()
.await
.expect_err("the guarded resolver must refuse a name resolving to loopback");
let refusal = dns_refusal_in_chain(&error);
assert!(
refusal
.as_deref()
.is_some_and(|r| r.contains("blocked address") && r.contains("localhost")),
"the guard's refusal must be recoverable from the error chain, and must \
name the host it refused: {error:?}"
);
assert!(
matches!(send_failure(&url, &error), Probe::Blocked { .. }),
"a DNS refusal is a block, not an unreachable host"
);
assert_never_dialled(&listener).await;
}
#[tokio::test]
async fn health_report_refuses_a_localhost_webvh_did_without_dialing() {
let (listener, port) = loopback_listener().await;
let did = format!("did:webvh:QmStandInScidAAAAAAAAAAAAAAAAAAAA:localhost%3A{port}:agent");
let report = build_report(&[Subject::new(Role::Vta, "stand-in VTA", did.clone())]).await;
let party = report
.parties
.iter()
.find(|p| p.did == did)
.expect("a subject is reported whether or not it resolves");
assert!(
party.resolved.is_none(),
"a refused host must not yield a resolved party: {party:?}"
);
assert!(
party
.error
.as_deref()
.is_some_and(|e| e.contains("BlockedHost")),
"the refusal must be reported as a blocked host, not as a timeout or \
an unreachable one: {:?}",
party.error
);
assert_never_dialled(&listener).await;
}
#[tokio::test]
async fn probe_does_not_follow_redirects() {
use wiremock::matchers::method;
use wiremock::{Mock, MockServer, ResponseTemplate};
let target = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(200))
.mount(&target)
.await;
let redirector = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(
ResponseTemplate::new(302)
.insert_header("Location", format!("{}/latest/meta-data/", target.uri())),
)
.mount(&redirector)
.await;
let client = probe_client(ProbePolicy::AllowPrivate).expect("the probe client builds");
let result = probe(&client, &redirector.uri(), ProbePolicy::AllowPrivate).await;
assert_eq!(
result,
Probe::Reachable {
url: redirector.uri(),
http_status: 302,
}
);
assert!(
target
.received_requests()
.await
.expect("request recording is on")
.is_empty(),
"the redirect target must never be requested"
);
}
#[tokio::test]
async fn health_report_blocks_did_peer_inline_internal_service() {
use base64::Engine as _;
use base64::engine::general_purpose::URL_SAFE_NO_PAD;
let (listener, port) = loopback_listener().await;
let did_with_service = |endpoint: String| {
let service = json!({ "t": "dm", "s": endpoint }).to_string();
format!(
"did:peer:2.Vz6MkhaXgBZDvotDkL5257faiztiGiC2QtKLGpbnnEGta2doK.S{}",
URL_SAFE_NO_PAD.encode(service)
)
};
let report = build_report(&[
Subject::new(
Role::Vtc,
"plaintext",
did_with_service(format!("http://127.0.0.1:{port}/latest/meta-data/")),
),
Subject::new(
Role::Vtc,
"https",
did_with_service(format!("https://127.0.0.1:{port}/latest/meta-data/")),
),
])
.await;
assert_eq!(report.parties.len(), 2, "{report:#?}");
for party in &report.parties {
let resolved = party
.resolved
.as_ref()
.unwrap_or_else(|| panic!("did:peer resolves offline: {party:#?}"));
assert!(
matches!(resolved.probes.as_slice(), [Probe::Blocked { .. }]),
"{}: the inline internal endpoint must be blocked: {:?}",
party.label,
resolved.probes
);
}
assert_eq!(
report
.notes
.iter()
.filter(|n| n.contains("not probed"))
.count(),
2,
"each blocked endpoint is a finding: {:?}",
report.notes
);
assert_never_dialled(&listener).await;
}
#[test]
fn probe_urls_are_vetted_after_canonicalisation() {
for url in [
"https://2130706433/",
"https://0x7f000001/",
"https://017700000001/",
"https://0177.0.0.1/",
"https://0x7f.0.0.1/",
"https://127.1/",
"https://127.0.1/",
"https://0/",
"https://%31%32%37.0.0.1/",
"https://169.254.169.254./",
"https://100.100.100.200/",
"https://[::ffff:127.0.0.1]/",
"https://[0:0:0:0:0:ffff:7f00:1]/",
"https://[::1]:8443/",
"https://[fd00:ec2::254]/",
"https://localhost/",
"https://LOCALHOST./",
"https://svc.localhost/",
"https://printer.local/",
"https://kube-dns.kube-system.svc.cluster.local/",
"https://example.com@127.0.0.1/",
"https://127.0.0.1\\@example.com/",
"https://user:pass@example.com/",
"https://0x100000000/",
"https://1.2.3.4.5/",
"https://[fe80::1%25en0]/",
"not a url",
"http://example.com/",
"ws://example.com/",
"ftp://example.com/",
"file:///etc/passwd",
"gopher://example.com/",
"data:text/plain,x",
"javascript:alert(1)",
"blob:https://x/y",
] {
assert!(
vet_probe_url(url, ProbePolicy::PublicOnly).is_err(),
"{url} must not be probed"
);
}
for url in [
"https://example.com/",
"https://example.com./",
"https://localhost.example.com/",
"https://8.8.8.8/",
"https://[2606:4700:4700::1111]/",
] {
assert!(
vet_probe_url(url, ProbePolicy::PublicOnly).is_ok(),
"{url} is a public HTTPS endpoint and must be probed"
);
}
}
#[test]
fn allow_private_admits_dev_endpoints_but_not_credentials() {
for url in [
"http://localhost:8000/",
"http://127.0.0.1:9099/",
"http://[::1]:7037/",
"https://10.0.0.5/",
] {
assert!(
vet_probe_url(url, ProbePolicy::AllowPrivate).is_ok(),
"{url} is a dev endpoint AllowPrivate exists for"
);
}
for url in [
"https://user:pass@10.0.0.5/",
"ftp://127.0.0.1/",
"not a url",
] {
assert!(
vet_probe_url(url, ProbePolicy::AllowPrivate).is_err(),
"{url} must still be refused"
);
}
}
#[test]
fn a_blocked_probe_is_a_finding_but_not_a_fault() {
let blocked = Probe::Blocked {
url: "https://127.0.0.1/".into(),
reason: "non-public host 127.0.0.1".into(),
};
assert_eq!(
serde_json::to_value(&blocked).expect("serialises"),
json!({
"status": "blocked",
"url": "https://127.0.0.1/",
"reason": "non-public host 127.0.0.1",
})
);
let mut party = resolved_party(
Role::Vtc,
"VTC",
"did:peer:2.x",
&json!({"service": [{
"id": "#dc", "type": "DIDCommMessaging",
"serviceEndpoint": "https://127.0.0.1/",
}]}),
);
party.resolved.as_mut().expect("resolved").probes = vec![blocked];
let parties = vec![party];
let notes = collect_notes(&parties, &[]);
assert!(
notes
.iter()
.any(|n| n.contains("https://127.0.0.1/ not probed (non-public host 127.0.0.1)")),
"{notes:?}"
);
assert!(
HealthReport {
parties,
links: Vec::new(),
notes,
}
.is_healthy()
);
}
}