use std::{
process::Child,
time::{Duration, SystemTime, UNIX_EPOCH},
};
use async_trait::async_trait;
use geph5_broker_protocol::{Credential, UserInfo};
use geph5_misc_rpc::{
client_control::{ConnInfo, ControlClient},
manager_control::{
AccountInfo, ConnState, ExitInfo, GephCtlProtocol, SessionContext, SettingsView, Status,
TunnelSettings,
},
};
use nanorpc::RpcTransport;
use tokio::sync::Mutex;
use crate::{
platform,
supervisor::{self, Settings},
vpn,
};
const CHILD_READY_TIMEOUT: Duration = Duration::from_secs(30);
const CHILD_HEALTH_INTERVAL: Duration = Duration::from_secs(3);
enum ChildRecovery {
Missing,
Exited(std::process::ExitStatus),
Uninspectable(std::io::Error),
}
fn child_recovery(child: Option<&mut Child>) -> Option<ChildRecovery> {
match child {
None => Some(ChildRecovery::Missing),
Some(child) => match child.try_wait() {
Ok(None) => None,
Ok(Some(status)) => Some(ChildRecovery::Exited(status)),
Err(error) => Some(ChildRecovery::Uninspectable(error)),
},
}
}
struct Inner {
settings: Settings,
query: Option<Child>,
tunnel: Option<Child>,
vpn: vpn::Vpn,
proxy_active: bool,
proxy_session: Option<SessionContext>,
desktop_session: Option<SessionContext>,
shutting_down: bool,
}
#[derive(Clone)]
pub struct ManagerImpl {
inner: std::sync::Arc<Mutex<Inner>>,
}
impl ManagerImpl {
pub async fn start() -> anyhow::Result<Self> {
vpn::cleanup_stale();
let settings = Settings::load()?;
let this = ManagerImpl {
inner: std::sync::Arc::new(Mutex::new(Inner {
settings,
query: None,
tunnel: None,
vpn: vpn::Vpn::new(),
proxy_active: false,
proxy_session: None,
desktop_session: None,
shutting_down: false,
})),
};
let mut inner = this.inner.lock().await;
if let Err(e) = this.ensure_query_engine(&mut inner).await {
tracing::warn!(err = %e, "query engine failed to start; broker queries unavailable");
}
if let Err(e) = Self::reconcile_tunnel(&mut inner).await {
tracing::warn!(err = %e, "could not restore persisted state; starting disconnected");
inner.settings.connected = false;
let _ = inner.settings.save();
let _ = Self::reconcile_tunnel(&mut inner).await;
}
drop(inner);
Ok(this)
}
async fn ensure_query_engine(&self, inner: &mut Inner) -> Result<(), String> {
if inner.query.is_some() {
return Ok(());
}
let service_user =
platform::ensure_service_user().map_err(|e| format!("engine service user: {e:#}"))?;
let child = supervisor::spawn_query(service_user).map_err(|e| format!("{e:?}"))?;
inner.query = Some(child);
if let Err(error) =
supervisor::wait_control_ready(&supervisor::query_control(), CHILD_READY_TIMEOUT).await
{
if let Some(child) = inner.query.take() {
geph5_rt::spawn_blocking(move || platform::kill_child(child)).await;
}
return Err(format!("{error:?}"));
}
Ok(())
}
async fn reconcile_tunnel(inner: &mut Inner) -> Result<(), String> {
let pre_resolve_broker_fronts = inner.settings.connected && inner.settings.vpn;
let tunnel_config = if inner.settings.connected {
let settings = inner.settings.clone();
let cfg = geph5_rt::spawn_blocking(move || {
supervisor::build_tunnel_config(&settings, pre_resolve_broker_fronts)
})
.await
.map_err(|e| format!("{e:#}"))?;
Some(cfg)
} else {
None
};
if let Some(child) = inner.tunnel.take() {
geph5_rt::spawn_blocking(move || platform::kill_child(child)).await;
}
inner.vpn.stop_transport();
if !inner.settings.connected {
reconcile_system_proxy(inner).await?;
inner.vpn.cleanup();
return Ok(());
}
let service_user =
platform::ensure_service_user().map_err(|e| format!("engine service user: {e:#}"))?;
let want_vpn = inner.settings.vpn;
if want_vpn {
inner
.vpn
.ensure_active(inner.settings.allow_lan, service_user)
.map_err(|e| format!("vpn reconcile: {e:#}"))?;
}
let tunnel_config = tunnel_config.expect("connected tunnel has a prepared config");
let spawned = supervisor::spawn_tunnel(
tunnel_config,
service_user,
inner.vpn.packet_mode(want_vpn),
inner.vpn.bind_indices(want_vpn),
)
.map_err(|e| format!("{e:?}"))?;
if let Err(error) = inner.vpn.attach_transport(want_vpn, spawned.transport) {
platform::kill_child(spawned.child);
return Err(format!("attaching VPN packet transport: {error:#}"));
}
let child = spawned.child;
if let Err(error) =
supervisor::wait_control_ready(&supervisor::live_control(), CHILD_READY_TIMEOUT).await
{
inner.vpn.stop_transport();
platform::kill_child(child);
return Err(format!("{error:?}"));
}
if let Err(error) = reconcile_system_proxy(inner).await {
inner.vpn.stop_transport();
platform::kill_child(child);
return Err(error);
}
if !want_vpn {
inner.vpn.cleanup();
}
inner.tunnel = Some(child);
Ok(())
}
async fn account_for_secret(&self, secret: &str) -> Result<AccountInfo, String> {
let cred = Credential::Secret(secret.to_string());
let params = vec![serde_json::to_value(&cred).map_err(|e| e.to_string())?];
let raw = supervisor::query_control()
.broker_rpc("get_user_info_by_cred".into(), params)
.await
.map_err(|e| format!("could not reach broker: {e:?}"))??;
let info: Option<UserInfo> = serde_json::from_value(raw).map_err(|e| e.to_string())?;
let info = info.ok_or_else(|| "incorrect secret".to_string())?;
Ok(account_info_from(info))
}
async fn forward_raw(
&self,
req: nanorpc::JrpcRequest,
) -> Result<nanorpc::JrpcResponse, String> {
let connected = self.inner.lock().await.settings.connected;
if connected {
match supervisor::live_control().0.call_raw(req.clone()).await {
Ok(resp) => return Ok(resp),
Err(e) => tracing::debug!(
err = ?e,
"tunnel engine unreachable; falling back to query engine"
),
}
}
supervisor::query_control()
.0
.call_raw(req)
.await
.map_err(|e| format!("could not reach engine: {e:?}"))
}
}
fn wants_auto_proxy(settings: &Settings) -> bool {
settings.proxy.as_ref().is_some_and(|p| p.autoconf)
}
async fn apply_proxy(session: Option<SessionContext>, connected: bool) -> Result<(), String> {
let url = format!("http://{}/proxy.pac", supervisor::PAC_ADDR);
let res = geph5_rt::spawn_blocking(move || {
platform::set_system_proxy(session.as_ref(), connected, &url)
})
.await;
res.map_err(|e| format!("system proxy config failed: {e:#}"))?;
tracing::info!(connected, "configured system proxy");
Ok(())
}
async fn reconcile_system_proxy(inner: &mut Inner) -> Result<(), String> {
let want_proxy =
inner.settings.connected && !inner.settings.vpn && wants_auto_proxy(&inner.settings);
let session = if want_proxy {
inner.desktop_session.clone()
} else {
inner
.proxy_session
.clone()
.or_else(|| inner.desktop_session.clone())
};
apply_proxy(session.clone(), want_proxy).await?;
inner.proxy_active = want_proxy;
inner.proxy_session = want_proxy.then_some(session).flatten();
Ok(())
}
fn now_unix() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or_default()
}
fn account_info_from(info: UserInfo) -> AccountInfo {
let is_plus = info
.plus_expires_unix
.map(|e| e > now_unix())
.unwrap_or(false);
AccountInfo {
user_id: info.user_id,
level: if is_plus { "plus" } else { "free" }.to_string(),
plus_expires_unix: info.plus_expires_unix,
bw_used_mb: info.bw_consumption.map(|b| b.mb_used),
bw_limit_mb: info.bw_consumption.map(|b| b.mb_limit),
}
}
fn exit_info_from(
hostname: String,
exit: &geph5_broker_protocol::ExitDescriptor,
meta: Option<&geph5_broker_protocol::ExitMetadata>,
) -> ExitInfo {
ExitInfo {
hostname,
country: exit.country.alpha2().to_string(),
city: exit.city.clone(),
load: exit.load,
allows_free: meta
.map(|m| {
m.allowed_levels
.contains(&geph5_broker_protocol::AccountLevel::Free)
})
.unwrap_or(false),
}
}
#[async_trait]
impl GephCtlProtocol for ManagerImpl {
async fn login(&self, secret: String) -> Result<AccountInfo, String> {
let secret = secret.trim().to_string();
let account = self.account_for_secret(&secret).await?;
let mut inner = self.inner.lock().await;
inner.settings.secret = Some(secret);
inner.settings.save().map_err(|e| format!("{e:?}"))?;
if inner.settings.connected {
Self::reconcile_tunnel(&mut inner).await?;
}
Ok(account)
}
async fn set_secret(&self, secret: String) -> Result<(), String> {
let secret = secret.trim().to_string();
let mut inner = self.inner.lock().await;
if inner.settings.secret.as_deref() != Some(secret.as_str()) {
inner.settings.secret = Some(secret);
inner.settings.save().map_err(|e| format!("{e:?}"))?;
if inner.settings.connected {
Self::reconcile_tunnel(&mut inner).await?;
}
}
Ok(())
}
async fn logout(&self, session: SessionContext) -> Result<(), String> {
let mut inner = self.inner.lock().await;
inner.desktop_session = Some(session);
inner.settings.secret = None;
inner.settings.connected = false;
inner.settings.save().map_err(|e| format!("{e:?}"))?;
Self::reconcile_tunnel(&mut inner).await
}
async fn account(&self) -> Result<AccountInfo, String> {
let secret = {
let inner = self.inner.lock().await;
inner
.settings
.secret
.clone()
.ok_or_else(|| "not logged in".to_string())?
};
self.account_for_secret(&secret).await
}
async fn connect(&self, session: SessionContext) -> Result<(), String> {
let mut inner = self.inner.lock().await;
if inner.settings.secret.is_none() {
return Err("not logged in".to_string());
}
inner.desktop_session = Some(session);
inner.settings.connected = true;
inner.settings.save().map_err(|e| format!("{e:?}"))?;
Self::reconcile_tunnel(&mut inner).await
}
async fn reconnect(&self, session: SessionContext) -> Result<(), String> {
let mut inner = self.inner.lock().await;
if !inner.settings.connected {
return Err("not connected".to_string());
}
inner.desktop_session = Some(session);
Self::reconcile_tunnel(&mut inner).await
}
async fn disconnect(&self, session: SessionContext) -> Result<(), String> {
let mut inner = self.inner.lock().await;
inner.desktop_session = Some(session);
inner.settings.connected = false;
inner.settings.save().map_err(|e| format!("{e:?}"))?;
Self::reconcile_tunnel(&mut inner).await
}
async fn status(&self) -> Result<Status, String> {
if !self.inner.lock().await.settings.connected {
return Ok(Status {
state: ConnState::Disconnected,
exit: None,
total_rx_bytes: 0.0,
total_tx_bytes: 0.0,
});
}
let client: ControlClient = supervisor::live_control();
let conn = match client.conn_info().await {
Ok(conn) => conn,
Err(_) => {
return Ok(Status {
state: ConnState::Connecting,
exit: None,
total_rx_bytes: 0.0,
total_tx_bytes: 0.0,
});
}
};
let (state, exit) = match conn {
ConnInfo::Disconnected => (ConnState::Disconnected, None),
ConnInfo::Connecting => (ConnState::Connecting, None),
ConnInfo::Connected { sessions } => {
let exit = sessions
.first()
.map(|s| exit_info_from(s.exit.country.alpha2().to_string(), &s.exit, None));
(ConnState::Connected, exit)
}
};
let total_rx_bytes = client
.stat_num("total_rx_bytes".into())
.await
.unwrap_or(0.0);
let total_tx_bytes = client
.stat_num("total_tx_bytes".into())
.await
.unwrap_or(0.0);
Ok(Status {
state,
exit,
total_rx_bytes,
total_tx_bytes,
})
}
async fn get_settings(&self) -> Result<SettingsView, String> {
let inner = self.inner.lock().await;
Ok(SettingsView {
logged_in: inner.settings.secret.is_some(),
exit_constraint: inner.settings.exit_constraint.clone(),
connected: inner.settings.connected,
proxy: inner.settings.proxy.clone(),
vpn: inner.settings.vpn,
allow_lan: inner.settings.allow_lan,
allow_direct: inner.settings.allow_direct,
passthrough_china: inner.settings.passthrough_china,
session_metadata: inner.settings.session_metadata.clone(),
})
}
async fn apply_settings(
&self,
settings: TunnelSettings,
session: SessionContext,
) -> Result<(), String> {
let mut inner = self.inner.lock().await;
tracing::debug!(?settings, "applying complete tunnel settings snapshot");
inner.desktop_session = Some(session);
inner.settings.apply_tunnel_settings(settings);
inner.settings.save().map_err(|e| format!("{e:?}"))?;
if inner.settings.connected {
Self::reconcile_tunnel(&mut inner).await
} else {
Ok(())
}
}
async fn list_exits(&self) -> Result<Vec<ExitInfo>, String> {
let net_status = supervisor::query_control()
.net_status()
.await
.map_err(|e| format!("could not reach broker: {e:?}"))??;
let mut out: Vec<ExitInfo> = net_status
.exits
.into_iter()
.map(|(hostname, (_pk, exit, meta))| exit_info_from(hostname, &exit, Some(&meta)))
.collect();
out.sort_by(|a, b| {
(a.country.as_str(), a.city.as_str()).cmp(&(b.country.as_str(), b.city.as_str()))
});
Ok(out)
}
async fn logs(&self, count: usize) -> Result<Vec<String>, String> {
let connected = self.inner.lock().await.settings.connected;
let client = if connected {
supervisor::live_control()
} else {
supervisor::query_control()
};
let mut logs = client
.recent_logs()
.await
.map_err(|e| format!("could not reach engine: {e:?}"))?;
if logs.len() > count {
logs = logs.split_off(logs.len() - count);
}
Ok(logs)
}
async fn daemon_rpc(
&self,
method: String,
params: Vec<serde_json::Value>,
) -> Result<serde_json::Value, String> {
let req = nanorpc::JrpcRequest {
jsonrpc: "2.0".into(),
method,
params,
id: nanorpc::JrpcId::Number(1),
};
let resp = self.forward_raw(req).await?;
match resp.error {
Some(err) => Err(err.message),
None => Ok(resp.result.unwrap_or(serde_json::Value::Null)),
}
}
}
async fn shutdown_teardown(inner: &Mutex<Inner>, teardown_lock: &Mutex<()>) {
let _teardown = teardown_lock.lock().await;
let (tunnel, query, proxy_active, proxy_session) = {
let mut inner = inner.lock().await;
inner.shutting_down = true;
let tunnel = inner.tunnel.take();
let query = inner.query.take();
let proxy_active = std::mem::take(&mut inner.proxy_active);
let proxy_session = inner.proxy_session.take();
(tunnel, query, proxy_active, proxy_session)
};
if let Some(child) = tunnel {
platform::kill_child(child);
}
if let Some(child) = query {
platform::kill_child(child);
}
inner.lock().await.vpn.cleanup();
if proxy_active {
let _ = apply_proxy(proxy_session, false).await;
}
}
pub async fn run_manager() -> anyhow::Result<()> {
let manager = ManagerImpl::start().await?;
let teardown_lock = std::sync::Arc::new(Mutex::new(()));
{
let inner = manager.inner.clone();
let teardown_lock = teardown_lock.clone();
geph5_rt::spawn(async move {
platform::shutdown_signal().await;
tracing::warn!("termination signal received; tearing down manager state and exiting");
shutdown_teardown(&inner, &teardown_lock).await;
std::process::exit(0);
})
.detach();
}
{
let inner = manager.inner.clone();
geph5_rt::spawn(async move {
loop {
vpn::wait_network_change().await;
let probe = {
let inner = inner.lock().await;
if inner.shutting_down || !(inner.settings.connected && inner.settings.vpn) {
continue;
}
match inner.vpn.network_probe() {
Some(probe) => probe,
None => continue,
}
};
let checked = geph5_rt::spawn_blocking(move || vpn::check_network(probe)).await;
let mut inner = inner.lock().await;
if checked.generation != inner.vpn.generation()
|| inner.shutting_down
|| !(inner.settings.connected && inner.settings.vpn)
{
continue;
}
match checked.action {
vpn::NetworkAction::Healthy => {}
vpn::NetworkAction::Reconcile => {
tracing::warn!("VPN network state changed; reconciling complete state");
if let Err(error) = ManagerImpl::reconcile_tunnel(&mut inner).await {
tracing::warn!(%error, "VPN network reconciliation failed");
}
}
}
}
})
.detach();
}
{
let manager = manager.clone();
geph5_rt::spawn(async move {
loop {
tokio::time::sleep(CHILD_HEALTH_INTERVAL).await;
let mut inner = manager.inner.lock().await;
if inner.shutting_down {
continue;
}
if let Some(reason) = child_recovery(inner.query.as_mut()) {
match &reason {
ChildRecovery::Missing => {}
ChildRecovery::Exited(status) => {
tracing::warn!(%status, "query engine exited; restarting");
}
ChildRecovery::Uninspectable(error) => {
tracing::warn!(%error, "could not inspect query engine; restarting");
}
}
if let Some(child) = inner.query.take()
&& matches!(reason, ChildRecovery::Uninspectable(_))
{
geph5_rt::spawn_blocking(move || platform::kill_child(child)).await;
}
if let Err(error) = manager.ensure_query_engine(&mut inner).await {
tracing::warn!(%error, "query engine restart failed");
}
}
if !inner.settings.connected {
continue;
}
if let Some(reason) = child_recovery(inner.tunnel.as_mut()) {
match &reason {
ChildRecovery::Missing => {}
ChildRecovery::Exited(status) => {
tracing::warn!(%status, "tunnel engine exited; reconciling");
}
ChildRecovery::Uninspectable(error) => {
tracing::warn!(%error, "could not inspect tunnel engine; reconciling");
}
}
if let Some(child) = inner.tunnel.take()
&& matches!(reason, ChildRecovery::Uninspectable(_))
{
geph5_rt::spawn_blocking(move || platform::kill_child(child)).await;
}
inner.vpn.stop_transport();
if let Err(error) = ManagerImpl::reconcile_tunnel(&mut inner).await {
tracing::warn!(%error, "tunnel engine recovery failed");
}
}
}
})
.detach();
}
let inner = manager.inner.clone();
let result = platform::serve_manager(manager).await;
tracing::warn!("manager control server exited; tearing down installed state");
shutdown_teardown(&inner, &teardown_lock).await;
result
}