use std::{
net::{IpAddr, Ipv4Addr, SocketAddr, ToSocketAddrs},
sync::mpsc,
time::Duration,
};
use anyhow::Context as _;
use geph5_broker_protocol::{Credential, ExitConstraint};
use geph5_misc_rpc::{
client_config::BrokerSource,
client_control::ControlClient,
manager_control::{ProxySettings, TunnelSettings},
};
use serde::{Deserialize, Serialize};
use crate::platform::{self, EngineLaunch, EngineRole, PacketMode, SpawnedEngine};
pub const PAC_ADDR: SocketAddr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1)), 12223);
const CONFIG_TEMPLATE: &str = include_str!("../default-config.yaml");
#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct Settings {
#[serde(default)]
pub secret: Option<String>,
#[serde(default = "default_exit_constraint")]
pub exit_constraint: ExitConstraint,
#[serde(default)]
pub connected: bool,
#[serde(default)]
pub proxy: Option<ProxySettings>,
#[serde(default)]
pub vpn: bool,
#[serde(default = "default_true")]
pub allow_lan: bool,
#[serde(default)]
pub allow_direct: bool,
#[serde(default)]
pub passthrough_china: bool,
#[serde(default)]
pub session_metadata: serde_json::Value,
}
fn default_exit_constraint() -> ExitConstraint {
ExitConstraint::Auto
}
fn default_true() -> bool {
true
}
impl Default for Settings {
fn default() -> Self {
Settings {
secret: None,
exit_constraint: ExitConstraint::Auto,
connected: false,
proxy: None,
vpn: true,
allow_lan: true,
allow_direct: false,
passthrough_china: false,
session_metadata: serde_json::Value::Null,
}
}
}
impl Settings {
pub fn apply_tunnel_settings(&mut self, settings: TunnelSettings) {
self.exit_constraint = settings.exit_constraint;
self.proxy = settings.proxy;
self.vpn = settings.vpn;
self.allow_lan = settings.allow_lan;
self.allow_direct = settings.allow_direct;
self.passthrough_china = settings.passthrough_china;
self.session_metadata = settings.session_metadata;
}
pub fn load() -> anyhow::Result<Self> {
let path = platform::settings_path();
match std::fs::read(&path) {
Ok(bytes) => Ok(serde_json::from_slice(&bytes)
.with_context(|| format!("could not parse {}", path.display()))?),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(Settings::default()),
Err(e) => Err(e).with_context(|| format!("could not read {}", path.display())),
}
}
pub fn save(&self) -> anyhow::Result<()> {
std::fs::create_dir_all(platform::state_dir())?;
let path = platform::settings_path();
let bytes = serde_json::to_vec_pretty(self)?;
std::fs::write(&path, bytes).with_context(|| format!("could not write {}", path.display()))
}
}
fn config_from_template() -> anyhow::Result<geph5_misc_rpc::client_config::Config> {
let template_val: serde_json::Value =
serde_yaml::from_str(CONFIG_TEMPLATE).context("bad embedded config template")?;
serde_json::from_value(template_val).context("bad embedded config template")
}
pub(crate) fn build_tunnel_config(
settings: &Settings,
pre_resolve_broker_fronts: bool,
) -> anyhow::Result<String> {
let mut cfg = config_from_template()?;
cfg.credentials = Credential::Secret(settings.secret.clone().unwrap_or_default());
cfg.exit_constraint = settings.exit_constraint.clone();
cfg.allow_lan = settings.allow_lan;
cfg.allow_direct = settings.allow_direct;
cfg.passthrough_china = settings.passthrough_china;
cfg.spoof_dns = settings.passthrough_china;
cfg.sess_metadata = settings.session_metadata.clone();
cfg.dry_run = !settings.connected;
platform::configure_engine_control(&mut cfg, EngineRole::Tunnel);
match &settings.proxy {
Some(p) => {
let ip: IpAddr = if p.listen_all {
Ipv4Addr::UNSPECIFIED.into()
} else {
Ipv4Addr::LOCALHOST.into()
};
cfg.socks5_listen = Some(SocketAddr::new(ip, p.socks5_port));
cfg.http_proxy_listen = Some(SocketAddr::new(ip, p.http_port));
cfg.pac_listen = Some(PAC_ADDR);
}
None => {
cfg.socks5_listen = None;
cfg.http_proxy_listen = None;
cfg.pac_listen = None;
}
}
let secret = settings.secret.as_deref().unwrap_or_default();
let cache_tag = blake3::hash(secret.as_bytes()).to_hex();
cfg.cache = Some(platform::cache_dir().join(format!("db-{cache_tag}")));
if pre_resolve_broker_fronts && let Some(broker) = cfg.broker.as_mut() {
pre_resolve_fronted_broker_sources(broker);
}
let val = serde_json::to_value(&cfg).context("could not serialize child config")?;
serde_yaml::to_string(&val).context("could not serialize child config")
}
fn pre_resolve_fronted_broker_sources(source: &mut BrokerSource) {
match source {
BrokerSource::Fronted {
front,
override_dns,
..
} if override_dns.is_none() => match resolve_front_override_dns(front) {
Ok(addrs) => {
tracing::info!(
front = %front,
addr_count = addrs.len(),
"pre-resolved fronted broker source for VPN bootstrap"
);
*override_dns = Some(addrs);
}
Err(err) => {
tracing::warn!(
front = %front,
err = %err,
"could not pre-resolve fronted broker source for VPN bootstrap"
);
}
},
BrokerSource::Race(sources) => {
for source in sources {
pre_resolve_fronted_broker_sources(source);
}
}
BrokerSource::PriorityRace(sources) => {
for source in sources.values_mut() {
pre_resolve_fronted_broker_sources(source);
}
}
_ => {}
}
}
const FRONT_RESOLVE_TIMEOUT: Duration = Duration::from_millis(500);
fn resolve_front_override_dns(front: &str) -> anyhow::Result<Vec<SocketAddr>> {
let (host, port) = front_lookup_target(front)?;
let (tx, rx) = mpsc::channel();
std::thread::spawn(move || {
let _ = tx.send(
(host.as_str(), port)
.to_socket_addrs()
.map(|it| it.collect::<Vec<_>>())
.with_context(|| format!("could not resolve {host}:{port}")),
);
});
let resolved = rx
.recv_timeout(FRONT_RESOLVE_TIMEOUT)
.context("front DNS pre-resolve timed out")??;
let mut addrs = Vec::new();
for addr in resolved {
if !addrs.contains(&addr) {
addrs.push(addr);
}
}
if addrs.is_empty() {
anyhow::bail!("front resolved to no socket addresses");
}
Ok(addrs)
}
fn front_lookup_target(front: &str) -> anyhow::Result<(String, u16)> {
let url = url::Url::parse(front).context("front is not a valid URL")?;
let host = url
.host_str()
.context("front URL has no hostname")?
.to_string();
let port = url
.port_or_known_default()
.context("front URL has no explicit or scheme-default port")?;
Ok((host, port))
}
#[cfg(test)]
mod tests {
use super::{Settings, build_tunnel_config, front_lookup_target};
use geph5_broker_protocol::ExitConstraint;
use geph5_misc_rpc::manager_control::{ProxySettings, TunnelSettings};
#[test]
fn front_lookup_target_uses_https_default_port() {
assert_eq!(
front_lookup_target("https://www.cdn77.com/").unwrap(),
("www.cdn77.com".to_string(), 443)
);
}
#[test]
fn front_lookup_target_preserves_explicit_port() {
assert_eq!(
front_lookup_target("http://example.com:8080/path").unwrap(),
("example.com".to_string(), 8080)
);
}
#[test]
fn front_lookup_target_rejects_urls_without_socket_target() {
assert!(front_lookup_target("file:///tmp/front").is_err());
}
#[test]
fn applying_tunnel_snapshot_preserves_lifecycle_and_credentials() {
let mut settings = Settings {
secret: Some("keep-me".into()),
connected: true,
..Settings::default()
};
settings.apply_tunnel_settings(TunnelSettings {
exit_constraint: ExitConstraint::Auto,
proxy: Some(ProxySettings::default()),
vpn: false,
allow_lan: false,
allow_direct: true,
passthrough_china: true,
session_metadata: serde_json::json!({"filter": {"ads": true}}),
});
assert_eq!(settings.secret.as_deref(), Some("keep-me"));
assert!(settings.connected);
assert!(!settings.vpn);
assert!(!settings.allow_lan);
assert!(settings.allow_direct);
assert!(settings.passthrough_china);
assert_eq!(settings.session_metadata["filter"]["ads"], true);
}
#[test]
fn tunnel_config_contains_the_complete_snapshot() {
let settings = Settings {
secret: Some("secret".into()),
connected: true,
allow_lan: false,
allow_direct: true,
passthrough_china: true,
session_metadata: serde_json::json!({"filter": {"nsfw": true}}),
..Settings::default()
};
let yaml = build_tunnel_config(&settings, false).unwrap();
let value: serde_json::Value = serde_yaml::from_str(&yaml).unwrap();
let config: geph5_misc_rpc::client_config::Config = serde_json::from_value(value).unwrap();
assert!(!config.allow_lan);
assert!(config.allow_direct);
assert!(config.passthrough_china);
assert_eq!(config.sess_metadata["filter"]["nsfw"], true);
}
}
fn build_query_config() -> anyhow::Result<String> {
let mut cfg = config_from_template()?;
cfg.credentials = Credential::Secret(String::new());
cfg.dry_run = true;
cfg.socks5_listen = None;
cfg.http_proxy_listen = None;
cfg.pac_listen = None;
platform::configure_engine_control(&mut cfg, EngineRole::Query);
cfg.cache = Some(platform::cache_dir().join("query-db"));
let val = serde_json::to_value(&cfg).context("could not serialize query config")?;
serde_yaml::to_string(&val).context("could not serialize query config")
}
pub fn spawn_tunnel(
config_yaml: String,
service_user: Option<(u32, u32)>,
packet_mode: PacketMode,
bind_indices: Option<(u32, u32)>,
) -> anyhow::Result<SpawnedEngine> {
let config_path = write_engine_config("child-config.yaml", config_yaml)?;
platform::spawn_engine(EngineLaunch {
role: EngineRole::Tunnel,
config_path,
service_user,
packet_mode,
bind_indices,
})
}
pub fn spawn_query(service_user: Option<(u32, u32)>) -> anyhow::Result<std::process::Child> {
let config_yaml = build_query_config()?;
let config_path = write_engine_config("query-config.yaml", config_yaml)?;
let spawned = platform::spawn_engine(EngineLaunch {
role: EngineRole::Query,
config_path,
service_user,
packet_mode: PacketMode::None,
bind_indices: None,
})?;
match spawned.transport {
platform::ChildTransport::None => Ok(spawned.child),
platform::ChildTransport::Stdio { .. } => {
platform::kill_child(spawned.child);
anyhow::bail!("query engine unexpectedly exposed a packet transport")
}
}
}
fn write_engine_config(
config_name: &str,
config_yaml: String,
) -> anyhow::Result<std::path::PathBuf> {
std::fs::create_dir_all(platform::state_dir())?;
let cache_dir = platform::cache_dir();
if cache_dir.is_file() {
let _ = std::fs::remove_file(&cache_dir);
for sfx in ["-shm", "-wal", "-journal"] {
let _ = std::fs::remove_file(format!("{}{sfx}", cache_dir.display()));
}
}
std::fs::create_dir_all(&cache_dir)?;
let config_path = platform::state_dir().join(config_name);
std::fs::write(&config_path, config_yaml)
.with_context(|| format!("could not write {}", config_path.display()))?;
Ok(config_path)
}
pub fn live_control() -> ControlClient {
platform::engine_control(EngineRole::Tunnel)
}
pub fn query_control() -> ControlClient {
platform::engine_control(EngineRole::Query)
}
pub async fn wait_control_ready(client: &ControlClient, timeout: Duration) -> anyhow::Result<()> {
let deadline = Duration::from_millis(100);
let start = std::time::Instant::now();
loop {
match client.start_time().await {
Ok(_) => return Ok(()),
Err(_) if start.elapsed() < timeout => {
tokio::time::sleep(deadline).await;
}
Err(e) => anyhow::bail!("engine never became reachable: {e:?}"),
}
}
}