DCTS/crates/node/src/worker.rs
Asfmq b91f1e4fa5 feat(server,dashboard): 引入多工作流数据隔离、安全中间件与前端 ESM 模块化重构
- server: 实现按 workflow_name 的多工作流数据隔离与旧数据库平滑迁移机制
- server: 新增 API Key 认证(auth)、限流中间件(rate_limit)与运维备份接口(admin)
- server: 统一 AppError 错误处理体系,重构调度器 scheduler 支持工作流级重置与抢占
- node: 节点 ID 缺失时自动生成随机 UUID,原生支持 `docker compose --scale node=N` 动态扩容
- dashboard: 前端模块化重构(state/api/components),升级 CSS 变量设计系统与 Toast 通知
- docker/docs: 更新 /healthz 健康检查、部署脚本 IP 配置及数据库设计文档
2026-07-28 21:54:02 +08:00

377 lines
15 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 crate::executor::execute_task;
use crate::reporter::report_result;
use anyhow::Result;
use common::config::NodeConfig;
use common::embedded::RuntimePaths;
use common::models::{NodeHeartbeatRequest, NodeRegisterRequest, TaskSpec};
use reqwest::Client;
use serde_json::Value;
use std::path::PathBuf;
use std::sync::atomic::{AtomicI32, Ordering};
use std::sync::Arc;
use tokio::time::{sleep, Duration};
use tracing::{info, warn};
pub struct NodeWorker {
config: NodeConfig,
client: Client,
runtime: RuntimePaths,
active_slots: Arc<AtomicI32>,
}
impl NodeWorker {
pub fn new(config: NodeConfig, runtime: RuntimePaths, client: Client) -> Self {
Self {
config,
client,
runtime,
active_slots: Arc::new(AtomicI32::new(0)),
}
}
/// 仅注册并领取专属 token供 main.rs 在本地无 token 时调用)。
/// 支持免凭据申请注册并轮询等待管理员在 Web Dashboard 上点击同意。
pub async fn register_and_fetch_token(
client: &Client,
server_url: &str,
node_id: &str,
) -> Result<String> {
info!(
"正在向服务端 {} 提交计算节点 {} 的注册申请...",
server_url, node_id
);
let req = NodeRegisterRequest {
node_id: node_id.to_string(),
host_name: gethostname::gethostname().to_string_lossy().to_string(),
max_slots: 0,
};
let resp = client
.post(format!("{}/api/node/register", server_url))
.json(&req)
.send()
.await?;
if !resp.status().is_success() {
anyhow::bail!("向服务端提交注册申请失败HTTP 状态码: {}", resp.status());
}
let json: Value = resp.json().await?;
let status = json.get("status").and_then(|v| v.as_str()).unwrap_or("");
if status == "approved" {
if let Some(t) = json.get("node_token").and_then(|v| v.as_str()) {
return Ok(t.to_string());
}
}
info!(
"⏳ 节点 {} 的注册申请已提交!等待管理员在管理 Dashboard 上点击【同意接入】...",
node_id
);
// 轮询等待管理员在 Dashboard 上的 Approve
loop {
sleep(Duration::from_secs(5)).await;
let check_req = serde_json::json!({ "node_id": node_id });
let resp = match client
.post(format!("{}/api/node/check_status", server_url))
.json(&check_req)
.send()
.await
{
Ok(r) => r,
Err(e) => {
warn!("轮询节点审批状态网络异常: {}", e);
continue;
}
};
if !resp.status().is_success() {
continue;
}
let body: Value = match resp.json().await {
Ok(b) => b,
Err(_) => continue,
};
let check_status = body.get("status").and_then(|v| v.as_str()).unwrap_or("");
if check_status == "approved" {
if let Some(token) = body.get("node_token").and_then(|v| v.as_str()) {
info!(
"🎉 节点 {} 已成功获取管理员授权!专属访问 Token 接收完成。",
node_id
);
return Ok(token.to_string());
}
} else if check_status == "rejected" {
anyhow::bail!("节点 {} 的注册申请已被管理员拒绝或清理", node_id);
}
}
}
/// 正式注册(带真实 slot 数),供 run() 启动时刷新节点信息用。
pub async fn register(&self) -> Result<()> {
info!(
"正在向服务端 {} 刷新节点 {} 注册信息...",
self.config.server_url, self.config.node_id
);
let req = NodeRegisterRequest {
node_id: self.config.node_id.clone(),
host_name: gethostname::gethostname().to_string_lossy().to_string(),
max_slots: self.config.max_slots as i32,
};
let resp = self
.client
.post(format!("{}/api/node/register", self.config.server_url))
.json(&req)
.send()
.await?;
if !resp.status().is_success() {
anyhow::bail!("向服务端注册节点失败HTTP 状态码: {}", resp.status());
}
Ok(())
}
pub async fn run(&self) -> Result<()> {
self.register().await?;
info!(
"计算节点已激活,最大并行 Slot 槽位数: {}",
self.config.max_slots
);
// Start background heartbeat loop
let hb_client = self.client.clone();
let hb_url = format!("{}/api/node/heartbeat", self.config.server_url);
let hb_node_id = self.config.node_id.clone();
let hb_slots = self.active_slots.clone();
let hb_interval = self.config.heartbeat_sec;
tokio::spawn(async move {
let sys_arc = std::sync::Arc::new(std::sync::Mutex::new(sysinfo::System::new_all()));
{
let s = sys_arc.clone();
let _ = tokio::task::spawn_blocking(move || {
if let Ok(mut sys) = s.lock() {
sys.refresh_cpu();
}
})
.await;
}
sleep(Duration::from_millis(200)).await;
{
let s = sys_arc.clone();
let _ = tokio::task::spawn_blocking(move || {
if let Ok(mut sys) = s.lock() {
sys.refresh_cpu();
}
})
.await;
}
loop {
sleep(Duration::from_secs(hb_interval)).await;
let s = sys_arc.clone();
let (cpu_usage, memory_usage) = tokio::task::spawn_blocking(move || {
let mut sys = match s.lock() {
Ok(guard) => guard,
Err(_) => return (0.0, 0.0),
};
sys.refresh_cpu();
sys.refresh_memory();
let cpu_usage = sys.global_cpu_info().cpu_usage();
let mem_total = sys.total_memory() as f32;
let mem_used = sys.used_memory() as f32;
let memory_usage = if mem_total > 0.0 {
(mem_used / mem_total) * 100.0
} else {
0.0
};
(cpu_usage, memory_usage)
})
.await
.unwrap_or((0.0, 0.0));
let active = hb_slots.load(Ordering::Acquire);
let req = NodeHeartbeatRequest {
node_id: hb_node_id.clone(),
active_slots: active,
cpu_usage,
memory_usage,
};
match hb_client.post(&hb_url).json(&req).send().await {
Ok(resp) => {
let status = resp.status();
// 401/403token 失效或被吊销。与 claim_task 口径统一:直接退出进程,
// 避免心跳线程持续发被拒请求刷日志、占用服务端限流计数。心跳通常比
// claim 更高频,往往先于 claim_task 发现吊销。
if status.as_u16() == 401 || status.as_u16() == 403 {
tracing::error!(
"节点 {} 心跳被服务端拒绝 (HTTP {})node token 已失效或被吊销。请清理 .node_token 文件后重启节点以重新发起注册审批。进程将退出,依赖编排系统重启。",
hb_node_id, status
);
std::process::exit(1);
}
}
Err(e) => warn!("节点 {} 心跳上报失败: {}", hb_node_id, e),
}
}
});
let work_dir = PathBuf::from(&self.config.work_dir);
tokio::fs::create_dir_all(&work_dir).await?;
let shutting_down = Arc::new(std::sync::atomic::AtomicBool::new(false));
let shutdown_signal = shutting_down.clone();
tokio::spawn(async move {
if tokio::signal::ctrl_c().await.is_ok() {
info!("收到 Ctrl+C 终止信号,停止领用新任务,准备优雅退出 (再次按 Ctrl+C 可强制立即退出)...");
shutdown_signal.store(true, Ordering::Release);
// 二次 Ctrl+C 强行立即退出
if tokio::signal::ctrl_c().await.is_ok() {
warn!("再次收到 Ctrl+C 终止信号,强行立即中断退出!");
std::process::exit(130);
}
}
});
let mut was_disconnected = false;
// 带有优雅退出信号响应的任务领用主循环
loop {
if shutting_down.load(Ordering::Acquire) {
break;
}
let active = self.active_slots.load(Ordering::Acquire);
if (active as usize) < self.config.max_slots {
match self.claim_task().await {
Ok(Some(task)) => {
if was_disconnected {
info!("与服务端恢复网络连接,已自动重新上线并开始领用计算任务!");
was_disconnected = false;
}
self.active_slots.fetch_add(1, Ordering::AcqRel);
let client = self.client.clone();
let server_url = self.config.server_url.clone();
let node_id = self.config.node_id.clone();
let runtime = self.runtime.clone();
let work_dir = work_dir.clone();
let slots_counter = self.active_slots.clone();
tokio::spawn(async move {
let slot_work_dir = work_dir.join(format!("task_{}", task.task_id));
let res =
execute_task(&client, &server_url, &runtime, &work_dir, &task)
.await
.map_err(|e| e.to_string());
let report_res =
report_result(&client, &server_url, &node_id, &task, res).await;
if report_res.is_ok() {
if let Err(e) =
crate::executor::cleanup_slot_work_dir(&slot_work_dir).await
{
warn!(
"清理任务 {} 的沙盒目录 {} 失败: {}",
task.task_id,
slot_work_dir.display(),
e
);
}
} else if let Err(ref e) = report_res {
warn!("向服务端上报任务 {} 计算结果失败: {}", task.task_id, e);
}
slots_counter.fetch_sub(1, Ordering::AcqRel);
});
}
Ok(None) => {
if was_disconnected {
info!("与服务端恢复网络连接,已自动重新上线 (当前暂无排队任务)。");
was_disconnected = false;
}
sleep(Duration::from_secs(5)).await;
}
Err(e) => {
was_disconnected = true;
warn!("向服务端请求领用计算任务时出错: {}", e);
sleep(Duration::from_secs(10)).await;
}
}
} else {
sleep(Duration::from_secs(2)).await;
}
}
// 等待在途任务完结(最多等待 30 秒)
if self.active_slots.load(Ordering::Acquire) > 0 {
info!(
"正在等待 {} 个在途计算任务优雅完结 (上限 30 秒,按二次 Ctrl+C 可强行中断)...",
self.active_slots.load(Ordering::Acquire)
);
}
let start_wait = std::time::Instant::now();
let mut last_log_time = std::time::Instant::now();
while self.active_slots.load(Ordering::Acquire) > 0 {
if start_wait.elapsed().as_secs() >= 30 {
warn!("在途任务等待超时 (30s),强制退出节点");
break;
}
if last_log_time.elapsed().as_secs() >= 5 {
info!(
"仍在等待 {} 个在途计算任务完结...",
self.active_slots.load(Ordering::Acquire)
);
last_log_time = std::time::Instant::now();
}
sleep(Duration::from_millis(500)).await;
}
info!("DCTS 计算节点安全退出。");
Ok(())
}
async fn claim_task(&self) -> Result<Option<TaskSpec>> {
let claim_url = format!("{}/api/task/claim", self.config.server_url);
let resp = self.client.post(&claim_url).send().await?;
let status = resp.status();
// 401/403 表明 node token 已被吊销或失效(区别于「暂无任务」与服务端 5xx 故障)。
// 服务端故障返回 5xx 会走 !is_success() 的 Ok(None) 分支,仅在网络层/鉴权层拒绝时
// 才是真正的吊销。此时继续轮询只会持续产生被拒请求并刷日志,故直接退出进程,
// 由编排系统Docker restart / systemd / k8s拉起新进程发现 .node_token 失效后
// 会自动走注册审批流程重新申请。
if status.as_u16() == 401 || status.as_u16() == 403 {
tracing::error!(
"领用任务被服务端拒绝 (HTTP {})node token 已失效或被吊销。请清理 .node_token 文件后重启节点以重新发起注册审批。进程将退出,依赖编排系统重启。",
status
);
std::process::exit(1);
}
if !status.is_success() {
return Ok(None);
}
let json: Value = resp.json().await?;
if json["status"] == "ok" && !json["task"].is_null() {
let task: TaskSpec = serde_json::from_value(json["task"].clone())?;
Ok(Some(task))
} else {
Ok(None)
}
}
}