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
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
- 实时监控:秒级实时数据展示
- 轻量高效:Rust 语言构建,低资源占用,极简高效
- 自托管:完全掌控数据隐私,部署简单
- 通知:节点掉线、流量、到期与登录,推送到 Telegram 或自定义 Webhook

## 组成

Expand Down
1 change: 1 addition & 0 deletions install-hub.sh
Original file line number Diff line number Diff line change
Expand Up @@ -222,6 +222,7 @@ User=$USER_NAME
WorkingDirectory=$ROOT
ReadWritePaths=$DATA
NoNewPrivileges=yes
RestrictSUIDSGID=yes
ProtectSystem=strict
ProtectHome=yes
PrivateTmp=yes
Expand Down
20 changes: 18 additions & 2 deletions src/api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -146,6 +146,7 @@ fn node_view(node: &Node, current: Option<&Agent>, traffic: &Traffic, full: bool
view["ipv6"] = json!(node.ipv6);
view["remark"] = json!(node.remark);
view["token"] = json!(node.token);
view["notify"] = json!(node.notify);
}
view
}
Expand Down Expand Up @@ -1395,6 +1396,7 @@ pub async fn settings(_: Admin, State(app): State<Shared>) -> Json<Value> {
for key in ["register_key", "register_until"] {
out.insert(key.into(), json!(app.db.get(key).unwrap_or_default()));
}
crate::notify::settings(&app, &mut out);
Json(Value::Object(out))
}

