use std::sync::Arc;
use anyhow::Result;
use tokio::sync::{broadcast, watch};
use tracing::{info, warn};
use crate::hub::{spawn_dispatcher, spawn_health_checker, HubState, InMemoryQueue, MessageQueue};
use crate::ilink::{LoginClient, UpstreamClient};
use crate::server::build_router;
use crate::store::Store;
#[derive(Debug, Clone)]
pub struct ServeOptions {
pub token: Option<String>,
pub addr: String,
pub ilink_base_url: Option<String>,
pub database_url: String,
}
pub async fn run_serve(opts: ServeOptions, mut shutdown_rx: watch::Receiver<bool>) -> Result<()> {
let ServeOptions {
token: token_arg,
addr,
ilink_base_url,
database_url,
} = opts;
info!(%addr, %database_url, "iLink Hub starting");
let store = Arc::new(Store::connect(&database_url).await?);
let (token, base_url) = resolve_token(token_arg, ilink_base_url, store.clone()).await?;
let upstream = Arc::new(UpstreamClient::new(token, Some(base_url)));
let queue = build_queue_backend();
let state = HubState::new(upstream.clone(), store.clone(), queue);
load_clients_from_db(state.clone(), store.clone()).await;
let (tx, rx) = broadcast::channel::<crate::ilink::types::WeixinMessage>(256);
spawn_dispatcher(state.clone(), rx);
spawn_health_checker(state.clone());
{
let upstream_clone = upstream.clone();
let tx_clone = tx.clone();
let shutdown_rx_clone = shutdown_rx.clone();
tokio::spawn(async move {
upstream_clone
.run_polling_loop(tx_clone, shutdown_rx_clone)
.await;
});
}
let identity = crate::relay::DeviceIdentity::load_or_create()?;
if crate::relay::relay_enabled() {
let hub_base = format!("http://{}", crate::relay::hub_loopback_addr(&addr));
let relay_ws = crate::relay::relay_ws_url();
let pair_url = crate::relay::resolve_pair_public_url(identity.device_id());
info!(
device_id = %identity.device_id(),
pair_url = %pair_url,
relay = %relay_ws,
"pairing relay enabled (zero-config)"
);
crate::relay::client::spawn_relay_client(identity, hub_base, relay_ws);
} else {
info!("pairing relay disabled (set HUB_PAIR_URL or ILINKHUB_RELAY=0)");
}
let router = build_router(state);
let listener = tokio::net::TcpListener::bind(&addr).await?;
info!(%addr, "iLink Hub listening");
axum::serve(listener, router)
.with_graceful_shutdown(async move {
while !*shutdown_rx.borrow() {
if shutdown_rx.changed().await.is_err() {
break;
}
}
})
.await?;
info!("iLink Hub stopped");
Ok(())
}
async fn resolve_token(
token_arg: Option<String>,
ilink_base_url: Option<String>,
store: Arc<Store>,
) -> Result<(String, String)> {
let default_base = "https://ilinkai.weixin.qq.com".to_string();
let from_cli = token_arg.is_some();
let (candidate, base) = if let Some(token) = token_arg {
let base = ilink_base_url
.clone()
.unwrap_or_else(|| default_base.clone());
(Some(token), base)
} else if let Some((token, base)) = store.load_credentials().await? {
info!("loaded bot token from database");
(Some(token), base)
} else {
(
None,
ilink_base_url
.clone()
.unwrap_or_else(|| default_base.clone()),
)
};
if let Some(token) = candidate {
let client = UpstreamClient::new(token.clone(), Some(base.clone()));
if client.validate_session().await {
if from_cli {
store.save_credentials(&token, &base).await?;
}
return Ok((token, base));
}
warn!("iLink token invalid or expired");
println!();
println!("⚠️ 未检测到有效的 iLink 微信登录态,请扫描下方二维码完成登录。");
println!();
} else {
info!("no iLink token found, starting QR login");
println!();
println!("首次启动需要绑定微信机器人,请扫描下方二维码登录。");
println!();
}
let login_base = ilink_base_url.clone();
let login_client = LoginClient::new(ilink_base_url);
let token = login_client.login_with_qr().await?;
let base = login_base.unwrap_or(default_base);
store.save_credentials(&token, &base).await?;
info!("iLink login successful, token saved");
Ok((token, base))
}
async fn load_clients_from_db(state: Arc<HubState>, store: Arc<Store>) {
match store.list_clients().await {
Ok(clients) => {
let count = clients.len();
let mut registry = state.registry.write().await;
for c in clients {
registry.register_with_vtoken(
c.name.clone(),
c.label.clone(),
Some(c.vtoken.clone()),
);
}
info!(count, "loaded clients from database");
}
Err(e) => {
tracing::warn!(error = %e, "failed to load clients from DB");
}
}
match store.list_routes().await {
Ok(routes) => {
let count = routes.len();
let mut router = state.router.lock().await;
for (from_user, vtoken) in routes {
router.set_route(&from_user, vtoken);
}
info!(count, "loaded routing state from database");
}
Err(e) => {
tracing::warn!(error = %e, "failed to load routing state from DB");
}
}
match store.list_recent_context_tokens(500).await {
Ok(entries) => {
let count = entries.len();
let mut ctx_map = state.ctx_map.lock().await;
for (vctx, real_ctx, peer_user_id) in entries {
ctx_map.seed_full(vctx, real_ctx, peer_user_id);
}
info!(count, "warmed context_token cache from database");
}
Err(e) => {
tracing::warn!(error = %e, "failed to load context_token cache from DB");
}
}
}
fn build_queue_backend() -> Arc<dyn MessageQueue> {
match std::env::var("ILINK_QUEUE_BACKEND")
.as_deref()
.unwrap_or("")
{
"memory" | "" => {
info!(backend = "memory", "queue backend initialized");
Arc::new(InMemoryQueue::new())
}
"redis" => {
tracing::warn!("Redis queue backend is not yet implemented; falling back to memory");
Arc::new(InMemoryQueue::new())
}
other => {
tracing::error!(
backend = other,
"Unknown ILINK_QUEUE_BACKEND; supported values: 'memory'. \
(redis: planned, not yet available). Falling back to memory."
);
Arc::new(InMemoryQueue::new())
}
}
}