`status` asked a running agent and failed without one; `doctor` checked
the host and ignored the agent. Between them they answered one question
in two halves. They are now one command: device, agent, networks, host
capability, local addresses, all graded and aligned the same way.
`id` was worse than either. It spawned a whole agent to print three
facts, which took the directory lock and so failed with "owned by
another running agent instance" exactly when the answer was most wanted.
The lock exists to keep one writer over the mandatory state; reading who
this device is needs no such thing.
So both commands now prefer the running agent, which is live and
authoritative, and fall back to the state store, which takes no lock.
StateStore gains device_identity(), which reads and never writes:
load_or_create had the side effect of deciding an identity as a
consequence of asking about one.
Presentation moved out of the library. StatusReport::render is gone and
the CLI renders the structured data, so there is one renderer rather than
two that would drift. Health gains an Info level for rows that are facts
rather than checks — an endpoint id is neither good nor bad, and a column
of green next to plain data teaches the eye to ignore the column. A
report with no checks in it now ends without a summary instead of
claiming that everything checked out.
PathAddr and TransportKind grew Display impls; both were reaching the
user through {:?}, which is how "Direct via Ip(88.198.17.44:49792)" got
printed. The transport grading compares without regard to case, because
the agent answering can be a different build from the client asking and
this field's spelling has now changed once.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
277 lines
10 KiB
Rust
277 lines
10 KiB
Rust
//! The local control interface: a client asking a running agent for status.
|
|
//!
|
|
//! Uses a real Unix socket on a temporary path, the real agent and the real
|
|
//! WireGuard data plane, so what a `tsunagi status` client would see is what
|
|
//! is checked here.
|
|
|
|
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
|
|
|
|
mod common;
|
|
|
|
use std::sync::Arc;
|
|
use std::time::Duration;
|
|
|
|
use common::{config_with, network, wait_for_peers, wait_until};
|
|
use tempfile::TempDir;
|
|
use tsunagi::dataplane::IpPlugin;
|
|
use tsunagi::dataplane::wireguard::{MemoryTunFactory, WireguardConfig, WireguardPlugin};
|
|
use tsunagi::discovery::SharedMemoryDiscovery;
|
|
use tsunagi::ipc::unix::{ControlSocket, request_status};
|
|
use tsunagi::ipc::{StatusReport, control_socket_path};
|
|
use tsunagi::{Agent, BoxFuture};
|
|
|
|
/// Builds the report the way the binary does, from the agent plus the plugin.
|
|
fn source(agent: Agent, plugin: Arc<WireguardPlugin>) -> Arc<dyn tsunagi::ipc::unix::ReportSource> {
|
|
Arc::new(move || -> BoxFuture<'static, StatusReport> {
|
|
let agent = agent.clone();
|
|
let plugin = Arc::clone(&plugin);
|
|
Box::pin(async move {
|
|
let status = agent.status().await.unwrap();
|
|
let networks = status
|
|
.networks
|
|
.iter()
|
|
.map(|net| tsunagi::ipc::NetworkReport {
|
|
name: net.name.to_string(),
|
|
network_id: net.network_id.to_string(),
|
|
active: true,
|
|
peers: net
|
|
.peers
|
|
.iter()
|
|
.map(|peer| tsunagi::ipc::PeerReport {
|
|
endpoint_id: peer.endpoint_id.to_string(),
|
|
hostname: peer.hostname.clone(),
|
|
transport: format!("{:?}", peer.transport),
|
|
rtt_ms: peer.rtt.map(|rtt| rtt.as_millis() as u64),
|
|
})
|
|
.collect(),
|
|
overlay: plugin.overview(net.network_id).map(|view| {
|
|
tsunagi::ipc::OverlayReport {
|
|
interface: view.interface.clone(),
|
|
mtu: view.mtu,
|
|
address: view.overlay_address.to_string(),
|
|
prefix: view.overlay_prefix.to_string(),
|
|
prefix_len: view.overlay_prefix_len,
|
|
peers: view
|
|
.peers
|
|
.iter()
|
|
.map(|peer| tsunagi::ipc::OverlayPeerReport {
|
|
public_key: peer.public_key.to_string(),
|
|
address: peer.overlay_address.to_string(),
|
|
handshake_secs_ago: peer
|
|
.tunnel
|
|
.as_ref()
|
|
.and_then(|t| t.health.since_handshake)
|
|
.map(|since| since.as_secs()),
|
|
..Default::default()
|
|
})
|
|
.collect(),
|
|
..Default::default()
|
|
}
|
|
}),
|
|
..Default::default()
|
|
})
|
|
.collect();
|
|
StatusReport {
|
|
endpoint_id: status.endpoint_id.to_string(),
|
|
hostname: status.hostname.clone(),
|
|
bound_sockets: status
|
|
.bound_sockets
|
|
.iter()
|
|
.map(ToString::to_string)
|
|
.collect(),
|
|
cache_healthy: status.cache_healthy,
|
|
networks,
|
|
}
|
|
})
|
|
})
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn a_client_sees_the_agent_and_its_overlay() {
|
|
let discovery = SharedMemoryDiscovery::new();
|
|
let (name, secret) = network("control-socket");
|
|
|
|
let dir_a = TempDir::new().unwrap();
|
|
let tuns = MemoryTunFactory::new();
|
|
let plugin = WireguardPlugin::open(
|
|
WireguardConfig::new(dir_a.path().join("wg"))
|
|
.with_interface_prefix("tca")
|
|
.with_reconcile(Duration::from_millis(20), Duration::from_millis(250)),
|
|
Arc::new(tuns),
|
|
)
|
|
.await
|
|
.unwrap();
|
|
let agent = Agent::spawn(
|
|
config_with(dir_a.path(), &discovery).with_plugin(plugin.clone() as Arc<dyn IpPlugin>),
|
|
)
|
|
.await
|
|
.unwrap();
|
|
|
|
let dir_b = TempDir::new().unwrap();
|
|
let plugin_b = WireguardPlugin::open(
|
|
WireguardConfig::new(dir_b.path().join("wg"))
|
|
.with_interface_prefix("tcb")
|
|
.with_reconcile(Duration::from_millis(20), Duration::from_millis(250)),
|
|
Arc::new(MemoryTunFactory::new()),
|
|
)
|
|
.await
|
|
.unwrap();
|
|
let agent_b = Agent::spawn(
|
|
config_with(dir_b.path(), &discovery).with_plugin(plugin_b.clone() as Arc<dyn IpPlugin>),
|
|
)
|
|
.await
|
|
.unwrap();
|
|
|
|
let network_id = agent.join_network(&name, &secret).await.unwrap();
|
|
agent_b.join_network(&name, &secret).await.unwrap();
|
|
wait_for_peers(&agent, network_id, 1).await;
|
|
|
|
// A short path: a Unix socket address is limited to about 100 bytes.
|
|
let socket_path = dir_a.path().join("agent.sock");
|
|
let control = ControlSocket::bind(&socket_path, source(agent.clone(), plugin.clone()))
|
|
.await
|
|
.unwrap();
|
|
|
|
let report = wait_until("the overlay is reported as up", || {
|
|
let socket_path = socket_path.clone();
|
|
async move {
|
|
let report = request_status(&socket_path).await.ok()?;
|
|
let overlay = report.networks.first()?.overlay.as_ref()?;
|
|
overlay
|
|
.peers
|
|
.iter()
|
|
.any(|peer| peer.is_up())
|
|
.then_some(report)
|
|
}
|
|
})
|
|
.await;
|
|
|
|
assert_eq!(report.endpoint_id, agent.endpoint_id().to_string());
|
|
assert_eq!(report.networks.len(), 1);
|
|
let net = &report.networks[0];
|
|
assert_eq!(net.network_id, network_id.to_string());
|
|
assert_eq!(net.peers.len(), 1);
|
|
assert_eq!(net.peers[0].endpoint_id, agent_b.endpoint_id().to_string());
|
|
|
|
let overlay = net.overlay.as_ref().unwrap();
|
|
assert!(overlay.interface.starts_with("tca"));
|
|
assert_eq!(overlay.mtu, 1280);
|
|
assert_eq!(overlay.peers.len(), 1);
|
|
|
|
// The report carries what a reader needs, in structured form: how it is
|
|
// laid out is the CLI's business and is tested there.
|
|
assert_eq!(report.endpoint_id, agent.endpoint_id().to_string());
|
|
assert_eq!(
|
|
overlay.peers.iter().filter(|peer| peer.is_up()).count(),
|
|
1,
|
|
"the tunnel is up: {:?}",
|
|
overlay.peers
|
|
);
|
|
assert!(!overlay.address.is_empty());
|
|
|
|
control.shutdown().await;
|
|
assert!(!socket_path.exists(), "the socket is removed on shutdown");
|
|
|
|
// With nothing listening, a client gets an error rather than hanging.
|
|
assert!(request_status(&socket_path).await.is_err());
|
|
|
|
agent.shutdown().await;
|
|
agent_b.shutdown().await;
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn a_leftover_socket_file_is_replaced_but_a_live_one_is_not() {
|
|
let dir = TempDir::new().unwrap();
|
|
let path = dir.path().join("agent.sock");
|
|
|
|
let empty: Arc<dyn tsunagi::ipc::unix::ReportSource> =
|
|
Arc::new(|| -> BoxFuture<'static, StatusReport> {
|
|
Box::pin(async { StatusReport::default() })
|
|
});
|
|
|
|
// A file with nobody listening is a leftover from a crash.
|
|
std::fs::write(&path, b"stale").unwrap();
|
|
let first = ControlSocket::bind(&path, Arc::clone(&empty))
|
|
.await
|
|
.unwrap();
|
|
assert!(request_status(&path).await.is_ok());
|
|
|
|
// A live socket is not stolen from the agent that owns it.
|
|
let second = ControlSocket::bind(&path, Arc::clone(&empty)).await;
|
|
assert!(
|
|
matches!(second, Err(tsunagi::Error::StateLocked { .. })),
|
|
"a second agent must not take over a live control socket"
|
|
);
|
|
|
|
first.shutdown().await;
|
|
}
|
|
|
|
#[test]
|
|
fn the_socket_path_is_derived_and_short_enough() {
|
|
let deep = std::path::PathBuf::from(
|
|
"/home/someone/.local/share/with/a/very/deeply/nested/directory/that/goes/on/and/on/and/on/tsunagi/state",
|
|
);
|
|
let path = control_socket_path(&deep);
|
|
|
|
// A Unix socket address is limited to roughly 100 bytes, so a deep state
|
|
// directory must not produce a path that cannot be bound.
|
|
if std::env::var_os("XDG_RUNTIME_DIR").is_some() {
|
|
assert!(
|
|
path.as_os_str().len() < 100,
|
|
"derived path is {} bytes: {}",
|
|
path.as_os_str().len(),
|
|
path.display()
|
|
);
|
|
}
|
|
|
|
// Deterministic, and different state directories never share a socket.
|
|
assert_eq!(path, control_socket_path(&deep));
|
|
assert_ne!(
|
|
path,
|
|
control_socket_path(&std::path::PathBuf::from("/somewhere/else"))
|
|
);
|
|
}
|
|
|
|
/// Asking who this device is must not need the directory lock.
|
|
///
|
|
/// The lock belongs to the one agent allowed to *write* the state. `id` only
|
|
/// reads, so making it take the lock would mean the question could never be
|
|
/// answered while an agent was running — which is exactly when you want to
|
|
/// ask it.
|
|
#[cfg(feature = "cli")]
|
|
#[tokio::test]
|
|
async fn identity_can_be_read_while_an_agent_holds_the_directory() {
|
|
let dir = TempDir::new().unwrap();
|
|
let discovery = SharedMemoryDiscovery::new();
|
|
|
|
// Hold the directory the way a running agent does.
|
|
let agent = Agent::spawn(config_with(dir.path(), &discovery))
|
|
.await
|
|
.unwrap();
|
|
let expected = agent.endpoint_id().to_string();
|
|
|
|
let output = tokio::process::Command::new(env!("CARGO_BIN_EXE_tsunagi"))
|
|
.arg("id")
|
|
// The same layout `StoragePaths::under` gives the agent above.
|
|
.arg("--state-dir")
|
|
.arg(dir.path().join("state"))
|
|
.arg("--cache-dir")
|
|
.arg(dir.path().join("cache"))
|
|
.output()
|
|
.await
|
|
.unwrap();
|
|
|
|
let stdout = String::from_utf8_lossy(&output.stdout);
|
|
let stderr = String::from_utf8_lossy(&output.stderr);
|
|
assert!(
|
|
output.status.success(),
|
|
"`id` failed while an agent was running:\n{stderr}"
|
|
);
|
|
assert!(
|
|
stdout.contains(&expected),
|
|
"expected {expected} in:\n{stdout}"
|
|
);
|
|
|
|
agent.shutdown().await;
|
|
}
|