Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -39,12 +39,15 @@ monitor-agent --server https://your-hub --token <token>

`src/collect.rs` 中的 `Facts` 与 `Metrics` 两个 struct 直接序列化为线上 JSON,是字段的权威定义。

- **`Facts`** 连接时上报一次:主机名、系统、内核、架构、虚拟化类型、CPU 型号与核数、内存与磁盘总量、本机 IPv4 / IPv6
- **`Facts`** 连接时上报一次:主机名、系统、内核、架构、虚拟化类型、CPU 型号与核数、内存与磁盘总量、本机 IPv4 / IPv6(每族一个,公网地址优先;IPv6 不取临时地址和已废弃地址)
- **`Metrics`** 每 `--interval` 秒上报:CPU、负载、内存、swap、磁盘、网卡收发速率与内核累计计数器、TCP / UDP 连接数、进程数、运行时间

`net_rx_total` / `net_tx_total` 为内核 lifetime 计数器,原样上报;`boot_id` 取自
`/proc/sys/kernel/random/boot_id`,是 hub 判定主机重启的唯一依据,**不要删**。

连接 hub 时逐个尝试解析出的地址,除最后一个外每个限 5 秒。网卡上只有内网 IPv4(NAT)时先连 hub 的
IPv4:NAT 的公网地址不在网卡上,hub 只有看到一条 IPv4 连接才知道它。

