refactor(proxy-nodes): 简化批量升级为直接写入 upgrade_to 目标

- 后端移除分波 rollout/探测逻辑,改为一次性给所有合格 tunnel 节点写入升级目标,并自动取消活动中的 rollout
- 前端移除 rollout 进度、阶段筛选、重试/跳过/冲突清理等相关 UI 和交互
- 同步更新批量升级接口类型与测试用例
This commit is contained in:
fawney19
2026-04-19 13:38:59 +08:00
parent e6fd95453d
commit 7007d7a556
4 changed files with 140 additions and 949 deletions

View File

@@ -6,7 +6,7 @@ use crate::handlers::admin::shared::query_param_value;
use crate::maintenance::{
cancel_proxy_upgrade_rollout, clear_proxy_upgrade_rollout_conflicts,
restore_proxy_upgrade_rollout_skipped_nodes, retry_proxy_upgrade_rollout_node,
skip_proxy_upgrade_rollout_node, start_proxy_upgrade_rollout, ProxyUpgradeRolloutProbeConfig,
skip_proxy_upgrade_rollout_node, ProxyUpgradeRolloutProbeConfig,
};
use crate::GatewayError;
use aether_admin::system::{
@@ -130,6 +130,16 @@ struct ProxyNodeBatchUpgradeRequest {
probe_timeout_secs: Option<u64>,
}
#[derive(Debug, Default)]
struct ProxyNodeBatchUpgradeDispatchSummary {
version: String,
eligible_total: usize,
updated: usize,
skipped: usize,
node_ids: Vec<String>,
rollout_cancelled: bool,
}
const JSON_OBJECT_REQUIRED_DETAIL: &str = "请求体必须是合法的 JSON 对象";
const DEFAULT_PROXY_UPGRADE_BATCH_SIZE: usize = 1;
const DEFAULT_PROXY_UPGRADE_COOLDOWN_SECS: u64 = 60;
@@ -506,10 +516,7 @@ pub(crate) async fn maybe_build_local_admin_proxy_nodes_response(
if decision.route_kind.as_deref() == Some("batch_upgrade_nodes")
&& request_context.method() == http::Method::POST
{
if !state.has_proxy_node_reader()
|| !state.has_proxy_node_writer()
|| !state.app().data.has_system_config_store()
{
if !state.has_proxy_node_reader() || !state.has_proxy_node_writer() {
return Ok(Some(build_admin_proxy_nodes_data_unavailable_response()));
}
let input = match parse_json_body::<ProxyNodeBatchUpgradeRequest>(request_body) {
@@ -520,42 +527,16 @@ pub(crate) async fn maybe_build_local_admin_proxy_nodes_response(
Ok(version) => version,
Err(response) => return Ok(Some(response)),
};
let batch_size = match validate_batch_size(input.batch_size) {
Ok(batch_size) => batch_size,
Err(response) => return Ok(Some(response)),
};
let cooldown_secs = match validate_cooldown_secs(input.cooldown_secs) {
Ok(cooldown_secs) => cooldown_secs,
Err(response) => return Ok(Some(response)),
};
let probe =
match validate_probe_config(input.probe_url.as_deref(), input.probe_timeout_secs) {
Ok(probe) => probe,
Err(response) => return Ok(Some(response)),
};
let rollout = start_proxy_upgrade_rollout(
&state.app().data,
version.clone(),
batch_size,
cooldown_secs,
probe,
)
.await
.map_err(|err| GatewayError::Internal(err.to_string()))?;
let summary = dispatch_proxy_node_upgrade_targets(state, &version).await?;
return Ok(Some(
Json(json!({
"version": version,
"batch_size": rollout.batch_size,
"cooldown_secs": rollout.cooldown_secs,
"updated": rollout.updated,
"skipped": rollout.skipped,
"node_ids": rollout.node_ids,
"blocked": rollout.blocked,
"pending_node_ids": rollout.pending_node_ids,
"rollout_active": rollout.rollout_active,
"completed": rollout.completed,
"remaining": rollout.remaining,
"version": summary.version,
"eligible_total": summary.eligible_total,
"updated": summary.updated,
"skipped": summary.skipped,
"node_ids": summary.node_ids,
"rollout_cancelled": summary.rollout_cancelled,
}))
.into_response(),
));
@@ -1480,6 +1461,79 @@ fn admin_proxy_node_test_node_id_from_path(path: &str) -> Option<String> {
}
}
fn normalize_proxy_upgrade_version(value: &str) -> String {
value
.trim()
.strip_prefix("proxy-v")
.unwrap_or(value.trim())
.to_ascii_lowercase()
}
async fn dispatch_proxy_node_upgrade_targets(
state: &AdminAppState<'_>,
version: &str,
) -> Result<ProxyNodeBatchUpgradeDispatchSummary, GatewayError> {
let mut nodes = state.list_proxy_nodes().await?;
nodes.sort_by(|left, right| left.name.cmp(&right.name).then(left.id.cmp(&right.id)));
let rollout_cancelled = if state.app().data.has_system_config_store() {
cancel_proxy_upgrade_rollout(&state.app().data)
.await
.map_err(|err| GatewayError::Internal(err.to_string()))?
.is_some()
} else {
false
};
let normalized_target = normalize_proxy_upgrade_version(version);
let mut summary = ProxyNodeBatchUpgradeDispatchSummary {
version: version.to_string(),
rollout_cancelled,
..Default::default()
};
for node in nodes {
if node.is_manual || !node.tunnel_mode {
continue;
}
summary.eligible_total = summary.eligible_total.saturating_add(1);
let current_version =
aether_data::repository::proxy_nodes::proxy_reported_version(node.proxy_metadata.as_ref());
let pending_target = aether_data::repository::proxy_nodes::remote_config_upgrade_target(
node.remote_config.as_ref(),
);
if pending_target.as_deref() == Some(normalized_target.as_str())
|| current_version.as_deref() == Some(normalized_target.as_str())
{
continue;
}
let Some(updated) = state
.update_proxy_node_remote_config(
&aether_data::repository::proxy_nodes::ProxyNodeRemoteConfigMutation {
node_id: node.id.clone(),
node_name: None,
allowed_ports: None,
log_level: None,
heartbeat_interval: None,
scheduling_state: None,
upgrade_to: Some(Some(version.to_string())),
},
)
.await?
else {
continue;
};
summary.node_ids.push(updated.id);
}
summary.updated = summary.node_ids.len();
summary.skipped = summary.eligible_total.saturating_sub(summary.updated);
Ok(summary)
}
fn validate_batch_size(batch_size: Option<usize>) -> Result<usize, Response<Body>> {
let batch_size = batch_size.unwrap_or(DEFAULT_PROXY_UPGRADE_BATCH_SIZE);
if (1..=100).contains(&batch_size) {

View File

@@ -1252,7 +1252,7 @@ async fn gateway_handles_admin_proxy_node_events_locally_with_trusted_admin_prin
}
#[tokio::test]
async fn gateway_updates_proxy_node_config_and_batches_upgrade_locally() {
async fn gateway_updates_proxy_node_config_and_dispatches_upgrade_targets_locally() {
let upstream_hits = Arc::new(Mutex::new(0usize));
let upstream_hits_clone = Arc::clone(&upstream_hits);
let upstream = Router::new()
@@ -1346,15 +1346,14 @@ async fn gateway_updates_proxy_node_config_and_batches_upgrade_locally() {
.await
.expect("json body should parse");
assert_eq!(upgrade_payload["version"], "2.0.0");
assert_eq!(upgrade_payload["batch_size"], 1);
assert_eq!(upgrade_payload["updated"], 1);
assert_eq!(upgrade_payload["skipped"], 1);
assert_eq!(upgrade_payload["blocked"], false);
assert_eq!(upgrade_payload["pending_node_ids"], json!(["node-online"]));
assert_eq!(upgrade_payload["node_ids"], json!(["node-online"]));
assert_eq!(upgrade_payload["completed"], 0);
assert_eq!(upgrade_payload["remaining"], 2);
assert_eq!(upgrade_payload["rollout_active"], true);
assert_eq!(upgrade_payload["eligible_total"], 3);
assert_eq!(upgrade_payload["updated"], 3);
assert_eq!(upgrade_payload["skipped"], 0);
assert_eq!(
upgrade_payload["node_ids"],
json!(["node-online", "node-offline", "node-zeta"])
);
assert_eq!(upgrade_payload["rollout_cancelled"], false);
let blocked_upgrade_response = client
.post(format!("{gateway_url}/api/admin/proxy-nodes/upgrade"))
@@ -1372,11 +1371,7 @@ async fn gateway_updates_proxy_node_config_and_batches_upgrade_locally() {
.await
.expect("json body should parse");
assert_eq!(blocked_upgrade_payload["updated"], 0);
assert_eq!(blocked_upgrade_payload["blocked"], true);
assert_eq!(
blocked_upgrade_payload["pending_node_ids"],
json!(["node-online"])
);
assert_eq!(blocked_upgrade_payload["skipped"], 3);
let heartbeat_response = client
.post(format!("{gateway_url}/api/internal/tunnel/heartbeat"))
@@ -1404,7 +1399,7 @@ async fn gateway_updates_proxy_node_config_and_batches_upgrade_locally() {
assert_eq!(heartbeat_payload["remote_config"]["allowed_ports"][1], 8443);
assert_eq!(heartbeat_payload["remote_config"]["log_level"], "info");
let second_upgrade_response = client
let post_heartbeat_upgrade_response = client
.post(format!("{gateway_url}/api/admin/proxy-nodes/upgrade"))
.header(GATEWAY_HEADER, "rust-phase3b")
.header(TRUSTED_ADMIN_USER_ID_HEADER, "admin-user-123")
@@ -1414,57 +1409,14 @@ async fn gateway_updates_proxy_node_config_and_batches_upgrade_locally() {
.send()
.await
.expect("request should succeed");
assert_eq!(second_upgrade_response.status(), StatusCode::OK);
let second_upgrade_payload: serde_json::Value = second_upgrade_response
assert_eq!(post_heartbeat_upgrade_response.status(), StatusCode::OK);
let post_heartbeat_upgrade_payload: serde_json::Value = post_heartbeat_upgrade_response
.json()
.await
.expect("json body should parse");
assert_eq!(second_upgrade_payload["batch_size"], 1);
assert_eq!(second_upgrade_payload["updated"], 0);
assert_eq!(second_upgrade_payload["skipped"], 2);
assert_eq!(second_upgrade_payload["blocked"], true);
assert_eq!(
second_upgrade_payload["pending_node_ids"],
json!(["node-online"])
);
assert_eq!(second_upgrade_payload["node_ids"], json!([]));
assert_eq!(second_upgrade_payload["completed"], 0);
assert_eq!(second_upgrade_payload["remaining"], 3);
assert_eq!(second_upgrade_payload["rollout_active"], true);
assert!(
record_proxy_upgrade_traffic_success(&data_state, "node-online")
.await
.expect("traffic confirmation should be recorded")
);
let third_upgrade_response = client
.post(format!("{gateway_url}/api/admin/proxy-nodes/upgrade"))
.header(GATEWAY_HEADER, "rust-phase3b")
.header(TRUSTED_ADMIN_USER_ID_HEADER, "admin-user-123")
.header(TRUSTED_ADMIN_USER_ROLE_HEADER, "admin")
.header(TRUSTED_ADMIN_SESSION_ID_HEADER, "session-123")
.json(&json!({ "version": "2.0.0", "cooldown_secs": 0 }))
.send()
.await
.expect("request should succeed");
assert_eq!(third_upgrade_response.status(), StatusCode::OK);
let third_upgrade_payload: serde_json::Value = third_upgrade_response
.json()
.await
.expect("json body should parse");
assert_eq!(third_upgrade_payload["batch_size"], 1);
assert_eq!(third_upgrade_payload["updated"], 1);
assert_eq!(third_upgrade_payload["skipped"], 1);
assert_eq!(third_upgrade_payload["blocked"], false);
assert_eq!(
third_upgrade_payload["pending_node_ids"],
json!(["node-zeta"])
);
assert_eq!(third_upgrade_payload["node_ids"], json!(["node-zeta"]));
assert_eq!(third_upgrade_payload["completed"], 1);
assert_eq!(third_upgrade_payload["remaining"], 1);
assert_eq!(third_upgrade_payload["rollout_active"], true);
assert_eq!(post_heartbeat_upgrade_payload["node_ids"], json!([]));
assert_eq!(post_heartbeat_upgrade_payload["updated"], 0);
assert_eq!(post_heartbeat_upgrade_payload["skipped"], 3);
assert_eq!(*upstream_hits.lock().expect("mutex should lock"), 0);
@@ -1473,7 +1425,7 @@ async fn gateway_updates_proxy_node_config_and_batches_upgrade_locally() {
}
#[tokio::test]
async fn gateway_marks_draining_proxy_nodes_unschedulable_and_rollout_skips_them() {
async fn gateway_still_dispatches_upgrade_targets_to_draining_proxy_nodes() {
let mut alpha = sample_proxy_node("node-alpha");
alpha.name = "alpha".to_string();
alpha.status = "online".to_string();
@@ -1562,11 +1514,11 @@ async fn gateway_marks_draining_proxy_nodes_unschedulable_and_rollout_skips_them
.json()
.await
.expect("json body should parse");
assert_eq!(upgrade_payload["updated"], 1);
assert_eq!(upgrade_payload["node_ids"], json!(["node-zeta"]));
assert_eq!(upgrade_payload["pending_node_ids"], json!(["node-zeta"]));
assert_eq!(upgrade_payload["eligible_total"], 2);
assert_eq!(upgrade_payload["updated"], 2);
assert_eq!(upgrade_payload["node_ids"], json!(["node-alpha", "node-zeta"]));
let rollout_response = client
let list_response_after_upgrade = client
.get(format!(
"{gateway_url}/api/admin/proxy-nodes?skip=0&limit=10"
))
@@ -1577,11 +1529,17 @@ async fn gateway_marks_draining_proxy_nodes_unschedulable_and_rollout_skips_them
.send()
.await
.expect("request should succeed");
let rollout_payload: serde_json::Value = rollout_response
let list_payload_after_upgrade: serde_json::Value = list_response_after_upgrade
.json()
.await
.expect("json body should parse");
assert_eq!(rollout_payload["rollout"]["online_eligible_total"], 1);
let alpha_after_upgrade = list_payload_after_upgrade["items"]
.as_array()
.expect("items should be array")
.iter()
.find(|item| item["id"] == "node-alpha")
.expect("alpha should exist after upgrade");
assert_eq!(alpha_after_upgrade["remote_config"]["upgrade_to"], "2.0.0");
gateway_handle.abort();
}