Expand Down Expand Up @@ -1432,6 +1434,7 @@ fn setting_error(app: &App, key: &str, value: &Value) -> Option<String> {
}
"admin_password" if value.len() < 12 => Some("password must be at least 12 characters".into()),
"admin_password" => None,
k if k.starts_with("notify_") => crate::notify::setting_error(k, value),
k if READABLE_SETTINGS.contains(&k) || k == "github_client_secret" => None,
_ => Some(format!("unknown setting: {key}")),
}
Expand Down Expand Up @@ -2549,6 +2552,13 @@ mod tests {
"retention_days": read["retention_days"],
"github_proxy": read["github_proxy"],
"public_page": "on",
"notify_grace": read["notify_grace"],
"notify_traffic": read["notify_traffic"],
"notify_expiry": read["notify_expiry"],
"notify_login": read["notify_login"],
"notify_telegram_chat": read["notify_telegram_chat"],
"notify_telegram_text": read["notify_telegram_text"],
"notify_webhook_body": read["notify_webhook_body"],
});
assert_eq!(
save_settings(Admin, State(app.clone()), HeaderMap::new(), Json(echoed)).await.status(),
Expand All @@ -2559,15 +2569,21 @@ mod tests {
}

#[tokio::test]
async fn settings_never_hand_back_the_github_secret() {
async fn settings_never_hand_back_a_secret() {
let app = app();
app.db.set("github_client_secret", "super-secret").unwrap();
app.db.set("github_client_id", "public-id").unwrap();
app.db.set("notify_telegram_token", "123:bot-secret").unwrap();
app.db.set("notify_webhook_url", "https://hooks.example/url-secret").unwrap();
app.db.set("notify_webhook_headers", "Authorization: header-secret").unwrap();

let Json(body) = settings(Admin, axum::extract::State(std::sync::Arc::new(app))).await;
assert_eq!(body["github_client_id"], "public-id");
assert_eq!(body["github_secret_set"], true);
assert_eq!(body["notify_webhook_url_set"], true);
assert!(body.get("github_client_secret").is_none());
assert!(!body.to_string().contains("super-secret"));
for secret in ["super-secret", "bot-secret", "url-secret", "header-secret"] {
assert!(!body.to_string().contains(secret), "{secret}");
}
}
}
21 changes: 14 additions & 7 deletions src/auth.rs
Original file line number Diff line number Diff line change
Expand Up @@ -183,7 +183,10 @@ pub async fn login(
}
app.throttle.clear(ip);
match issue_session(&app, &headers) {
Ok(cookie) => with_cookies(Json(serde_json::json!({"ok": true})), [cookie]),
Ok(cookie) => {
crate::notify::signed_in(&app, "应急密码", ip);
with_cookies(Json(serde_json::json!({"ok": true})), [cookie])
}
Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()).into_response(),
}
}
Expand Down Expand Up @@ -228,6 +231,7 @@ pub struct Callback {

pub async fn github_callback(
State(app): State<crate::Shared>,
ConnectInfo(peer): ConnectInfo<std::net::SocketAddr>,
headers: HeaderMap,
Query(query): Query<Callback>,
) -> Response {
Expand All @@ -248,13 +252,15 @@ pub async fn github_callback(
let Some(code) = query.code.as_deref().filter(|c| !c.is_empty()) else {
return sign_in_failed(&app, &headers, "GitHub sent no authorization code");
};
if let Err(e) = github_login(&app, code).await {
return sign_in_failed(&app, &headers, &e.to_string());
}
let user = match github_login(&app, code).await {
Ok(user) => user,
Err(e) => return sign_in_failed(&app, &headers, &e.to_string()),
};
let session = match issue_session(&app, &headers) {
Ok(cookie) => cookie,
Err(e) => return sign_in_failed(&app, &headers, &e.to_string()),
};
crate::notify::signed_in(&app, &format!("GitHub {user}"), client_ip(&headers, peer.ip()));
with_cookies(Redirect::to("/admin"), [clear_state(&app, &headers), session])
}

Expand Down Expand Up @@ -303,8 +309,9 @@ fn urlencode(value: &str) -> String {
.collect()
}

/// Exchanges the code for a token and checks the login against the allow list.
async fn github_login(app: &App, code: &str) -> Result<()> {
/// Exchanges the code for a token and checks the login against the allow list,
/// returning the accepted login.
async fn github_login(app: &App, code: &str) -> Result<String> {
let (Some(id), Some(secret)) = (app.db.get("github_client_id"), app.db.get("github_client_secret"))
else {
bail!("not configured");
Expand Down Expand Up @@ -365,7 +372,7 @@ async fn github_login(app: &App, code: &str) -> Result<()> {
bail!("GitHub user {} is not on the allowed list", user.login);
}
info!("GitHub sign-in accepted for {}", user.login);
Ok(())
Ok(user.login)
}

/// Peer address, or the last hop in X-Forwarded-For when the request arrived
Expand Down
35 changes: 32 additions & 3 deletions src/db.rs
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,12 @@ CREATE TABLE IF NOT EXISTS node (
-- Survives the disconnection it describes, unlike the in-memory live entry:
-- an offline node's page is exactly where "since when" is worth reading.
last_seen INTEGER NOT NULL DEFAULT 0,
-- Opt-in, as the operator decides which machines are worth an alert.
notify INTEGER NOT NULL DEFAULT 0,
-- `last_seen` as of the offline alert, zero while none is outstanding. Stored
-- rather than held in memory so that a hub restart neither repeats the alert
-- nor loses the recovery that pairs with it.
down_since INTEGER NOT NULL DEFAULT 0,
created_at INTEGER NOT NULL
);

Expand Down Expand Up @@ -122,7 +128,7 @@ CREATE TABLE IF NOT EXISTS session (
/// Schema revision this build expects, stamped into `PRAGMA user_version`.
/// Increment it and add a `migrate_to_N` when the schema changes under a
/// database already in service.
const SCHEMA_VERSION: i64 = 3;
const SCHEMA_VERSION: i64 = 4;

/// Adds a column older databases lack. A duplicate column indicates the
/// migration has already run; every other error must propagate.
Expand Down Expand Up @@ -219,6 +225,11 @@ fn migrate_to_3(conn: &Connection) -> Result<()> {
add_column(conn, "node", "country TEXT NOT NULL DEFAULT ''")
}

fn migrate_to_4(conn: &Connection) -> Result<()> {
add_column(conn, "node", "notify INTEGER NOT NULL DEFAULT 0")?;
add_column(conn, "node", "down_since INTEGER NOT NULL DEFAULT 0")
}

/// Brings a database already in service up to `SCHEMA_VERSION` and stamps it.
/// `from` is its current version, so a fresh file passes `SCHEMA_VERSION` and
/// receives only the stamp.
Expand All @@ -235,6 +246,9 @@ fn migrate(conn: &Connection, from: i64) -> Result<()> {
if from < 3 {
migrate_to_3(conn)?;
}
if from < 4 {
migrate_to_4(conn)?;
}
conn.execute_batch(&format!("PRAGMA user_version = {SCHEMA_VERSION}"))?;
Ok(())
}
Expand Down Expand Up @@ -309,6 +323,11 @@ pub struct Node {
/// the metric row. Zero for a node that has never reported.
#[serde(default)]
pub last_seen: i64,
/// Whether going offline and coming back are announced. See `notify`.
#[serde(default)]
pub notify: bool,
#[serde(default)]
pub down_since: i64,
/// What the agent authenticates with. Readable so the panel can display an
/// install command on demand; it never leaves the admin view.
#[serde(default)]
Expand All @@ -334,6 +353,7 @@ pub struct NodePatch {
pub traffic_limit: Option<i64>,
pub traffic_mode: Option<String>,
pub traffic_reset_day: Option<u32>,
pub notify: Option<bool>,
}

fn expiry_patch<'de, D: serde::Deserializer<'de>>(d: D) -> Result<Option<Option<String>>, D::Error> {
Expand Down Expand Up @@ -553,7 +573,8 @@ impl Db {
expires_at=CASE WHEN ?8 THEN ?9 ELSE expires_at END,
remark=COALESCE(?10,remark), traffic_limit=COALESCE(?11,traffic_limit),
traffic_mode=COALESCE(?12,traffic_mode),
traffic_reset_day=COALESCE(?13,traffic_reset_day)
traffic_reset_day=COALESCE(?13,traffic_reset_day),
notify=COALESCE(?14,notify)
WHERE id=?1",
params![
id,
Expand All @@ -568,7 +589,8 @@ impl Db {
n.remark,
n.traffic_limit,
n.traffic_mode,
n.traffic_reset_day
n.traffic_reset_day,
n.notify
],
)?;
Ok(())
Expand All @@ -579,6 +601,11 @@ impl Db {
Ok(())
}

pub fn set_down_since(&self, id: i64, ts: i64) -> Result<()> {
self.conn().execute("UPDATE node SET down_since=?2 WHERE id=?1", params![id, ts])?;
Ok(())
}

pub fn reorder_nodes(&self, ids: &[i64]) -> Result<()> {
let unique: HashSet<_> = ids.iter().collect();
if unique.len() != ids.len() {
Expand Down Expand Up @@ -1526,6 +1553,8 @@ fn row_to_node(r: &rusqlite::Row<'_>) -> Node {
ipv6: s("ipv6"),
country: s("country"),
last_seen: n("last_seen"),
notify: n("notify") != 0,
down_since: n("down_since"),
token: s("token"),
}
}
Expand Down
37 changes: 31 additions & 6 deletions src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ mod api;
mod auth;
mod db;
mod frontend;
mod notify;

use std::collections::HashMap;
use std::net::{IpAddr, SocketAddr};
Expand Down Expand Up @@ -55,10 +56,12 @@ pub struct App {
pub site: String,
/// Parent directory containing one folder per installed public theme.
pub themes: PathBuf,
/// Alerts on their way out; see `notify::send`.
pub notes: tokio::sync::mpsc::Sender<notify::Note>,
}

impl App {
fn new(db: Db, site: String, themes: PathBuf) -> Self {
fn new(db: Db, site: String, themes: PathBuf, notes: tokio::sync::mpsc::Sender<notify::Note>) -> Self {
Self {
db,
agents: RwLock::default(),
Expand All @@ -71,12 +74,14 @@ impl App {
.expect("http client"),
site,
themes,
notes,
}
}

#[cfg(test)]
pub fn for_test(db: Db) -> Self {
Self::new(db, String::new(), PathBuf::from("themes"))
// Nothing delivers in tests; `notify::send` drops into the closed channel.
Self::new(db, String::new(), PathBuf::from("themes"), tokio::sync::mpsc::channel(1).0)
}

pub fn public_page(&self) -> bool {
Expand Down Expand Up @@ -324,7 +329,8 @@ async fn main() -> Result<()> {

let args = parse_args()?;
std::fs::create_dir_all(&args.themes)?;
let app = Arc::new(App::new(Db::open(&args.database)?, args.site.clone(), args.themes));
let (notes, inbox) = tokio::sync::mpsc::channel(notify::QUEUE);
let app = Arc::new(App::new(Db::open(&args.database)?, args.site.clone(), args.themes, notes));
let url = advertised_url(&args.site, args.listen);
first_run(&app, &url)?;
if exposed_over_plain_http(&url) {
Expand Down Expand Up @@ -369,6 +375,8 @@ async fn main() -> Result<()> {
}

tokio::spawn(housekeeping(app.clone()));
tokio::spawn(notify::deliver(app.clone(), inbox));
tokio::spawn(notify::watch(app.clone()));

let router = Router::new()
// Agents.
Expand Down Expand Up @@ -398,6 +406,7 @@ async fn main() -> Result<()> {
.route("/api/sessions", get(api::sessions))
.route("/api/sessions/{id}", delete(api::delete_session))
.route("/api/settings", get(api::settings).put(api::save_settings))
.route("/api/notify/test", post(notify::test))
.route("/api/themes", get(api::themes))
.route("/api/themes/{short}", delete(api::delete_theme))
.route("/api/themes/{short}/preview", get(api::theme_preview))
Expand Down Expand Up @@ -580,7 +589,9 @@ fn renew_online_nodes(app: &App) -> Result<()> {
// until 08:00 while the panel already shows it expired.
let today = Local::now().date_naive();
let online: Vec<i64> = app.agents.read().unwrap_or_else(|e| e.into_inner()).keys().copied().collect();
for node in app.db.nodes()? {
let nodes = app.db.nodes()?;
let mut rolled = Vec::new();
for node in &nodes {
if !online.contains(&node.id) {
continue;
}
Expand All @@ -590,11 +601,14 @@ fn renew_online_nodes(app: &App) -> Result<()> {
let Some(next) = renewed(expires, &node.billing_cycle, today) else { continue };
app.db.set_expiry(node.id, &next.to_string())?;
info!("node {} is still up past {expires}, expiry rolled to {next}", node.name);
rolled.push((node.name.as_str(), format!("{expires} → {next}")));
}
notify::renewed(app, rolled);
Ok(())
}

/// Expires sessions, trims history and rolls over expiry dates once an hour.
/// Expires sessions, trims history, rolls over expiry dates and sends the daily
/// expiry digest, once an hour.
async fn housekeeping(app: Shared) {
let mut ticker = tokio::time::interval(std::time::Duration::from_secs(3_600));
loop {
Expand All @@ -609,6 +623,12 @@ async fn housekeeping(app: Shared) {
if let Err(e) = renew_online_nodes(&app) {
warn!("rolling expiry dates failed: {e:#}");
}
// After the roll-over, so the digest lists dates as they now stand.
match notify::expiry_digest(&app, Local::now()) {
Ok(Some(note)) => notify::send(&app, note),
Ok(None) => {}
Err(e) => warn!("expiry digest failed: {e:#}"),
}
}
}

Expand All @@ -618,7 +638,12 @@ mod tests {
use axum::http::{StatusCode, Uri};

fn app(site: &str) -> App {
App::new(Db::open(":memory:").unwrap(), site.into(), PathBuf::from("themes"))
App::new(
Db::open(":memory:").unwrap(),
site.into(),
PathBuf::from("themes"),
tokio::sync::mpsc::channel(1).0,
)
}

/// A request as a reverse proxy would forward it, or as it arrives with none
Expand Down
Loading