协议说明见 [hub 仓库](https://github.com/monitor-probe/monitor)。

## 构建
Expand Down
159 changes: 144 additions & 15 deletions src/collect.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@

use std::collections::HashMap;
use std::fs;
use std::net::{IpAddr, Ipv6Addr};
use std::time::Instant;

use serde::Serialize;
Expand Down Expand Up @@ -93,8 +94,8 @@ pub struct Facts {
pub swap_total: u64,
pub disk_total: u64,
pub agent_version: String,
/// The host's own addresses. The hub sees only the family the agent
/// connected over, which on a dual-stack host is usually v6.
/// The host's own addresses, public ones first; see [`pick`]. The hub sees
/// only the family the agent connected over.
pub ipv4: String,
pub ipv6: String,
}
Expand Down Expand Up @@ -311,27 +312,81 @@ fn uptime() -> u64 {
.unwrap_or(0.0) as u64
}

/// The first real address of each family the kernel reports. On a VPS these
/// are the public ones; behind NAT the v4 is private, which is what the machine
/// actually holds -- no external service is consulted.
/// One address of each family the machine holds. Behind NAT the v4 is private,
/// which is what the machine actually holds -- no external service is
/// consulted.
///
/// Filtered by [`SKIP_IFACES`] alone, so a docker bridge cannot pass for the
/// machine's address. [`is_stacked`] is not applied here: it answers whether
/// bytes were already counted lower down, and a bridge holding the host address
/// is both stacked and this machine.
fn addresses() -> (String, String) {
let (mut v4, mut v6) = (String::new(), String::new());
for iface in if_addrs::get_if_addrs().unwrap_or_default() {
if skip_iface(&iface.name) || iface.is_link_local() || !iface.is_oper_up() {
continue;
}
match iface.ip() {
std::net::IpAddr::V4(ip) if v4.is_empty() => v4 = ip.to_string(),
std::net::IpAddr::V6(ip) if v6.is_empty() => v6 = ip.to_string(),
_ => {}
let transient = transient_v6(&fs::read_to_string("/proc/net/if_inet6").unwrap_or_default());
let held: Vec<IpAddr> = if_addrs::get_if_addrs()
.unwrap_or_default()
.into_iter()
.filter(|i| !skip_iface(&i.name) && !i.is_link_local() && i.is_oper_up())
.map(|i| i.ip())
.collect();
pick(&held, &transient)
}

/// A public address before any other, then a stable v6 before a transient
/// one; ties keep the kernel's order. Taking the first address instead would
/// report a ULA or a proxy's TUN address whenever its interface is listed
/// ahead of the one holding the public address: an LXC guest with a ULA on
/// eth0 and its public /128 on eth1 would report the ULA.
fn pick(held: &[IpAddr], transient: &[Ipv6Addr]) -> (String, String) {
let best = |v6: bool| {
held.iter()
.filter(|ip| ip.is_ipv6() == v6)
.min_by_key(|ip| (!is_public(**ip), matches!(ip, IpAddr::V6(a) if transient.contains(a))))
.map_or_else(String::new, ToString::to_string)
};
(best(false), best(true))
}

/// IPv6 addresses held but not worth reporting: temporary (privacy extensions,
/// replaced daily), deprecated (past their preferred lifetime, as an old
/// prefix is after a home line redials), tentative, or failed duplicate
/// detection. Read from /proc/net/if_inet6, whose fields are address, ifindex,
/// prefix length, scope, flags and name; if_addrs does not expose the flags.
fn transient_v6(text: &str) -> Vec<Ipv6Addr> {
// IFA_F_TEMPORARY | IFA_F_DADFAILED | IFA_F_DEPRECATED | IFA_F_TENTATIVE
const TRANSIENT: u8 = 0x01 | 0x08 | 0x20 | 0x40;
text.lines()
.filter_map(|line| {
let mut f = line.split_whitespace();
let addr = u128::from_str_radix(f.next()?, 16).ok()?;
let flags = u8::from_str_radix(f.nth(3)?, 16).ok()?;
(flags & TRANSIENT != 0).then(|| Ipv6Addr::from(addr))
})
.collect()
}

/// Globally routable. Excluded on the v4 side: RFC 1918, CGNAT (100.64/10),
/// loopback, link-local, 0/8, 192.0.0/24 (where 464XLAT places its CLAT),
/// 198.18/15 (the fake-IP range TUN-mode proxies such as Clash assign to
/// themselves), multicast and reserved. On the v6 side only 2000::/3 counts,
/// which leaves out ULA (fc00::/7), link-local and loopback.
///
/// The hub and its panel apply the same ranges to the addresses this agent
/// reports; the three lists are to be changed together.
pub fn is_public(ip: IpAddr) -> bool {
match ip {
IpAddr::V4(v4) => {
let [a, b, c, _] = v4.octets();
!(v4.is_private()
|| v4.is_loopback()
|| v4.is_link_local()
|| a == 0
|| a >= 224
|| (a == 100 && b & 0xc0 == 64)
|| (a == 192 && b == 0 && c == 0)
|| (a == 198 && b & 0xfe == 18))
}
IpAddr::V6(v6) => v6.segments()[0] & 0xe000 == 0x2000,
}
(v4, v6)
}

/// Sums the kernel's lifetime byte counters, one count per byte on the wire.
Expand Down Expand Up @@ -675,6 +730,80 @@ mod tests {
assert!(!skip_iface("eth0") && !is_stacked("eth0"), "the wire itself is what gets counted");
}

#[test]
fn the_reported_address_is_the_public_one_whatever_the_kernel_lists_first() {
let ip = |s: &str| s.parse::<IpAddr>().unwrap();
let picked = |held: &[&str], transient: &[Ipv6Addr]| {
let held: Vec<IpAddr> = held.iter().map(|s| ip(s)).collect();
pick(&held, transient)
};
let pair = |v4: &str, v6: &str| (v4.to_owned(), v6.to_owned());

// An LXC NAT guest: private v4 and a ULA on eth0, its public /128 on eth1.
assert_eq!(
picked(&["10.10.1.5", "fd42:43af:6613:5936::1", "2401:b60:1c::5"], &[]),
pair("10.10.1.5", "2401:b60:1c::5")
);
// A TUN-mode proxy, a CGNAT overlay and the LAN all listed before the
// public address.
assert_eq!(
picked(&["198.18.0.1", "100.64.0.9", "192.168.1.5", "203.0.113.7"], &[]),
pair("203.0.113.7", "")
);
// With nothing public the kernel's order stands.
assert_eq!(picked(&["192.168.1.5", "172.19.0.1"], &[]), pair("192.168.1.5", ""));

// SLAAC with privacy extensions lists the temporary address first.
let inet6 = "24098a1e3b4179f0a1b2c3d4e5f60718 02 40 00 01 eth0\n\
24098a1e3b4179f00211223344556677 02 40 00 00 eth0\n\
24098a1e3b4100000211223344556677 02 40 00 20 eth0\n\
fe80000000000000021122fffe334455 02 40 20 80 eth0\n";
let transient = transient_v6(inet6);
let v6 = |s: &str| s.parse::<Ipv6Addr>().unwrap();
assert_eq!(
transient,
[v6("2409:8a1e:3b41:79f0:a1b2:c3d4:e5f6:718"), v6("2409:8a1e:3b41::211:2233:4455:6677")]
);
let slaac = [
"2409:8a1e:3b41:79f0:a1b2:c3d4:e5f6:718",
"2409:8a1e:3b41::211:2233:4455:6677",
"2409:8a1e:3b41:79f0:211:2233:4455:6677",
];
assert_eq!(picked(&slaac, &transient), pair("", "2409:8a1e:3b41:79f0:211:2233:4455:6677"));
// A transient address is still better than none.
assert_eq!(picked(&slaac[..1], &transient), pair("", slaac[0]));
}

#[test]
fn only_globally_routable_addresses_count_as_public() {
let ip = |s: &str| s.parse::<IpAddr>().unwrap();
for s in [
"10.0.0.1",
"172.16.0.1",
"192.168.1.1",
"100.64.0.1",
"100.127.255.1",
"127.0.0.1",
"169.254.1.1",
"0.0.0.1",
"192.0.0.4",
"198.18.0.1",
"198.19.255.1",
"224.0.0.1",
"fd42::1",
"fc00::1",
"fe80::1",
"::1",
] {
assert!(!is_public(ip(s)), "{s}");
}
for s in
["1.1.1.1", "100.128.0.1", "198.20.0.1", "192.0.1.1", "223.5.5.5", "2401:b60:1c::5", "3fff::1"]
{
assert!(is_public(ip(s)), "{s}");
}
}

#[test]
fn net_rate_is_zero_on_first_sample_and_after_a_reboot() {
let mut c = Collector::new();
Expand Down
132 changes: 124 additions & 8 deletions src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -211,18 +211,31 @@ fn reconnect_wait(previous: u64, lasted: Duration) -> u64 {
}
}

/// Deadline covering all three stages of establishing a connection.
/// Deadline covering every stage of establishing a connection: resolution,
/// the TCP handshake, the TLS exchange and the HTTP upgrade.
///
/// Only the TCP handshake has a deadline of its own; the TLS exchange and the
/// HTTP upgrade have none, so a peer that accepts and then goes silent would
/// leave `connect_async` pending indefinitely, and the agent running without
/// leave the handshake pending indefinitely, and the agent running without
/// reporting or logging.
///
/// Deliberately generous: a healthy connect takes a quarter of a second, the
/// slowest measured sixty. This is not a latency budget but the point past
/// which nothing is expected to arrive.
const CONNECT_DEADLINE: Duration = Duration::from_secs(120);

/// How long one of the hub's addresses may take to accept a TCP connection
/// while another remains to be tried: the first SYN and its retransmits at one
/// and three seconds. The last address has no limit of its own and is bounded
/// by [`CONNECT_DEADLINE`] alone, so a hub with a single slow address is
/// reached as before.
///
/// Without it a black-holed address -- typically an AAAA record over a v6
/// route that leads nowhere -- holds the connect for the kernel's 127 seconds
/// of SYN retries, beyond CONNECT_DEADLINE, and the address behind it is never
/// tried.
const DIAL_FALLBACK: Duration = Duration::from_secs(5);

/// The hub sends one kind of message, a probe list a few hundred bytes long.
/// Tungstenite's 64 MiB default would hand the peer this process's entire
/// memory budget.
Expand Down Expand Up @@ -253,18 +266,37 @@ async fn session(
.insert("authorization", format!("Bearer {token}").parse().context("token is not header-safe")?);
let config =
WebSocketConfig::default().max_message_size(Some(MAX_MESSAGE)).max_frame_size(Some(MAX_MESSAGE));
let connect = tokio_tungstenite::connect_async_with_config(request, Some(config), false);
let (mut ws, _) = tokio::time::timeout(CONNECT_DEADLINE, connect)
let uri = request.uri();
// Brackets off an IPv6 literal, which `lookup_host` parses bare.
let host = uri
.host()
.context("server URL has no host")?
.trim_start_matches('[')
.trim_end_matches(']')
.to_owned();
let port = uri.port_u16().unwrap_or(if uri.scheme_str() == Some("wss") { 443 } else { 80 });
// Collected before dialing, as the v4 it reports decides which family is
// tried first.
let facts = collector.facts();
let behind_nat = facts.ipv4.parse().is_ok_and(|ip| !collect::is_public(ip));
let connect = async {
let stream = dial(&host, port, behind_nat).await?;
let peer = stream.peer_addr()?;
let (ws, _) = tokio_tungstenite::client_async_tls_with_config(request, stream, Some(config), None)
.await
.context("handshake")?;
anyhow::Ok((ws, peer))
};
let (mut ws, peer) = tokio::time::timeout(CONNECT_DEADLINE, connect)
.await
.with_context(|| format!("no connection after {}s", CONNECT_DEADLINE.as_secs()))?
.context("connect")?;
eprintln!("connected");
.with_context(|| format!("no connection after {}s", CONNECT_DEADLINE.as_secs()))??;
eprintln!("connected to {peer}");
*connected = Some(Instant::now());
// The clock starts at the handshake and the hello below draws from it like
// every other write, so no two writes can each claim a full HUB_SILENCE.
let mut last_frame = Instant::now();

send(&mut ws, notify("hello", serde_json::to_value(collector.facts())?), remaining(last_frame)).await?;
send(&mut ws, notify("hello", serde_json::to_value(facts)?), remaining(last_frame)).await?;

let (result_tx, mut result_rx) = mpsc::channel::<Message>(64);
let mut ping_tasks: Vec<(PingTask, tokio::task::JoinHandle<()>)> = Vec::new();
Expand Down Expand Up @@ -316,6 +348,47 @@ async fn session(
result
}

/// Opens the TCP connection to the hub, trying its addresses in turn.
///
/// A host whose IPv4 is private tries the hub's IPv4 addresses first. Its
/// public IPv4 exists only on the NAT in front of it, and the hub, which knows
/// no more than the address a connection arrives from, learns it only from a
/// connection made over v4. A public v6 needs no such route: it sits on the
/// interface and travels in the hello. Every other host keeps the resolver's
/// order.
async fn dial(host: &str, port: u16, prefer_v4: bool) -> Result<TcpStream> {
let mut addrs: Vec<std::net::SocketAddr> =
tokio::net::lookup_host((host, port)).await.with_context(|| format!("resolve {host}"))?.collect();
if prefer_v4 {
addrs.sort_by_key(|a| !a.is_ipv4());
}
connect_first(&addrs).await.with_context(|| format!("connect {host}"))
}

/// The first of `addrs` to accept, each but the last given [`DIAL_FALLBACK`].
/// Every failure is kept, so the log names which family failed and how.
async fn connect_first(addrs: &[std::net::SocketAddr]) -> Result<TcpStream> {
let mut failures = Vec::new();
for (i, addr) in addrs.iter().enumerate() {
let attempt = TcpStream::connect(addr);
let result = if i + 1 < addrs.len() {
tokio::time::timeout(DIAL_FALLBACK, attempt)
.await
.unwrap_or_else(|_| Err(std::io::ErrorKind::TimedOut.into()))
} else {
attempt.await
};
match result {
Ok(stream) => return Ok(stream),
Err(e) => failures.push(format!("{addr}: {e}")),
}
}
if failures.is_empty() {
bail!("no address");
}
bail!("{}", failures.join("; "))
}

/// Ceiling on concurrent probe loops.
///
/// A task serialises to about forty bytes, so one [`MAX_MESSAGE`] frame could
Expand Down Expand Up @@ -550,6 +623,49 @@ mod tests {
);
}

/// A listener whose accept queue is full drops further SYNs rather than
/// refusing them: a black hole on loopback. The clock is paused, so the
/// fallback deadline elapses as soon as nothing else can progress, while
/// without it the connect would wait out the kernel's SYN retries.
#[tokio::test(start_paused = true)]
async fn a_black_holed_address_gives_way_to_the_next_within_the_fallback() {
let hole = tokio::net::TcpSocket::new_v4().unwrap();
hole.bind("127.0.0.1:0".parse().unwrap()).unwrap();
let hole = hole.listen(0).unwrap();
let dead = hole.local_addr().unwrap();
let _queued = std::net::TcpStream::connect(dead).unwrap();
let live = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let live = live.local_addr().unwrap();

let started = Instant::now();
let stream = connect_first(&[dead, live]).await.unwrap();
assert_eq!(stream.peer_addr().unwrap(), live);
assert_eq!(started.elapsed(), DIAL_FALLBACK, "the dead address costs the fallback and no more");
let e = connect_first(&[dead, "127.0.0.1:1".parse().unwrap()]).await.unwrap_err().to_string();
assert!(e.contains("timed out") && e.contains("refused"), "each failure is named: {e}");
}

/// Where `localhost` resolves to ::1 before 127.0.0.1, only the preference
/// can land this connect on v4. Elsewhere -- no v6 loopback, or a resolver
/// ordering v4 first -- the preference is not observable and the test says
/// so rather than failing.
#[tokio::test]
async fn a_host_behind_nat_dials_the_hubs_v4_first() {
let v4 = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let port = v4.local_addr().unwrap().port();
let Ok(_v6) = std::net::TcpListener::bind(("::1", port)) else {
return eprintln!("skipped: no IPv6 loopback");
};
if !dial("localhost", port, false).await.unwrap().peer_addr().unwrap().is_ipv6() {
return eprintln!("skipped: the resolver lists 127.0.0.1 first");
}
assert!(dial("localhost", port, true).await.unwrap().peer_addr().unwrap().is_ipv4());
assert!(
dial("::1", port, true).await.unwrap().peer_addr().unwrap().is_ipv6(),
"a literal is dialed as given"
);
}

/// The deadline must stay under the kernel's first SYN retransmit, or a
/// dropped SYN returns as roughly 1200ms of apparent latency -- the 1s timer
/// plus the round trip. Such readings cluster at 1200ms and 3200ms, the
Expand Down
Loading