use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, OnceLock};
use serde::Serialize;
#[async_trait::async_trait]
pub trait MeshHost: Send + Sync {
async fn status_json(&self) -> anyhow::Result<serde_json::Value>;
async fn ensure_funnel(&self, port: u16) -> anyhow::Result<String>;
async fn funnel_url(&self, port: u16) -> Option<String>;
}
fn host_slot() -> &'static OnceLock<Arc<dyn MeshHost>> {
static HOST: OnceLock<Arc<dyn MeshHost>> = OnceLock::new();
&HOST
}
pub fn set_global_host(host: Arc<dyn MeshHost>) {
let _ = host_slot().set(host);
}
fn host() -> Option<Arc<dyn MeshHost>> {
host_slot().get().cloned()
}
pub fn is_insecure_auth_token_placeholder(token: &str) -> bool {
const PLACEHOLDERS: &[&str] = &[
"CHANGE_ME",
"CHANGEME",
"REPLACE_ME",
"REPLACEME",
"YOUR_TOKEN_HERE",
"TOKEN",
"SECRET",
"PASSWORD",
];
let trimmed = token.trim();
PLACEHOLDERS
.iter()
.any(|placeholder| trimmed.eq_ignore_ascii_case(placeholder))
}
#[derive(Clone, Default)]
pub struct MeshHandle;
impl MeshHandle {
pub fn new() -> Self {
Self
}
pub async fn status(&self) -> MeshStatus {
query_status().await
}
pub fn enabled(&self) -> bool {
is_enabled()
}
}
static MESH_PREF_ENABLED: AtomicBool = AtomicBool::new(false);
pub fn set_pref_enabled(enabled: bool) {
MESH_PREF_ENABLED.store(enabled, Ordering::Relaxed);
}
pub fn parse_enabled(value: Option<&str>) -> bool {
match value.map(|v| v.trim().to_ascii_lowercase()).as_deref() {
None | Some("") | Some("0") | Some("false") | Some("no") => false,
Some(_) => true,
}
}
pub fn is_enabled() -> bool {
match std::env::var("RYU_MESH_ENABLED").ok() {
Some(v) => parse_enabled(Some(&v)),
None => MESH_PREF_ENABLED.load(Ordering::Relaxed),
}
}
#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
pub struct MeshPeer {
pub name: String,
pub host_or_dns: String,
pub magic_dns_name: String,
pub tailscale_ips: Vec<String>,
pub online: bool,
pub os: String,
}
#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
pub struct MeshStatus {
pub enabled: bool,
pub reachable: bool,
pub up: bool,
pub backend: Option<String>,
pub backend_state: String,
pub control_server: Option<String>,
pub magic_dns_name: Option<String>,
pub tailscale_ips: Vec<String>,
pub peers: Vec<MeshPeer>,
pub webhook_ingress_mode: Option<String>,
}
impl Default for MeshStatus {
fn default() -> Self {
Self {
enabled: false,
reachable: false,
up: false,
backend: None,
backend_state: "Stopped".to_owned(),
control_server: None,
magic_dns_name: None,
tailscale_ips: Vec::new(),
peers: Vec::new(),
webhook_ingress_mode: None,
}
}
}
const TAILSCALE_SAAS_CONTROL: &str = "controlplane.tailscale.com";
fn classify_backend(control_url: Option<&str>) -> Option<String> {
match control_url {
None => None,
Some(url) if url.contains(TAILSCALE_SAAS_CONTROL) => Some("tailscale".to_owned()),
Some(_) => Some("headscale".to_owned()),
}
}
pub fn parse_status_json(enabled: bool, raw: &serde_json::Value) -> MeshStatus {
let backend_state = raw
.get("BackendState")
.and_then(|v| v.as_str())
.unwrap_or("Stopped")
.to_owned();
let reachable = backend_state == "Running";
let control_server = raw
.get("Self")
.and_then(|s| s.get("ControlURL"))
.and_then(|v| v.as_str())
.or_else(|| raw.get("ControlURL").and_then(|v| v.as_str()))
.or_else(|| {
raw.get("CurrentTailnet")
.and_then(|t| t.get("Name"))
.and_then(|v| v.as_str())
})
.filter(|s| !s.is_empty())
.map(str::to_owned);
let backend = if backend_state == "Stopped" || backend_state == "NoState" {
None
} else {
classify_backend(control_server.as_deref())
};
let self_node = raw.get("Self");
let magic_dns_name = self_node
.and_then(|s| s.get("DNSName"))
.and_then(|v| v.as_str())
.map(|s| s.trim_end_matches('.').to_owned())
.filter(|s| !s.is_empty());
let tailscale_ips = self_node
.and_then(|s| s.get("TailscaleIPs"))
.and_then(|v| v.as_array())
.map(|arr| {
arr.iter()
.filter_map(|v| v.as_str().map(str::to_owned))
.collect()
})
.unwrap_or_default();
let peers = raw
.get("Peer")
.and_then(|v| v.as_object())
.map(|map| map.values().map(parse_peer).collect::<Vec<_>>())
.unwrap_or_default();
MeshStatus {
enabled,
reachable,
up: reachable,
backend,
backend_state,
control_server,
magic_dns_name,
tailscale_ips,
peers,
webhook_ingress_mode: None,
}
}
fn parse_peer(peer: &serde_json::Value) -> MeshPeer {
let dns = peer
.get("DNSName")
.and_then(|v| v.as_str())
.map(|s| s.trim_end_matches('.').to_owned())
.unwrap_or_default();
let host = peer
.get("HostName")
.and_then(|v| v.as_str())
.unwrap_or_default()
.to_owned();
let tailscale_ips: Vec<String> = peer
.get("TailscaleIPs")
.and_then(|v| v.as_array())
.map(|arr| {
arr.iter()
.filter_map(|v| v.as_str().map(str::to_owned))
.collect()
})
.unwrap_or_default();
let online = peer
.get("Online")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let os = peer
.get("OS")
.and_then(|v| v.as_str())
.unwrap_or_default()
.to_owned();
let host_or_dns = if !dns.is_empty() {
dns.clone()
} else if let Some(ip) = tailscale_ips.first() {
ip.clone()
} else {
host.clone()
};
let name = if !host.is_empty() {
host
} else {
dns.split('.').next().unwrap_or_default().to_owned()
};
MeshPeer {
name,
host_or_dns,
magic_dns_name: dns,
tailscale_ips,
online,
os,
}
}
pub async fn query_status() -> MeshStatus {
let enabled = is_enabled();
if !enabled {
return MeshStatus::default();
}
let Some(h) = host() else {
return MeshStatus {
enabled: true,
..Default::default()
};
};
match h.status_json().await {
Ok(raw) => parse_status_json(true, &raw),
Err(e) => {
tracing::debug!("mesh: status query failed: {e}");
MeshStatus {
enabled: true,
..Default::default()
}
}
}
}
pub async fn ensure_funnel(port: u16) -> anyhow::Result<String> {
if !is_enabled() {
anyhow::bail!("mesh disabled: set RYU_MESH_ENABLED to use Tailscale Funnel");
}
let h = host().ok_or_else(|| anyhow::anyhow!("mesh host not installed"))?;
h.ensure_funnel(port).await
}
pub async fn funnel_url(port: u16) -> Option<String> {
if !is_enabled() {
return None;
}
host()?.funnel_url(port).await
}
pub const BEARER_SOURCE_SHARED: &str = "shared-mesh-token";
pub const BEARER_SOURCE_NONE: &str = "none";
pub const BEARER_NONE_NOTE: &str =
"No usable RYU_TOKEN on this node. Provision every mesh node with the SAME strong \
RYU_TOKEN (the shared node-admittance secret) so a peer's require_auth accepts it; \
otherwise supply the target peer's own RYU_TOKEN when adding it.";
const DEFAULT_CORE_PORT: u16 = 7980;
fn peer_core_port() -> u16 {
std::env::var("RYU_MESH_PEER_PORT")
.ok()
.and_then(|v| v.trim().parse::<u16>().ok())
.unwrap_or(DEFAULT_CORE_PORT)
}
fn peer_url(peer: &MeshPeer, port: u16) -> String {
let host = if peer.magic_dns_name.is_empty() {
peer.host_or_dns.as_str()
} else {
peer.magic_dns_name.as_str()
};
format!("http://{host}:{port}")
}
pub fn resolve_mesh_bearer(node_token: Option<&str>) -> Option<String> {
let token = node_token?.trim();
if token.is_empty() || is_insecure_auth_token_placeholder(token) {
return None;
}
Some(token.to_owned())
}
#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
pub struct MeshPeerEntry {
pub name: String,
pub url: String,
pub magic_dns_name: String,
pub host_or_dns: String,
pub port: u16,
pub online: bool,
pub os: String,
pub bearer_available: bool,
pub bearer: Option<String>,
}
#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
pub struct MeshPeersResponse {
pub enabled: bool,
pub reachable: bool,
pub peers: Vec<MeshPeerEntry>,
pub bearer_source: String,
pub note: Option<String>,
}
pub fn build_peers_response(status: &MeshStatus, node_token: Option<&str>) -> MeshPeersResponse {
let bearer = resolve_mesh_bearer(node_token);
let bearer_available = bearer.is_some();
let port = peer_core_port();
let peers = status
.peers
.iter()
.map(|p| MeshPeerEntry {
name: p.name.clone(),
url: peer_url(p, port),
magic_dns_name: p.magic_dns_name.clone(),
host_or_dns: p.host_or_dns.clone(),
port,
online: p.online,
os: p.os.clone(),
bearer_available,
bearer: bearer.clone(),
})
.collect();
MeshPeersResponse {
enabled: status.enabled,
reachable: status.reachable,
peers,
bearer_source: if bearer_available {
BEARER_SOURCE_SHARED.to_owned()
} else {
BEARER_SOURCE_NONE.to_owned()
},
note: if bearer_available {
None
} else {
Some(BEARER_NONE_NOTE.to_owned())
},
}
}
#[cfg(test)]
mod tests {
use super::*;
static MESH_ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
fn lock_env() -> std::sync::MutexGuard<'static, ()> {
MESH_ENV_LOCK.lock().unwrap_or_else(|e| e.into_inner())
}
struct EnvGuard {
prev: Option<String>,
}
impl EnvGuard {
fn set(key: &str, val: &str) -> Self {
let prev = std::env::var(key).ok();
std::env::set_var(key, val);
Self { prev }
}
}
impl Drop for EnvGuard {
fn drop(&mut self) {
match &self.prev {
Some(v) => std::env::set_var("RYU_MESH_ENABLED", v),
None => std::env::remove_var("RYU_MESH_ENABLED"),
}
}
}
fn running_status_json() -> serde_json::Value {
serde_json::json!({
"BackendState": "Running",
"Self": {
"DNSName": "ryu-host.tailnet-x.ts.net.",
"TailscaleIPs": ["100.64.0.1", "fd7a:115c::1"],
"ControlURL": "https://controlplane.tailscale.com"
},
"Peer": {
"nodekey:abc": {
"HostName": "ryu-pi",
"DNSName": "ryu-pi.tailnet-x.ts.net.",
"TailscaleIPs": ["100.64.0.8"],
"Online": true,
"OS": "macOS"
}
}
})
}
#[test]
fn parse_status_json_running() {
let status = parse_status_json(true, &running_status_json());
assert!(status.enabled);
assert!(status.reachable);
assert!(status.up);
assert_eq!(status.reachable, status.up);
assert_eq!(status.backend.as_deref(), Some("tailscale"));
assert_eq!(status.backend_state, "Running");
assert_eq!(
status.magic_dns_name.as_deref(),
Some("ryu-host.tailnet-x.ts.net")
);
assert_eq!(status.tailscale_ips.len(), 2);
assert_eq!(status.peers.len(), 1);
let peer = &status.peers[0];
assert_eq!(peer.name, "ryu-pi");
assert_eq!(peer.host_or_dns, "ryu-pi.tailnet-x.ts.net");
assert_eq!(peer.magic_dns_name, "ryu-pi.tailnet-x.ts.net");
assert_eq!(peer.tailscale_ips, vec!["100.64.0.8".to_owned()]);
assert!(peer.online);
assert_eq!(peer.os, "macOS");
}
#[test]
fn parse_status_json_needs_login() {
let raw = serde_json::json!({ "BackendState": "NeedsLogin", "Self": {} });
let status = parse_status_json(true, &raw);
assert!(status.enabled);
assert!(!status.reachable);
assert!(!status.up);
assert_eq!(status.backend_state, "NeedsLogin");
assert!(status.backend.is_none());
assert!(status.peers.is_empty());
assert!(status.tailscale_ips.is_empty());
}
#[test]
fn parse_status_json_headscale_backend() {
let mut raw = running_status_json();
raw["Self"]["ControlURL"] = serde_json::json!("https://headscale.example.org");
let status = parse_status_json(true, &raw);
assert_eq!(status.backend.as_deref(), Some("headscale"));
assert_eq!(
status.control_server.as_deref(),
Some("https://headscale.example.org")
);
}
#[test]
fn disabled_shape_is_all_default() {
let status = MeshStatus::default();
assert!(!status.enabled);
assert!(!status.reachable);
assert!(!status.up);
assert!(status.backend.is_none());
assert_eq!(status.backend_state, "Stopped");
assert!(status.control_server.is_none());
assert!(status.magic_dns_name.is_none());
assert!(status.tailscale_ips.is_empty());
assert!(status.peers.is_empty());
assert!(status.webhook_ingress_mode.is_none());
}
#[test]
fn disabled_shape_serializes_to_contract6() {
let json = serde_json::to_value(MeshStatus::default()).unwrap();
assert_eq!(json["enabled"], serde_json::json!(false));
assert_eq!(json["reachable"], serde_json::json!(false));
assert_eq!(json["up"], serde_json::json!(false));
assert_eq!(json["backend"], serde_json::Value::Null);
assert_eq!(json["backend_state"], serde_json::json!("Stopped"));
assert_eq!(json["control_server"], serde_json::Value::Null);
assert_eq!(json["magic_dns_name"], serde_json::Value::Null);
assert_eq!(json["tailscale_ips"], serde_json::json!([]));
assert_eq!(json["peers"], serde_json::json!([]));
assert_eq!(json["webhook_ingress_mode"], serde_json::Value::Null);
}
#[test]
fn is_enabled_default_off() {
if std::env::var("RYU_MESH_ENABLED").is_err() {
assert!(!is_enabled());
}
}
struct PrefGuard {
prev: bool,
}
impl PrefGuard {
fn set(v: bool) -> Self {
let prev = MESH_PREF_ENABLED.load(Ordering::Relaxed);
set_pref_enabled(v);
Self { prev }
}
}
impl Drop for PrefGuard {
fn drop(&mut self) {
set_pref_enabled(self.prev);
}
}
#[test]
fn pref_enable_drives_is_enabled_when_env_unset() {
let _lock = lock_env();
if std::env::var("RYU_MESH_ENABLED").is_err() {
let _p = PrefGuard::set(true);
assert!(is_enabled());
set_pref_enabled(false);
assert!(!is_enabled());
}
}
#[test]
fn env_wins_over_pref() {
let _lock = lock_env();
let _p = PrefGuard::set(true);
let _e = EnvGuard::set("RYU_MESH_ENABLED", "0");
assert!(!is_enabled());
let _e = EnvGuard::set("RYU_MESH_ENABLED", "1");
assert!(is_enabled());
}
#[test]
fn parse_enabled_matches_env_truthiness() {
assert!(!parse_enabled(None));
assert!(!parse_enabled(Some("")));
assert!(!parse_enabled(Some("0")));
assert!(!parse_enabled(Some("false")));
assert!(!parse_enabled(Some("FALSE")));
assert!(!parse_enabled(Some("no")));
assert!(parse_enabled(Some("1")));
assert!(parse_enabled(Some("true")));
assert!(parse_enabled(Some("yes")));
assert!(parse_enabled(Some(" 1 ")));
}
#[test]
fn peer_host_or_dns_falls_back_to_ip() {
let peer = serde_json::json!({
"HostName": "",
"DNSName": "",
"TailscaleIPs": ["100.64.0.9"],
"Online": false,
"OS": "linux"
});
let parsed = parse_peer(&peer);
assert_eq!(parsed.host_or_dns, "100.64.0.9");
assert!(!parsed.online);
}
#[test]
fn resolve_mesh_bearer_returns_real_token() {
assert_eq!(
resolve_mesh_bearer(Some("ryu_shared_secret")).as_deref(),
Some("ryu_shared_secret")
);
}
#[test]
fn resolve_mesh_bearer_is_fail_closed_without_a_real_token() {
assert!(resolve_mesh_bearer(None).is_none());
assert!(resolve_mesh_bearer(Some("")).is_none());
assert!(resolve_mesh_bearer(Some(" ")).is_none());
assert!(resolve_mesh_bearer(Some("CHANGE_ME")).is_none());
assert!(resolve_mesh_bearer(Some("change_me")).is_none());
assert!(resolve_mesh_bearer(Some("REPLACE_ME")).is_none());
assert!(resolve_mesh_bearer(Some("SECRET")).is_none());
}
#[test]
fn placeholder_predicate_matches_known_weak_tokens() {
assert!(is_insecure_auth_token_placeholder("CHANGE_ME"));
assert!(is_insecure_auth_token_placeholder(" changeme "));
assert!(is_insecure_auth_token_placeholder("PASSWORD"));
assert!(!is_insecure_auth_token_placeholder("ryu_strong_random"));
assert!(!is_insecure_auth_token_placeholder(""));
}
#[test]
fn peers_response_carries_shared_bearer_and_urls() {
let status = parse_status_json(true, &running_status_json());
let resp = build_peers_response(&status, Some("ryu_shared_secret"));
assert!(resp.enabled);
assert_eq!(resp.bearer_source, BEARER_SOURCE_SHARED);
assert!(resp.note.is_none());
assert_eq!(resp.peers.len(), 1);
let peer = &resp.peers[0];
assert_eq!(peer.name, "ryu-pi");
assert_eq!(peer.url, "http://ryu-pi.tailnet-x.ts.net:7980");
assert_eq!(peer.port, 7980);
assert!(peer.bearer_available);
assert_eq!(peer.bearer.as_deref(), Some("ryu_shared_secret"));
}
#[test]
fn peers_response_without_token_is_honest_and_documents_secret() {
let status = parse_status_json(true, &running_status_json());
let resp = build_peers_response(&status, None);
assert_eq!(resp.bearer_source, BEARER_SOURCE_NONE);
assert_eq!(resp.note.as_deref(), Some(BEARER_NONE_NOTE));
let peer = &resp.peers[0];
assert!(!peer.bearer_available);
assert!(peer.bearer.is_none());
assert_eq!(peer.url, "http://ryu-pi.tailnet-x.ts.net:7980");
}
#[test]
fn disabled_mesh_yields_empty_peers() {
let resp = build_peers_response(&MeshStatus::default(), Some("ryu_shared_secret"));
assert!(!resp.enabled);
assert!(resp.peers.is_empty());
assert_eq!(resp.bearer_source, BEARER_SOURCE_SHARED);
}
#[tokio::test]
async fn disabled_query_status_never_touches_host() {
let _lock = lock_env();
let _p = PrefGuard::set(false);
if std::env::var("RYU_MESH_ENABLED").is_err() {
let status = query_status().await;
assert_eq!(status, MeshStatus::default());
assert!(ensure_funnel(443).await.is_err());
assert!(funnel_url(443).await.is_none());
}
}
}