use anyhow::Result; use server::api::{self, AppState}; use server::db::Database; use server::scheduler::GridScheduler; use axum::{ extract::DefaultBodyLimit, http::HeaderValue, routing::{get, post}, Router, }; use clap::Parser; use common::config::{GridConfig, ServerConfig}; use common::logging::init_logging; use mq::sqlite_queue::SqliteTaskQueue; use std::net::SocketAddr; use std::path::{Path, PathBuf}; use std::sync::Arc; use tokio::time::{sleep, Duration}; use tower_http::services::{ServeDir, ServeFile}; use tracing::info; #[derive(Parser, Debug)] #[command( name = "server", version = "0.1.0", about = "Distributed Computing TLUSTY/SYNSPEC (DCTS) Server" )] struct CliArgs { /// Optional path to workflow configuration YAML file to auto-register on startup #[arg(short = 'w', long = "workflow")] workflow: Option, /// Listen port (overrides DCTS_PORT env var) #[arg(short = 'p', long = "port")] port: Option, } #[tokio::main] async fn main() -> Result<()> { dotenvy::dotenv().ok(); let _logging_guards = init_logging("server", "info,server=debug")?; let cli = CliArgs::parse(); info!("启动 DCTS 服务端 (Distributed Computing TLUSTY/SYNSPEC Server)..."); let mut server_cfg = ServerConfig::default(); if let Some(port) = cli.port { server_cfg.bind_addr = format!("0.0.0.0:{}", port); } if let Some(wf) = cli.workflow { server_cfg.grid_config = wf.to_string_lossy().to_string(); } let db = Database::new(&server_cfg.db_path).await?; let queue = Arc::new(SqliteTaskQueue::new(&server_cfg.queue_db_path).await?); let scheduler = Arc::new(GridScheduler::new(db.clone(), queue.clone())); // Auto-register sdB_cno.yaml if exists and not yet in DB. // 仅在 DB 中尚无该工作流时注册(INSERT),绝不覆盖已存在的配置—— // 旧实现用 upsert 每次启动都用文件内容覆盖 config_yaml/description/status, // 导致管理员通过 API 编辑过的配置在重启后被静默回退。 let default_wf_path = Path::new(&server_cfg.grid_config); if default_wf_path.is_file() { if let Ok(yaml_content) = std::fs::read_to_string(default_wf_path) { let already_exists = matches!(db.get_workflow("sdB_cno").await, Ok(Some(_))); if !already_exists { if let Err(e) = db .upsert_workflow( "sdB_cno", Some("sdB CNO 6D Stellar Atmosphere Grid"), &yaml_content, "idle", ) .await { tracing::warn!("预注册默认工作流失败: {}", e); } else { info!("已在数据库中成功预注册默认工作流 'sdB_cno'"); } } else { info!("默认工作流 'sdB_cno' 已存在于数据库,保留现有配置(不覆盖 API 编辑)"); } } } // 弱口令凭据安全警告检测:仅按强度阈值判断(短于 16 字节视为弱口令)。 // 推荐用 `openssl rand -hex 32`(64 字符)生成。 let is_weak_token = |t: Option<&str>| -> bool { t.map(|s| s.len() < 16).unwrap_or(false) }; if is_weak_token(server_cfg.admin_token.as_deref()) { tracing::warn!("⚠️ 检测到系统当前正在使用弱口令凭据或默认 Token!建议生产环境在 .env 中配置使用 openssl rand -hex 32 生成的高强度 Token!"); } // 启动恢复:把卡在 `initializing` 态的工作流重新初始化。 // 背景:`start_workflow` 把状态切到 `initializing` 后在后台 spawn `initialize_grid`, // 若进程在初始化中途崩溃/重启,工作流会永久卡在 `initializing`——`get_running_workflow_names` // 仍把它视为可调度,但网格展开未完成,导致半初始化网格被调度。 // `initialize_grid` 是幂等的(upsert ON CONFLICT DO NOTHING),重跑可补齐缺失点并把 // 状态推进到 `running`。失败则回退为 `idle` 等待人工重启(与 start_workflow 口径一致)。 match db.get_initializing_workflows().await { Ok(stuck) if !stuck.is_empty() => { info!( "检测到 {} 个卡在 initializing 态的工作流(上次启动未完成即重启),开始重新初始化...", stuck.len() ); for (wf_name, wf_yaml) in &stuck { match GridConfig::from_yaml_str(wf_yaml) { Ok(cfg) => match scheduler.initialize_grid(&cfg, wf_name).await { Ok(_) => { let _ = db.update_workflow_status(wf_name, "running").await; info!( "启动恢复:工作流 {} 已完成重新初始化并切回 running", wf_name ); } Err(e) => { tracing::warn!( "启动恢复:工作流 {} 重新初始化失败,回退为 idle: {}", wf_name, e ); let _ = db.update_workflow_status(wf_name, "idle").await; } }, Err(e) => { tracing::warn!( "启动恢复:工作流 {} 的 YAML 配置解析失败,回退为 idle: {}", wf_name, e ); let _ = db.update_workflow_status(wf_name, "idle").await; } } } } Ok(_) => {} Err(e) => tracing::warn!("启动恢复:查询 initializing 工作流失败: {}", e), } let rate_limiter = api::rate_limit::RateLimiter::new(5, std::time::Duration::from_secs(300)); let state = AppState { db, queue: queue.clone(), scheduler: scheduler.clone(), seeds_dir: server_cfg.seeds_dir.clone(), rate_limiter, admin_token: server_cfg.admin_token.clone(), auth_disabled: server_cfg.auth_disabled, admin_sessions: std::sync::Arc::new(tokio::sync::RwLock::new( std::collections::HashMap::new(), )), }; // Background maintenance & scheduling with Exponential Backoff let bg_db = state.db.clone(); let bg_queue = queue.clone(); let bg_scheduler = scheduler.clone(); let stale_sec = server_cfg.stale_sec; let node_stale_sec = server_cfg.node_stale_sec; tokio::spawn(async move { let mut fail_count: u32 = 0; let mut first_run = true; loop { if first_run { first_run = false; } else { let base_delay = 30u64; let current_delay = if fail_count == 0 { base_delay } else { (base_delay * (1u64 << fail_count.min(4))).min(300) }; sleep(Duration::from_secs(current_delay)).await; } let bg_db_clone = bg_db.clone(); let bg_queue_clone = bg_queue.clone(); let bg_scheduler_clone = bg_scheduler.clone(); let join_handle = tokio::spawn(async move { let mut has_error = false; match bg_queue_clone.requeue_stale_tasks(stale_sec).await { Ok(requeued) => { if !requeued.is_empty() { info!("重新将 {} 个超时/掉线任务放回待计算队列", requeued.len()); let mut by_wf: std::collections::HashMap> = std::collections::HashMap::new(); for (point, wf) in &requeued { // 旧版任务(payload 无 workflow_name)归一到 '__legacy__', // 使 reset_specific_grid_points_to_pending 命中 legacy 网格点 // (H1 修复:否则按 '' 更新 0 行,重投后的 legacy 点永远重置不回 pending)。 by_wf .entry(server::db::normalize_workflow_name(wf.as_deref())) .or_default() .push(point.clone()); } for (wf, points) in by_wf { let _ = bg_db_clone .reset_specific_grid_points_to_pending(&points, &wf) .await; } } } Err(e) => { tracing::warn!("重投超时任务失败: {}", e); has_error = true; } } match bg_db_clone.mark_stale_nodes_offline(node_stale_sec).await { Ok(offline) => { if offline > 0 { info!("已标记 {} 个心跳超时的计算节点为离线状态", offline); } } Err(e) => { tracing::warn!("标记超时节点离线失败: {}", e); has_error = true; } } // 孤儿点回收(#6 修复兜底,2026-08-02 涡旋事故重构):queue 凭证已消失 // (误删/崩溃丢队列/insert 后 push 前崩溃)但 grid_points 仍卡在 // running/queued 的点,requeue_stale_tasks 找不到它们。经 MQ 活性交叉 // 校验确认真孤儿后重置为 pending 让调度器重新派发,并清除作为判据的 // 僵尸 tasks 行(旧实现仅凭 tasks 表 stale pending 行判定,僵尸行使 // 判据恒真 → 重复派发涡旋,已废弃)。 // // 审查修复 #M4:reclaim 用独立且更大的阈值(2 * stale_sec)。reclaim 的语义是 // 「孤儿回收」(队列凭证完全丢失),时间尺度应比 requeue 的「claim 超时重投」 // 更宽松:刚被 claim 的任务 tasks 行仍 pending,过小阈值会把它误判孤儿候选、 // 在 requeue 把队列行打回 pending 到节点重新 claim 的窗口内增加抖动。2x 给 // 正常长任务足够缓冲。 let reclaim_threshold = stale_sec.saturating_mul(2); match bg_scheduler_clone .reclaim_orphaned_points(reclaim_threshold) .await { Ok(reset) => { if reset > 0 { info!( "已回收 {} 个孤儿网格点(领用凭证丢失,重置为 pending)", reset ); } } Err(e) => { tracing::warn!("回收孤儿网格点失败: {}", e); has_error = true; } } if let Err(e) = bg_scheduler_clone.schedule_pending_tasks().await { tracing::warn!("后台定时性任务调度检测失败: {}", e); has_error = true; } if let Err(e) = bg_db_clone.sync_all_running_workflows_completion().await { tracing::warn!("后台同步已完成工作流状态失败: {}", e); has_error = true; } // P3 进度快照:对每个运行中工作流记录计数(record_progress_snapshot // 内部去重——计数无变化不落库);顺带清理超过 7 天的旧快照。 // 观测性写入失败不回退调度退避(不置 has_error)。 match bg_db_clone.get_running_workflow_names().await { Ok(names) => { for wf in names { if let Err(e) = bg_db_clone.record_progress_snapshot(&wf).await { tracing::warn!("记录工作流 {} 进度快照失败: {}", wf, e); } } } Err(e) => tracing::warn!("获取运行中工作流列表失败: {}", e), } if let Err(e) = bg_db_clone.purge_progress_snapshots(7).await { tracing::warn!("清理过期进度快照失败: {}", e); } has_error }); match join_handle.await { Ok(has_error) => { if has_error { fail_count = fail_count.saturating_add(1); } else { fail_count = 0; } } Err(e) => { tracing::error!("后台维护任务内部发生 Panic: {:?}", e); fail_count = fail_count.saturating_add(1); } } } }); // 每天自动触发一次数据库备份。 // 备份目录跟随 server_cfg.backup_dir(DCTS_BACKUP_DIR,默认 data/backups), // 与 DB_PATH 解耦,避免 DB 卷与备份卷不一致时备份落到未持久化层。 // 首次延迟 1 小时,避免频繁重启(如调试阶段)短时间堆积备份文件;backup_database // 自身还带有 7 天保留期清理兜底。 let backup_db = state.db.clone(); let backup_dir = server_cfg.backup_dir.clone(); tokio::spawn(async move { sleep(Duration::from_secs(3600)).await; loop { if let Err(e) = backup_db.backup_database(&backup_dir).await { tracing::warn!("自动备份数据库失败: {}", e); } sleep(Duration::from_secs(24 * 3600)).await; } }); // 大体积上传端点单独拎出,套用更宽松的 body limit(256MB,覆盖收敛种子 .7 文件量级) // 并限制并发数:每个 report 请求最多 256MB 驻留内存,无并发上限时 N 个请求可耗尽内存。 // 限流后超出并发数的请求排队等待(而非直接拒绝),保证正常业务不被误伤。 const REPORT_BODY_LIMIT: usize = 256 * 1024 * 1024; const DEFAULT_BODY_LIMIT: usize = 10 * 1024 * 1024; const REPORT_MAX_CONCURRENCY: usize = 4; let report_router = Router::new() .route("/task/report", post(api::task::report_task)) .layer(DefaultBodyLimit::max(REPORT_BODY_LIMIT)) .layer(tower::ServiceBuilder::new().concurrency_limit(REPORT_MAX_CONCURRENCY)); // 节点注册接口独立 IP 限流保护(每分钟最多 10 次申请,无论成败都计数,防恶意频繁注册) // 使用 new_count_all:此 limiter 专挂 /node/register,对注册路径的所有响应计入窗口。 // 通用 API 限流器(见下方 auth_enabled 分支)用 new 构造(count_all=false),不会因 // 成功注册把 IP 锁出整个 /api/*,避免跨端点连锁限流。 let register_limiter = api::rate_limit::RateLimiter::new_count_all(10, std::time::Duration::from_secs(60)); let register_rate_limit_layer = axum::middleware::from_fn_with_state( register_limiter, api::rate_limit::rate_limit_middleware, ); // 其余 API(小体积)套用 10MB 默认上限,防止大文件内存耗尽 DoS。 // 注意:body limit layer 从外到内执行、先接触原始 body 流的层先生效。 // 必须把 10MB 限制只套在"小体积子 router"上,再与 report_router 合并, // 合并后的外层不能再套任何全局 limit —— 否则外层 10MB 会截断 report 的 256MB body 流, // 导致收敛种子 .7(常 >10MB)上报被 413 拒绝、结果反复重算。 let small_body_router = Router::new() // Auth API .route("/login", post(api::auth::login)) .route("/auth/check", get(api::auth::check_auth)) .route("/auth/logout", post(api::auth::logout)) // Core Node & Task API .route( "/node/register", post(api::node::register_node).layer(register_rate_limit_layer), ) .route("/node/check_status", post(api::node::check_node_status)) .route("/node/heartbeat", post(api::node::heartbeat_node)) .route("/task/claim", post(api::task::claim_task)) .route("/seed/:name", get(api::seed::download_seed)) .route("/status", get(api::status::get_status)) // Static Data API .route( "/data/file/*filename", get(api::data::download_single_data_file), ) .route("/data/linelist", get(api::data::download_linelist)) // Workflow Management CRUD API .route( "/workflows", get(api::workflow::list_workflows).post(api::workflow::save_workflow), ) .route( "/workflows/:name", get(api::workflow::get_workflow) .put(api::workflow::save_workflow) .delete(api::workflow::delete_workflow), ) .route( "/workflows/:name/start", post(api::workflow::start_workflow), ) .route("/workflows/:name/stop", post(api::workflow::stop_workflow)) // 工作流执行观测 API(进度统计 / 逐点明细 / 单点诊断,均要求 Admin 角色) .route( "/workflows/:name/stats", get(api::workflow::get_workflow_stats), ) .route( "/workflows/:name/progress", get(api::workflow::get_workflow_progress), ) .route( "/workflows/:name/points", get(api::workflow::get_workflow_points), ) .route( "/workflows/:name/points/:point", get(api::workflow::get_workflow_point_detail), ) // Admin Management API(节点凭据查看/审批/重发/停用/启用,均要求 Admin 角色) .route("/admin/nodes", get(api::admin::list_nodes)) .route( "/admin/nodes/:node_id/approve", post(api::admin::approve_node), ) .route( "/admin/nodes/:node_id/reject", post(api::admin::reject_node), ) .route( "/admin/nodes/:node_id/reissue", post(api::admin::reissue_node), ) .route( "/admin/nodes/:node_id/disable", post(api::admin::disable_node), ) .route( "/admin/nodes/:node_id/enable", post(api::admin::enable_node), ) .route( "/admin/nodes/:node_id/quota", post(api::admin::set_node_quota), ) // 临时迁移端点:扫 seeds_dir 的 conv.json → 写 grid_points.summary_json。 // 迁移完成后删除本路由 + api/migrate.rs 即可。 .route("/admin/migrate_conv", post(api::migrate::migrate_conv)) .layer(DefaultBodyLimit::max(DEFAULT_BODY_LIMIT)); // 合并两个子 router:各自携带自己的 body limit,互不覆盖。 let api_router = small_body_router.merge(report_router); // 鉴权策略:fail-closed。 // - 配置了 DCTS_ADMIN_TOKEN → 启用完整鉴权。 // - 显式 DCTS_AUTH_DISABLE=1 → 无鉴权(仅本地调试,需运维主动声明承担风险)。 // - 既未配置 token、又未显式 disable → **拒绝启动**。 // 避免 .env 缺失/变量名拼错/容器未注入环境变量时服务静默退化为完全无鉴权裸奔。 let auth_enabled = !state.auth_disabled && state.admin_token.is_some(); if !auth_enabled && !state.auth_disabled { anyhow::bail!( "拒绝启动:未配置 DCTS_ADMIN_TOKEN 且未显式设置 DCTS_AUTH_DISABLE=1。\ 生产部署必须在 .env 中配置 DCTS_ADMIN_TOKEN;若确为本地调试,\ 请显式设置 DCTS_AUTH_DISABLE=1 以承担无鉴权风险。" ); } let api_router = if auth_enabled { info!("已启用 API 身份鉴权保护(Admin 端点需 admin token 验证;Node 节点免 Token 提交申请,经 Dashboard 管理员审批授权下发)"); // 鉴权失败限流(防 token 在线暴力):外层先判 IP 限流,内层再做鉴权。 // 限流状态为 20 次/分钟(按 IP),超阈值返回 429。 let limiter = api::rate_limit::RateLimiter::new(20, std::time::Duration::from_secs(60)); let rate_limit_layer = axum::middleware::from_fn_with_state(limiter, api::rate_limit::rate_limit_middleware); let auth_layer = axum::middleware::from_fn_with_state(state.clone(), api::auth_middleware); api_router.layer(auth_layer).layer(rate_limit_layer) } else { info!("DCTS_AUTH_DISABLE=1 已生效:服务端运行在无鉴权模式(仅限本地调试,切勿用于生产)。"); api_router }; // Host Dashboard SPA static files from dashboard/dist if directory exists or fallback to index.html let serve_dir = ServeDir::new("dashboard/dist").fallback(ServeFile::new("dashboard/dist/index.html")); // 安全响应头(CSP / nosniff / DENY / Referrer-Policy)。 let security_headers = axum::middleware::from_fn(security_headers_middleware); let app = Router::new() // 独立健康检查端点:不走鉴权、不走 CORS/body 限制,专供 docker healthcheck 与外部监控探测。 // 开启鉴权后 /api/status 会返回 401,导致容器被判定不健康而反复重启,故单独提供 /healthz。 .route("/healthz", get(api::status::healthz)) .nest("/api", api_router) .layer(server::cors::build_cors_layer()) .layer(security_headers) .fallback_service(serve_dir) .with_state(state); let addr: SocketAddr = server_cfg.bind_addr.parse()?; info!("DCTS 服务端已在 http://{} 启动监听", addr); let listener = tokio::net::TcpListener::bind(addr).await?; // into_make_service_with_connect_info:让限流中间件能从连接拿到客户端 IP(反代场景则用 X-Forwarded-For) axum::serve( listener, app.into_make_service_with_connect_info::(), ) .with_graceful_shutdown(async { let _ = tokio::signal::ctrl_c().await; info!("收到 Ctrl+C 终止信号,DCTS 服务端准备优雅关闭..."); }) .await?; info!("DCTS 服务端已安全关闭。"); Ok(()) } /// 注入安全响应头的中间件函数。 async fn security_headers_middleware( req: axum::http::Request, next: axum::middleware::Next, ) -> axum::response::Response { let mut resp = next.run(req).await; let headers = resp.headers_mut(); // CSP:default-src 'self';放行 Google Fonts(index.html 引用);允许 data: 图片。 // 已移除 'unsafe-eval'(dashboard 不用 eval/new Function)与 script-src 'unsafe-inline' // (主题初始化脚本已移为外链 public/theme-init.js,见 dashboard/index.html)。无内联 // 脚本/内联事件处理器,外链脚本走 'self' 即可,脚本注入防御最大化。 // style-src 保留 'unsafe-inline':index.html 含内联 style 属性(SVG flex 布局等),移除会破坏渲染。 headers .entry(axum::http::header::CONTENT_SECURITY_POLICY) .or_insert_with(|| { HeaderValue::from_static( "default-src 'self'; script-src 'self'; \ style-src 'self' 'unsafe-inline' https://fonts.googleapis.com; \ font-src 'self' data: https://fonts.gstatic.com; \ connect-src 'self'; img-src 'self' data: blob:; \ frame-ancestors 'none'", ) }); headers .entry(axum::http::header::X_CONTENT_TYPE_OPTIONS) .or_insert_with(|| HeaderValue::from_static("nosniff")); headers .entry(axum::http::header::X_FRAME_OPTIONS) .or_insert_with(|| HeaderValue::from_static("DENY")); headers .entry(axum::http::HeaderName::from_static("referrer-policy")) .or_insert_with(|| HeaderValue::from_static("strict-origin-when-cross-origin")); resp }