Files
DCTS/crates/server/src/main.rs
T
fmq d16b3d3cdc feat(all): 数据库模块化拆分与版本化迁移、任务引擎命名体系收敛、物理输出校验加固与用户配置接通
- server/db: 拆 4929 行 db.rs 单体为 db/ 目录,migrations.rs 引入 PRAGMA user_version
    版本化迁移运行器(M1~M13)
  - 任务引擎 Phase 6/7b/7c 改名收敛:EngineStageConfig→PhaseConfig、StagePolicy→ResumePolicy、
    Converged→Completed、删除 task_type 列、success_method 拆 tlusty_/synspec_ 双列、
    新增 tlusty_status/synspec_status 半失败阶段守卫
  - 科学正确性加固:conv_check 任意行 NaN/Inf/溢出判无效(0 行容忍)、新增 spec_is_valid
    校验 SYNSPEC 脏谱、itek_history 逐次迭代全量保真、fmt_abn powf 溢出饱和
  - 用户配置真正接通:tlusty_chain/tlusty_input 由死字段经 调度器→TaskSpec→executor→runner
    透传生效;config 加载期 validate + deny_unknown_fields + 解析失败记 warn
  - 调度修复:H1 活锁(pending_strategies 跳过已失败策略)、种子查找错误不再静默降级冷启动
  - dashboard: 阶段配置面板 tlusty_stage/synspec_stage、"已完成"标签、迭代诊断展示
  - docs: 新增 database_refactor_design.md,同步 database/api/PIPELINE/workflow_detail
2026-08-06 20:51:21 +08:00

534 lines
24 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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<PathBuf>,
/// Listen port (overrides DCTS_PORT env var)
#[arg(short = 'p', long = "port")]
port: Option<u16>,
}
#[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<String, Vec<String>> =
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_dirDCTS_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))
// 历史种子导入同样上传 .7 大气文件,并入宽松 body limit / 并发限流组。
.route("/admin/import_seed", post(api::task::import_seed))
.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),
)
.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::<SocketAddr>(),
)
.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<axum::body::Body>,
next: axum::middleware::Next,
) -> axum::response::Response {
let mut resp = next.run(req).await;
let headers = resp.headers_mut();
// CSPdefault-src 'self';放行 Google Fontsindex.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
}