diff --git a/Cargo.lock b/Cargo.lock index acab1463..8f99dd62 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -965,7 +965,7 @@ checksum = "c8d4a3bb8b1e0c1050499d1815f5ab16d04f0959b233085fb31653fbfc9d98f9" [[package]] name = "cluster_mgr" -version = "1.11.1" +version = "1.12.2" dependencies = [ "anyhow", "async-trait", diff --git a/src/cluster_mgr/src/cli/mod.rs b/src/cluster_mgr/src/cli/mod.rs index d826f586..f44bc40e 100644 --- a/src/cluster_mgr/src/cli/mod.rs +++ b/src/cluster_mgr/src/cli/mod.rs @@ -327,6 +327,12 @@ Behavior:\n\ default_value_t = false )] force: bool, + #[arg( + long, + default_value_t = false, + help = "Skip stopping and restarting the standalone log service during rolling update" + )] + skip_log_restart: bool, }, #[command( diff --git a/src/cluster_mgr/src/cli/task/db_update_task.rs b/src/cluster_mgr/src/cli/task/db_update_task.rs index ec22d0d0..6a6f43d2 100644 --- a/src/cluster_mgr/src/cli/task/db_update_task.rs +++ b/src/cluster_mgr/src/cli/task/db_update_task.rs @@ -281,6 +281,7 @@ impl TaskExecutor for DbDeploymentUpdateTask { nodes: ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), }, cluster_config: Some(content), } diff --git a/src/cluster_mgr/src/cli/task/eloq_tx_ctl_task.rs b/src/cluster_mgr/src/cli/task/eloq_tx_ctl_task.rs index 3244e846..7fa4a134 100644 --- a/src/cluster_mgr/src/cli/task/eloq_tx_ctl_task.rs +++ b/src/cluster_mgr/src/cli/task/eloq_tx_ctl_task.rs @@ -216,9 +216,8 @@ impl RedisProbe { } } Err(err) => { - if err.is_connection_refusal() { - maybe_continue_probe!(wait_secs); - } + info!("Redis connection to {} failed: {}, retrying", url, err); + maybe_continue_probe!(wait_secs); return Ok(HashMap::from([ ( CMD.to_string(), @@ -730,6 +729,7 @@ impl TaskExecutor for EloqTxCtlTask { } "stop" | "force_stop" => { let stop_cmd = self.ctl_cmd.cmd_value(); + let has_topology_filter = self.receiver.is_some(); let mut target_ports: Vec = Vec::new(); match server_type { "txservice" => { @@ -766,10 +766,28 @@ impl TaskExecutor for EloqTxCtlTask { unreachable!("Unknown server type: {}", server_type); } } - if target_ports.is_empty() { + if target_ports.is_empty() && !has_topology_filter { target_ports.push(port.to_string()); } + if target_ports.is_empty() { + info!( + "No matching {} ports found for host {} in current topology; skipping {}.", + server_type, self.task_id.host, ctl_cmd_ref + ); + return Ok(Some(HashMap::from([ + (CMD.to_string(), TaskArgValue::Str(stop_cmd)), + (CMD_STATUS.to_string(), TaskArgValue::Number(0)), + ( + CMD_OUTPUT.to_string(), + TaskArgValue::Str(format!( + "No matching {server_type} ports found for host {} in current topology; skipped.", + self.task_id.host + )), + ), + ]))); + } + let runtime_ports = target_ports .iter() .filter_map(|p| p.parse::().ok()) @@ -920,7 +938,7 @@ impl TaskExecutor for EloqTxCtlTask { (CMD_STATUS.to_string(), TaskArgValue::Number(0)), ( CMD_OUTPUT.to_string(), - TaskArgValue::Str(format!("eloqkv service is down: {err}")), + TaskArgValue::Str(format!("eloqkv service status unknown: {err}")), ), ]), Err(err) => return Err(err), diff --git a/src/cluster_mgr/src/cli/task/failover_op_task.rs b/src/cluster_mgr/src/cli/task/failover_op_task.rs index a110ce00..12f036bd 100644 --- a/src/cluster_mgr/src/cli/task/failover_op_task.rs +++ b/src/cluster_mgr/src/cli/task/failover_op_task.rs @@ -22,6 +22,7 @@ pub struct FailoverOpTask { receiver: watch::Receiver, password: Option, service_endpoints: Option>, + require_explicit_target: bool, } impl FailoverOpTask { @@ -43,6 +44,7 @@ impl FailoverOpTask { receiver, password, service_endpoints: None, + require_explicit_target: false, } } @@ -54,27 +56,82 @@ impl FailoverOpTask { self } + pub fn require_explicit_target(mut self) -> Self { + self.require_explicit_target = true; + self + } + // Helper function to find the best replica for failover fn find_best_replica(&self, cluster_nodes: &ClusterNodes) -> Option<(String, u16)> { - // If we have explicit new_leader_host/port set and it's in the replicas list, use it + // If the caller selected a new_leader_host/port, require it to be a valid replica. if !self.new_leader_host.is_empty() && self.new_leader_port > 0 { - let specified_replica = cluster_nodes + let Some(specified_replica) = cluster_nodes .replicas .iter() - .find(|r| r.ip == self.new_leader_host && r.port == self.new_leader_port); + .find(|r| r.ip == self.new_leader_host && r.port == self.new_leader_port) + else { + if self.require_explicit_target { + error!( + "Selected new leader {}:{} not found as replica", + self.new_leader_host, self.new_leader_port + ); + return None; + } + info!( + "Selected new leader {}:{} not found as replica, will choose another one", + self.new_leader_host, self.new_leader_port + ); + return cluster_nodes + .replicas + .iter() + .find(|replica| replica.connected) + .map(|replica| (replica.ip.clone(), replica.port)); + }; + + let old_leader = cluster_nodes + .masters + .iter() + .find(|node| node.ip == self.old_leader_host && node.port == self.old_leader_port); + if let (Some(old_leader), Some(master_id)) = + (old_leader, specified_replica.master_id.as_ref()) + { + if old_leader.node_id.as_ref() != Some(master_id) { + error!( + "Selected new leader {}:{} is not a replica of old leader {}:{}", + self.new_leader_host, + self.new_leader_port, + self.old_leader_host, + self.old_leader_port + ); + return None; + } + } - if specified_replica.is_some() { + if specified_replica.connected { return Some((self.new_leader_host.clone(), self.new_leader_port)); } - info!( - "Specified new leader {}:{} not found as replica, will choose another one", + error!( + "Selected new leader {}:{} is not connected as replica", self.new_leader_host, self.new_leader_port ); + return None; + } + + if self.require_explicit_target { + error!( + "Internal failover target is required for {}:{}, but no target was selected", + self.old_leader_host, self.old_leader_port + ); + return None; } // Select first available replica - if let Some(replica) = cluster_nodes.replicas.first() { + if let Some(replica) = cluster_nodes + .replicas + .iter() + .find(|replica| replica.connected) + { return Some((replica.ip.clone(), replica.port)); } diff --git a/src/cluster_mgr/src/cli/task/group/db_cluster_ctrl_group.rs b/src/cluster_mgr/src/cli/task/group/db_cluster_ctrl_group.rs index ae898616..c1cdfac5 100644 --- a/src/cluster_mgr/src/cli/task/group/db_cluster_ctrl_group.rs +++ b/src/cluster_mgr/src/cli/task/group/db_cluster_ctrl_group.rs @@ -181,6 +181,7 @@ impl TaskGroup for CtrlDBTaskGroup { let (redis_tx, redis_rx) = watch::channel(ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), }); let redis_task_id = TaskId { cmd: "topology".to_string(), @@ -517,6 +518,7 @@ impl CtrlDBTaskGroup { let (topology_tx, _) = watch::channel(ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), }); let mut topology_nodes = config.get_host_port_list(DeploymentPackage::EloqTx); topology_nodes.extend(config.get_host_port_list(DeploymentPackage::EloqStandby)); diff --git a/src/cluster_mgr/src/cli/task/group/failover_group.rs b/src/cluster_mgr/src/cli/task/group/failover_group.rs index a786b96a..e7b76001 100644 --- a/src/cluster_mgr/src/cli/task/group/failover_group.rs +++ b/src/cluster_mgr/src/cli/task/group/failover_group.rs @@ -52,6 +52,7 @@ fn failover_task_group( let (pre_failover_tx, pre_failover_rx) = watch::channel::(ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), }); // Create another channel for the post-failover task @@ -59,6 +60,7 @@ fn failover_task_group( let (post_failover_tx, post_failover_rx) = watch::channel::(ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), }); // Keep the receiver alive to prevent the channel from closing diff --git a/src/cluster_mgr/src/cli/task/group/launch_group.rs b/src/cluster_mgr/src/cli/task/group/launch_group.rs index ae48da6f..c42977ce 100644 --- a/src/cluster_mgr/src/cli/task/group/launch_group.rs +++ b/src/cluster_mgr/src/cli/task/group/launch_group.rs @@ -201,6 +201,7 @@ impl TaskGroup for LaunchTaskGroup { let empty_cluster_nodes = ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), }; let (redis_op_tx, redis_op_rx) = watch::channel(empty_cluster_nodes.clone()); diff --git a/src/cluster_mgr/src/cli/task/group/scale_group.rs b/src/cluster_mgr/src/cli/task/group/scale_group.rs index e789be1b..bb2c2fa3 100644 --- a/src/cluster_mgr/src/cli/task/group/scale_group.rs +++ b/src/cluster_mgr/src/cli/task/group/scale_group.rs @@ -162,6 +162,7 @@ impl super::TaskGroup for ScaleTaskGroup { let empty_cluster_nodes = ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), }; let empty_cluster_nodes_with_config = ClusterNodesWithConfig { @@ -713,6 +714,7 @@ impl super::TaskGroup for ScaleTaskGroup { let empty_cluster_nodes = ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), }; let (final_topology_tx, final_topology_rx) = watch::channel(empty_cluster_nodes.clone()); diff --git a/src/cluster_mgr/src/cli/task/group/scale_log_group.rs b/src/cluster_mgr/src/cli/task/group/scale_log_group.rs index 0375348a..4c0ad1fe 100644 --- a/src/cluster_mgr/src/cli/task/group/scale_log_group.rs +++ b/src/cluster_mgr/src/cli/task/group/scale_log_group.rs @@ -841,6 +841,7 @@ impl TaskGroup for ScaleLogTaskGroup { let empty_cluster_nodes = ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), }; let (final_topology_tx, final_topology_rx) = watch::channel(empty_cluster_nodes.clone()); diff --git a/src/cluster_mgr/src/cli/task/group/update_config_group.rs b/src/cluster_mgr/src/cli/task/group/update_config_group.rs index e84836c2..c2eaf50a 100644 --- a/src/cluster_mgr/src/cli/task/group/update_config_group.rs +++ b/src/cluster_mgr/src/cli/task/group/update_config_group.rs @@ -96,6 +96,7 @@ fn build_config_update( let (redis_op_tx, redis_op_rx) = watch::channel(ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), }); let redis_task_id = TaskId { diff --git a/src/cluster_mgr/src/cli/task/redis_op_task.rs b/src/cluster_mgr/src/cli/task/redis_op_task.rs index 23ef6141..dd0df052 100644 --- a/src/cluster_mgr/src/cli/task/redis_op_task.rs +++ b/src/cluster_mgr/src/cli/task/redis_op_task.rs @@ -137,6 +137,7 @@ mod tests { let (tx, _rx) = watch::channel(ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), }); let task = RedisOpTask::new( TaskId { @@ -175,6 +176,10 @@ pub struct NodeInfo { pub ip: String, pub port: u16, pub connected: bool, + #[serde(default)] + pub node_id: Option, + #[serde(default)] + pub master_id: Option, } // Implement Eq and Hash for NodeInfo @@ -204,6 +209,8 @@ impl fmt::Display for NodeInfo { pub struct ClusterNodes { pub masters: Vec, pub replicas: Vec, + #[serde(default)] + pub voters: Vec, } const MAX_RETRIES: usize = 500; @@ -302,6 +309,7 @@ pub fn parse_cluster_nodes(value: Value) -> RedisResult> { let mut cluster_nodes_list = Vec::new(); let mut masters = Vec::new(); let mut replicas = Vec::new(); + let mut voters = Vec::new(); // Parse each line of the CLUSTER NODES output for line in nodes_str.lines() { @@ -332,9 +340,17 @@ pub fn parse_cluster_nodes(value: Value) -> RedisResult> { Err(_) => continue, // Skip if port is not a valid u16 }; - // Check if node is master or replica + // Check node role from comma-delimited flags. let flags = parts[2]; - let is_master = !flags.contains("slave"); + let is_master = flags + .split(',') + .any(|flag| flag.eq_ignore_ascii_case("master")); + let is_replica = flags + .split(',') + .any(|flag| flag.eq_ignore_ascii_case("slave")); + let is_voter = flags + .split(',') + .any(|flag| flag.eq_ignore_ascii_case("voter")); // Add node to appropriate list let connected = parts[7].eq_ignore_ascii_case("connected"); @@ -342,18 +358,26 @@ pub fn parse_cluster_nodes(value: Value) -> RedisResult> { ip, port, connected, + node_id: Some(parts[0].to_string()), + master_id: (parts[3] != "-").then(|| parts[3].to_string()), }; if is_master { masters.push(node_info); - } else { + } else if is_replica { replicas.push(node_info); + } else if is_voter { + voters.push(node_info); } } // Group all masters and replicas into one ClusterNodes - if !masters.is_empty() || !replicas.is_empty() { - cluster_nodes_list.push(ClusterNodes { masters, replicas }); + if !masters.is_empty() || !replicas.is_empty() || !voters.is_empty() { + cluster_nodes_list.push(ClusterNodes { + masters, + replicas, + voters, + }); } Ok(cluster_nodes_list) @@ -402,8 +426,11 @@ pub fn parse_cluster_nodes_single(value: Value, default_host: &str) -> RedisResu ip, port, connected: true, + node_id: None, + master_id: None, }], replicas: vec![], + voters: vec![], }) } @@ -588,6 +615,7 @@ impl TaskExecutor for RedisOpTask { }; let mut unique_masters = HashSet::new(); let mut unique_replicas = HashSet::new(); + let mut unique_voters = HashSet::new(); for slot in &cluster_nodes { for master in &slot.masters { @@ -596,6 +624,9 @@ impl TaskExecutor for RedisOpTask { for replica in &slot.replicas { unique_replicas.insert(replica.clone()); } + for voter in &slot.voters { + unique_voters.insert(voter.clone()); + } } if unique_masters.is_empty() { @@ -708,6 +739,7 @@ impl TaskExecutor for RedisOpTask { // Convert HashSets to Vectors let unique_masters: Vec = unique_masters.into_iter().collect(); let unique_replicas: Vec = unique_replicas.into_iter().collect(); + let unique_voters: Vec = unique_voters.into_iter().collect(); // For debugging: print the unique masters and replicas for master in &unique_masters { @@ -716,10 +748,14 @@ impl TaskExecutor for RedisOpTask { for replica in &unique_replicas { info!("Replicas: {}:{}", replica.ip, replica.port); } + for voter in &unique_voters { + info!("Voters: {}:{}", voter.ip, voter.port); + } let cluster_nodes = ClusterNodes { masters: unique_masters, replicas: unique_replicas, + voters: unique_voters, }; let response_str = serde_json::to_string(&cluster_nodes)?; diff --git a/src/cluster_mgr/src/cli/task/rolling_upgrade/mod.rs b/src/cluster_mgr/src/cli/task/rolling_upgrade/mod.rs index 35d88aee..85db8306 100644 --- a/src/cluster_mgr/src/cli/task/rolling_upgrade/mod.rs +++ b/src/cluster_mgr/src/cli/task/rolling_upgrade/mod.rs @@ -134,6 +134,12 @@ impl RollingUpgrade { fn friendly_step_name(name: &str) -> &str { match name { "DownloadAndExtract" => "Prepare upgrade package", + "UploadToAllNodes" => "Upload binaries to all nodes", + "SelectStandbyForFailover" => "Select standby for failover", + "RestartSelectedStandby" => "Restart selected standby node", + "FailoverToStandby" => "Fail over to upgraded standby", + "RestartNonLeaderNodes" => "Restart non-leader nodes one by one", + "RestartTemporaryLeaders" => "Restart temporary leader nodes", "UploadToStandby" => "Upload binaries to standby nodes", "UploadToMaster" => "Upload binaries to old master nodes", "StopStandbyOnly" => "Stop standby nodes", diff --git a/src/cluster_mgr/src/cli/task/rolling_upgrade/steps.rs b/src/cluster_mgr/src/cli/task/rolling_upgrade/steps.rs index 6c6d8ccb..6085c8bb 100644 --- a/src/cluster_mgr/src/cli/task/rolling_upgrade/steps.rs +++ b/src/cluster_mgr/src/cli/task/rolling_upgrade/steps.rs @@ -8,17 +8,20 @@ use crate::cli::task::exec_custom_cmd::ExecCustomCommand; use crate::cli::task::failover_op_task::FailoverOpTask; use crate::cli::task::group::Config; use crate::cli::task::local_extract_task::LocalExtractTask; -use crate::cli::task::redis_op_task::{ClusterNodes, RedisOpTask}; -use crate::cli::task::task_base::{TaskExecutionContext, TaskHost, TaskId, TaskInstance}; -use crate::cli::task::wait_replica_ready_task::WaitReplicaReadyTask; -use crate::cli::SubCommand; +use crate::cli::task::redis_op_task::{ClusterNodes, NodeInfo, RedisOpTask}; +use crate::cli::task::task_base::{ + TaskArgValue, TaskExecutionContext, TaskExecutor, TaskHost, TaskId, TaskInstance, +}; +use crate::cli::task::wait_replica_ready_task::{WaitNodeReadyTask, WaitReplicaReadyTask}; +use crate::cli::{SubCommand, CMD_OUTPUT, CMD_STATUS}; use crate::config::config_base::{DeployConfig, ELOQ_FILE_KEY, ELOQ_LOG_FILE_KEY}; use crate::config::storage_service_config::DataStoreServiceBackend; use crate::config::DeploymentPackage; use anyhow::bail; use async_trait::async_trait; use indexmap::IndexMap; -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; +use std::sync::{Arc, Mutex}; use tokio::sync::watch; use tracing::info; @@ -36,6 +39,387 @@ fn single_barrier_ctx( } } +fn node_addr(host: &str, port: u16) -> String { + format!("{host}:{port}") +} + +fn hosts_from_host_ports(host_ports: &[String]) -> Vec { + host_ports + .iter() + .filter_map(|host_port| host_port.split_once(':').map(|(host, _)| host.to_string())) + .collect::>() + .into_iter() + .collect() +} + +fn connected_managed_nodes( + nodes: &[crate::cli::task::redis_op_task::NodeInfo], + managed_nodes: &HashSet, +) -> Vec { + nodes + .iter() + .filter(|node| node.connected) + .map(|node| node_addr(&node.ip, node.port)) + .filter(|node| managed_nodes.contains(node)) + .collect() +} + +fn connected_managed_node_infos( + nodes: &[NodeInfo], + managed_nodes: &HashSet, +) -> Vec { + nodes + .iter() + .filter(|node| node.connected) + .filter(|node| managed_nodes.contains(&node_addr(&node.ip, node.port))) + .cloned() + .collect() +} + +fn parse_host_port(value: &str, field: &str) -> anyhow::Result<(String, u16)> { + let Some((host, port_str)) = value.split_once(':') else { + bail!("invalid {field} host:port: '{value}'"); + }; + if host.trim().is_empty() { + bail!("invalid {field} host:port with empty host: '{value}'"); + } + let port = port_str + .parse::() + .map_err(|_| anyhow::anyhow!("invalid {field} port in host:port: '{value}'"))?; + Ok((host.to_string(), port)) +} + +fn ordered_without(nodes: Vec, excluded: &HashSet) -> Vec { + let mut seen = HashSet::new(); + nodes + .into_iter() + .filter(|node| !excluded.contains(node)) + .filter(|node| seen.insert(node.clone())) + .collect() +} + +fn select_connected_failover_targets( + topology: &ClusterNodes, + managed_tx_standby: &HashSet, +) -> anyhow::Result> { + let current_masters = connected_managed_node_infos(&topology.masters, managed_tx_standby); + let current_replicas = connected_managed_node_infos(&topology.replicas, managed_tx_standby); + if current_masters.is_empty() { + bail!( + "rolling update could not find connected current master; topology reported masters={:?}, replicas={:?}", + topology.masters, + topology.replicas + ); + } + if current_replicas.is_empty() { + bail!( + "rolling update requires at least one connected standby/replica before failover; topology reported masters={:?}, replicas={:?}", + topology.masters, + topology.replicas + ); + } + + let mut failover_pairs: Vec<(String, String)> = Vec::new(); + let mut selected_targets = HashSet::new(); + for master in ¤t_masters { + let master_addr = node_addr(&master.ip, master.port); + let target = if let Some(master_id) = master.node_id.as_ref() { + current_replicas + .iter() + .find(|replica| replica.master_id.as_ref() == Some(master_id)) + .ok_or_else(|| { + anyhow::anyhow!( + "no connected replica available for master {master_addr} ({master_id})" + ) + })? + } else if current_masters.len() == 1 { + current_replicas.first().ok_or_else(|| { + anyhow::anyhow!("no connected replica available for master {master_addr}") + })? + } else { + bail!( + "CLUSTER NODES output did not include node id for master {master_addr}; cannot safely map replicas to masters in a multi-master cluster" + ); + }; + let target_addr = node_addr(&target.ip, target.port); + if selected_targets.contains(&target_addr) { + bail!( + "replica {target_addr} was selected for more than one current master; cannot safely fail over" + ); + } + selected_targets.insert(target_addr.clone()); + failover_pairs.push((master_addr, target_addr)); + } + + Ok(failover_pairs) +} + +async fn fetch_cluster_nodes( + ctx: &UpgradeContext, + task_name: &str, +) -> anyhow::Result { + let task_id = TaskId { + cmd: "topology".to_string(), + task: task_name.to_string(), + host: "_local".to_string(), + }; + let (topology_tx, _) = watch::channel(ClusterNodes { + masters: Vec::new(), + replicas: Vec::new(), + voters: Vec::new(), + }); + let result = RedisOpTask::new( + task_id, + ctx.redis_cluster_startup_nodes(), + "cluster topology".to_string(), + topology_tx, + ctx.redis_password.clone(), + true, + ) + .with_service_endpoints(ctx.deploy.connection.service_endpoints.clone()) + .execute(TaskHost::Local, HashMap::default()) + .await?; + + let values = result.ok_or_else(|| anyhow::anyhow!("missing topology task result"))?; + let status = values + .get(CMD_STATUS) + .cloned() + .unwrap_or(TaskArgValue::Number(1)); + let output = values + .get(CMD_OUTPUT) + .cloned() + .unwrap_or_else(|| TaskArgValue::Str("missing cluster topology output".to_string())); + + match (status, output) { + (TaskArgValue::Number(0), TaskArgValue::Str(json)) => { + Ok(serde_json::from_str::(&json)?) + } + (_, TaskArgValue::Str(err)) => Err(anyhow::anyhow!(err)), + _ => Err(anyhow::anyhow!("unexpected topology task output")), + } +} + +fn build_stop_node_tasks( + ctx: &UpgradeContext, + task_group: &str, + nodes: Vec, +) -> anyhow::Result { + if nodes.is_empty() { + return Ok(TaskExecutionContext::dummy()); + } + let stop = EloqTxCtlTask::from_config_with_channel( + SubCommand::Stop { + cluster: ctx.cluster.clone(), + tx: Some(true), + log: true, + store: false, + monitor: false, + force: true, + all: false, + password: ctx.redis_password.clone(), + nodes, + }, + &ctx.deploy, + ServerType::Node, + None, + )?; + Ok(single_barrier_ctx(task_group, stop)) +} + +fn build_start_node_tasks( + ctx: &UpgradeContext, + task_group: &str, + nodes: Vec, +) -> TaskExecutionContext { + if nodes.is_empty() { + return TaskExecutionContext::dummy(); + } + let start = EloqTxCtlTask::from_config( + SubCommand::Start { + cluster: ctx.cluster.clone(), + nodes, + }, + &ctx.deploy, + ServerType::Node, + ); + single_barrier_ctx(task_group, start) +} + +fn append_step_context( + barrier: &mut Vec, + executable: &mut IndexMap, + ctx: TaskExecutionContext, +) { + if ctx.executable.is_empty() { + return; + } + if let Some(step_barrier) = ctx.barrier { + barrier.extend(step_barrier); + } else { + barrier.push(ctx.executable.len()); + } + executable.extend(ctx.executable); +} + +fn build_wait_replica_ready_tasks( + ctx: &UpgradeContext, + task_group: &str, + task_prefix: &str, + source_master: &str, + target_replicas: &[String], +) -> anyhow::Result { + if target_replicas.is_empty() { + return Ok(TaskExecutionContext::dummy()); + } + + let Some((source_host, source_port_str)) = source_master.split_once(':') else { + bail!("invalid host:port in current master: '{source_master}'"); + }; + let Ok(source_port) = source_port_str.parse::() else { + bail!("invalid port in current master: '{source_master}'"); + }; + + let mut executable = IndexMap::new(); + for target_addr in target_replicas { + let Some((target_host, target_port_str)) = target_addr.split_once(':') else { + bail!("invalid host:port in target replica list: '{target_addr}'"); + }; + let Ok(target_port) = target_port_str.parse::() else { + bail!("invalid port in target replica list: '{target_addr}'"); + }; + let task_id = TaskId { + cmd: "topology".to_string(), + task: format!("{task_prefix}-{target_port}"), + host: target_host.to_string(), + }; + executable.insert( + task_id.clone(), + TaskInstance { + task_input: HashMap::default(), + task: Box::new( + WaitReplicaReadyTask::new( + task_id, + ctx.redis_cluster_startup_nodes(), + source_host.to_string(), + source_port, + target_host.to_string(), + target_port, + ctx.redis_password.clone(), + ) + .with_service_endpoints(ctx.deploy.connection.service_endpoints.clone()), + ), + task_host: TaskHost::Local, + }, + ); + } + + Ok(single_barrier_ctx(task_group, executable)) +} + +fn build_wait_node_ready_tasks( + ctx: &UpgradeContext, + task_group: &str, + task_prefix: &str, + target_nodes: &[String], + require_master: bool, +) -> anyhow::Result { + if target_nodes.is_empty() { + return Ok(TaskExecutionContext::dummy()); + } + + let mut executable = IndexMap::new(); + for target_addr in target_nodes { + let (target_host, target_port) = parse_host_port(target_addr, "target node")?; + let task_id = TaskId { + cmd: "topology".to_string(), + task: format!("{task_prefix}-{target_port}"), + host: target_host.clone(), + }; + let mut task = WaitNodeReadyTask::new( + task_id.clone(), + ctx.redis_cluster_startup_nodes(), + target_host, + target_port, + ctx.redis_password.clone(), + ) + .with_service_endpoints(ctx.deploy.connection.service_endpoints.clone()); + if require_master { + task = task.require_master(); + } + + executable.insert( + task_id.clone(), + TaskInstance { + task_input: HashMap::default(), + task: Box::new(task), + task_host: TaskHost::Local, + }, + ); + } + + Ok(single_barrier_ctx(task_group, executable)) +} + +fn build_restart_nodes_sequence_ctx( + ctx: &UpgradeContext, + task_group: &str, + task_prefix: &str, + nodes: Vec, +) -> anyhow::Result { + if nodes.is_empty() { + return Ok(TaskExecutionContext::dummy()); + } + + let mut barrier = Vec::new(); + let mut executable = IndexMap::new(); + for node in nodes { + let (_host, port) = parse_host_port(&node, "restart node")?; + append_step_context( + &mut barrier, + &mut executable, + build_stop_node_tasks( + ctx, + &format!("{task_prefix}-stop-{port}"), + vec![node.clone()], + )?, + ); + append_step_context( + &mut barrier, + &mut executable, + build_start_node_tasks( + ctx, + &format!("{task_prefix}-start-{port}"), + vec![node.clone()], + ), + ); + append_step_context( + &mut barrier, + &mut executable, + build_wait_node_ready_tasks( + ctx, + &format!("{task_prefix}-wait-{port}"), + &format!("{task_prefix}-wait"), + &[node], + false, + )?, + ); + } + + Ok(TaskExecutionContext { + task_group: task_group.to_string(), + barrier: Some(barrier), + executable, + }) +} + +#[derive(Clone, Default)] +struct RollingUpdateState { + first_batch_nodes: Vec, + second_batch_nodes: Vec, + temporary_leader_nodes: Vec, + restart_nodes: Vec, +} + // ── Context ─────────────────────────────────────────────────────────────────── /// All configuration extracted upfront from CLI args + deploy config. @@ -47,6 +431,8 @@ pub struct UpgradeContext { pub cluster: String, pub redis_password: Option, pub force: bool, + pub skip_log_restart: bool, + update_state: Arc>, } impl UpgradeContext { @@ -55,14 +441,21 @@ impl UpgradeContext { pub(crate) fn new(cmd_arg: &SubCommand, config: Config) -> Self { let Config::Cluster(ref deploy) = config; let deploy = deploy.clone(); - let (redis_password, force) = match cmd_arg { + let (redis_password, force, skip_log_restart) = match cmd_arg { SubCommand::Update { - password, force, .. - } => (deploy.redis_password(password.clone()), *force), + password, + force, + skip_log_restart, + .. + } => ( + deploy.redis_password(password.clone()), + *force, + *skip_log_restart, + ), SubCommand::UpdateConf { password, .. } => { - (deploy.redis_password(password.clone()), false) + (deploy.redis_password(password.clone()), false, false) } - _ => (None, false), + _ => (None, false, false), }; Self { cluster: deploy.deployment.cluster_name.clone(), @@ -70,6 +463,8 @@ impl UpgradeContext { deploy, redis_password, force, + skip_log_restart, + update_state: Arc::new(Mutex::new(RollingUpdateState::default())), } } @@ -107,6 +502,80 @@ impl UpgradeContext { host_ports.extend(self.standby_host_ports()); host_ports } + + fn managed_tx_and_standby_nodes(&self) -> Vec { + let mut host_ports = self.tx_host_ports(); + host_ports.extend(self.standby_host_ports()); + host_ports + } + + fn rolling_update_kv_nodes(&self) -> Vec { + self.managed_tx_and_standby_nodes() + } + + fn managed_tx_and_standby_set(&self) -> HashSet { + self.managed_tx_and_standby_nodes().into_iter().collect() + } + + fn set_first_batch_nodes(&self, nodes: Vec) { + self.update_state + .lock() + .expect("rolling update state poisoned") + .first_batch_nodes = nodes; + } + + fn first_batch_nodes(&self) -> Vec { + self.update_state + .lock() + .expect("rolling update state poisoned") + .first_batch_nodes + .clone() + } + + fn set_second_batch_nodes(&self, nodes: Vec) { + self.update_state + .lock() + .expect("rolling update state poisoned") + .second_batch_nodes = nodes; + } + + fn second_batch_nodes(&self) -> Vec { + self.update_state + .lock() + .expect("rolling update state poisoned") + .second_batch_nodes + .clone() + } + + fn set_temporary_leader_nodes(&self, nodes: Vec) { + self.update_state + .lock() + .expect("rolling update state poisoned") + .temporary_leader_nodes = nodes; + } + + fn temporary_leader_nodes(&self) -> Vec { + self.update_state + .lock() + .expect("rolling update state poisoned") + .temporary_leader_nodes + .clone() + } + + fn set_restart_nodes(&self, nodes: Vec) { + self.update_state + .lock() + .expect("rolling update state poisoned") + .restart_nodes = nodes; + } + + fn restart_nodes(&self) -> Vec { + self.update_state + .lock() + .expect("rolling update state poisoned") + .restart_nodes + .clone() + } } // ── Helper: build a round of topo→failover→stop ───────────────────────────── @@ -131,22 +600,23 @@ impl Step for StopStandbyOnly { if !self.ctx.has_standby() { return Ok(TaskExecutionContext::dummy()); } - let stop = EloqTxCtlTask::from_config( - SubCommand::Stop { - cluster: self.ctx.cluster.clone(), - tx: None, - log: false, - store: false, - monitor: false, - force: true, - all: false, - password: self.ctx.redis_password.clone(), - nodes: Vec::new(), - }, - &self.ctx.deploy, - ServerType::Standby, + let topology = fetch_cluster_nodes(&self.ctx, "rolling-update-initial-topology").await?; + let managed_nodes = self.ctx.managed_tx_and_standby_set(); + let current_replicas = connected_managed_nodes(&topology.replicas, &managed_nodes); + let first_batch_nodes = current_replicas; + if first_batch_nodes.is_empty() { + bail!( + "rolling update requires a connected replica/standby node, but topology reported masters={:?}, replicas={:?}", + topology.masters, + topology.replicas + ); + } + info!( + "Rolling update first batch nodes selected from current replicas: {:?}", + first_batch_nodes ); - Ok(single_barrier_ctx("stop-standby", stop)) + self.ctx.set_first_batch_nodes(first_batch_nodes.clone()); + build_stop_node_tasks(&self.ctx, "stop-standby", first_batch_nodes) } } @@ -186,21 +656,300 @@ impl Step for FailoverAndStopOldMaster { return Ok(single_barrier_ctx("stop-old-master", stop_tx)); } - let tx_host_ports = self.ctx.tx_host_ports(); - let standby_host_ports = self.ctx.standby_host_ports(); - let mut all_nodes = tx_host_ports.clone(); - all_nodes.extend(standby_host_ports); + let topology = + fetch_cluster_nodes(&self.ctx, "rolling-update-pre-failover-topology").await?; + let managed_nodes = self.ctx.managed_tx_and_standby_set(); + let current_masters = connected_managed_nodes(&topology.masters, &managed_nodes); + let current_replicas = connected_managed_nodes(&topology.replicas, &managed_nodes); + if current_masters.is_empty() { + bail!( + "rolling update could not find a connected current master before failover; topology reported masters={:?}, replicas={:?}", + topology.masters, + topology.replicas + ); + } + if current_replicas.is_empty() { + bail!( + "rolling update could not find a connected current replica to fail over to; topology reported masters={:?}, replicas={:?}", + topology.masters, + topology.replicas + ); + } + info!( + "Rolling update failover sources selected from current masters: {:?}; current replicas: {:?}", + current_masters, + current_replicas + ); + self.ctx.set_second_batch_nodes(current_masters.clone()); + let all_nodes = self.ctx.managed_tx_and_standby_nodes(); build_round( "failover-stop-master", - &tx_host_ports, - &tx_host_ports, + ¤t_masters, + ¤t_masters, &all_nodes, &self.ctx, ) } } +pub struct FailoverToStandby { + ctx: UpgradeContext, +} + +pub struct SelectStandbyForFailover { + ctx: UpgradeContext, +} + +impl SelectStandbyForFailover { + pub fn new(ctx: UpgradeContext) -> Self { + Self { ctx } + } +} + +#[async_trait] +impl Step for SelectStandbyForFailover { + fn name(&self) -> &str { + "SelectStandbyForFailover" + } + + async fn build(&self) -> anyhow::Result { + if !self.ctx.has_standby() { + return Ok(TaskExecutionContext::dummy()); + } + + let topology = + fetch_cluster_nodes(&self.ctx, "rolling-update-select-failover-target").await?; + let managed_tx_standby = self.ctx.managed_tx_and_standby_set(); + let failover_pairs = select_connected_failover_targets(&topology, &managed_tx_standby)?; + let temporary_leaders: Vec = failover_pairs + .iter() + .map(|(_old_leader, new_leader)| new_leader.clone()) + .collect(); + let temporary_leader_set: HashSet = temporary_leaders.iter().cloned().collect(); + let restart_nodes = + ordered_without(self.ctx.rolling_update_kv_nodes(), &temporary_leader_set); + + self.ctx + .set_temporary_leader_nodes(temporary_leaders.clone()); + self.ctx.set_restart_nodes(restart_nodes.clone()); + self.ctx.set_second_batch_nodes(restart_nodes); + info!( + "Rolling update selected standby nodes to restart before failover: {:?}; initial failover pairs: {:?}", + temporary_leaders, failover_pairs + ); + + Ok(TaskExecutionContext::dummy()) + } +} + +impl FailoverToStandby { + pub fn new(ctx: UpgradeContext) -> Self { + Self { ctx } + } +} + +#[async_trait] +impl Step for FailoverToStandby { + fn name(&self) -> &str { + "FailoverToStandby" + } + + async fn build(&self) -> anyhow::Result { + if !self.ctx.has_standby() { + return Ok(TaskExecutionContext::dummy()); + } + + let topology = + fetch_cluster_nodes(&self.ctx, "rolling-update-pre-failover-topology").await?; + let managed_tx_standby = self.ctx.managed_tx_and_standby_set(); + let selected_targets = self.ctx.temporary_leader_nodes(); + let failover_pairs = if selected_targets.is_empty() { + let pairs = select_connected_failover_targets(&topology, &managed_tx_standby)?; + let targets: Vec = pairs + .iter() + .map(|(_old_leader, new_leader)| new_leader.clone()) + .collect(); + let target_set: HashSet = targets.iter().cloned().collect(); + let restart_nodes = ordered_without(self.ctx.rolling_update_kv_nodes(), &target_set); + self.ctx.set_temporary_leader_nodes(targets.clone()); + self.ctx.set_restart_nodes(restart_nodes.clone()); + self.ctx.set_second_batch_nodes(restart_nodes); + pairs + } else { + let connected_masters = + connected_managed_node_infos(&topology.masters, &managed_tx_standby); + if connected_masters.is_empty() { + bail!( + "rolling update could not find connected current master before failover; topology reported masters={:?}, replicas={:?}", + topology.masters, + topology.replicas + ); + } + + let mut pairs = Vec::new(); + for target_addr in &selected_targets { + let (target_host, target_port) = + parse_host_port(target_addr, "selected failover target")?; + if connected_masters + .iter() + .any(|node| node.ip == target_host && node.port == target_port) + { + info!( + "Selected failover target {target_addr} is already connected as master; no failover command needed" + ); + continue; + } + + let target = topology + .replicas + .iter() + .find(|node| { + node.ip == target_host && node.port == target_port && node.connected + }) + .ok_or_else(|| { + anyhow::anyhow!( + "selected standby {target_addr} is not a connected replica before failover; topology reported masters={:?}, replicas={:?}", + topology.masters, + topology.replicas + ) + })?; + + let master = if let Some(master_id) = target.master_id.as_ref() { + connected_masters + .iter() + .find(|master| master.node_id.as_ref() == Some(master_id)) + .ok_or_else(|| { + anyhow::anyhow!( + "selected standby {target_addr} follows master id {master_id}, but that master is not connected; topology reported masters={:?}, replicas={:?}", + topology.masters, + topology.replicas + ) + })? + } else if connected_masters.len() == 1 { + connected_masters.first().expect("checked non-empty") + } else { + bail!( + "CLUSTER NODES output did not include master id for selected standby {target_addr}; cannot safely map it to a master in a multi-master cluster" + ); + }; + pairs.push((node_addr(&master.ip, master.port), target_addr.clone())); + } + pairs + }; + + if failover_pairs.is_empty() { + return build_wait_node_ready_tasks( + &self.ctx, + "wait-failover-target-masters", + "wait-failover-target-master", + &self.ctx.temporary_leader_nodes(), + true, + ); + } + info!( + "Rolling update failover pairs after selected standby restart: {:?}; temporary leaders held until final restart: {:?}", + failover_pairs, + self.ctx.temporary_leader_nodes() + ); + + build_failover_pairs_ctx("failover-to-standby", &failover_pairs, &self.ctx) + } +} + +fn build_failover_pairs_ctx( + round_label: &str, + failover_pairs: &[(String, String)], + ctx: &UpgradeContext, +) -> anyhow::Result { + let mut barrier = vec![]; + let mut executable = IndexMap::new(); + + let topo_task_id = TaskId { + cmd: "topology".to_string(), + task: format!("check-topology-{round_label}"), + host: "_local".to_string(), + }; + let (topo_tx, failover_rx) = watch::channel::(ClusterNodes { + masters: Vec::new(), + replicas: Vec::new(), + voters: Vec::new(), + }); + + executable.insert( + topo_task_id.clone(), + TaskInstance { + task_input: HashMap::default(), + task: Box::new( + RedisOpTask::new( + topo_task_id, + ctx.redis_cluster_startup_nodes(), + "cluster topology".to_string(), + topo_tx, + ctx.redis_password.clone(), + true, + ) + .with_service_endpoints(ctx.deploy.connection.service_endpoints.clone()), + ), + task_host: TaskHost::Local, + }, + ); + barrier.push(1); + + for (old_leader, new_leader) in failover_pairs { + let (old_host, old_port) = parse_host_port(old_leader, "old leader")?; + let (new_host, new_port) = parse_host_port(new_leader, "new leader")?; + let fid = TaskId { + cmd: "failover".to_string(), + task: format!("failover-{round_label}-{old_port}-to-{new_port}"), + host: old_host.clone(), + }; + executable.insert( + fid.clone(), + TaskInstance { + task_input: HashMap::default(), + task: Box::new( + FailoverOpTask::new( + fid, + old_host, + old_port, + new_host, + new_port, + failover_rx.clone(), + ctx.redis_password.clone(), + ) + .require_explicit_target() + .with_service_endpoints(ctx.deploy.connection.service_endpoints.clone()), + ), + task_host: TaskHost::Local, + }, + ); + } + barrier.push(failover_pairs.len()); + + let failover_targets: Vec = failover_pairs + .iter() + .map(|(_old_leader, new_leader)| new_leader.clone()) + .collect(); + append_step_context( + &mut barrier, + &mut executable, + build_wait_node_ready_tasks( + ctx, + "wait-failover-target-masters", + "wait-failover-target-master", + &failover_targets, + true, + )?, + ); + + Ok(TaskExecutionContext { + task_group: format!("rolling-restart-{round_label}"), + barrier: Some(barrier), + executable, + }) +} + fn build_round( round_label: &str, nodes_to_failover: &[String], @@ -219,6 +968,7 @@ fn build_round( let (topo_tx, failover_rx) = watch::channel::(ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), }); let stop_rx = failover_rx.clone(); @@ -377,11 +1127,9 @@ impl Step for UploadToStandby { &self.ctx.deploy.deployment, ); - let standby_hosts: std::collections::HashSet = self - .ctx - .standby_host_ports() - .iter() - .filter_map(|hp| hp.split(':').next().map(|h| h.to_string())) + let first_batch_nodes = self.ctx.first_batch_nodes(); + let standby_hosts: HashSet = hosts_from_host_ports(&first_batch_nodes) + .into_iter() .collect(); let log_hosts: std::collections::HashSet = self @@ -440,11 +1188,9 @@ impl Step for UploadToMaster { &self.ctx.deploy.deployment, ); - let tx_hosts: std::collections::HashSet = self - .ctx - .tx_host_ports() - .iter() - .filter_map(|hp| hp.split(':').next().map(|h| h.to_string())) + let second_batch_nodes = self.ctx.second_batch_nodes(); + let tx_hosts: HashSet = hosts_from_host_ports(&second_batch_nodes) + .into_iter() .collect(); let voter_hosts: std::collections::HashSet = self @@ -474,6 +1220,39 @@ impl Step for UploadToMaster { } } +pub struct UploadToAllNodes { + ctx: UpgradeContext, +} + +impl UploadToAllNodes { + pub fn new(ctx: UpgradeContext) -> Self { + Self { ctx } + } +} + +#[async_trait] +impl Step for UploadToAllNodes { + fn name(&self) -> &str { + "UploadToAllNodes" + } + + async fn build(&self) -> anyhow::Result { + let all_uploads = + crate::cli::task::upload::eloq_upload_builder::EloqUpload::eloq_image_upload( + &self.ctx.deploy.deployment, + ); + + let upload_tasks = crate::cli::task::upload::eloq_upload_builder::EloqUpload::build_tasks( + &self.ctx.config, + "update", + "upload_to_all_nodes", + all_uploads, + ); + + Ok(single_barrier_ctx("upload-to-all-nodes", upload_tasks)) + } +} + pub struct StopTxNodes { ctx: UpgradeContext, } @@ -526,6 +1305,84 @@ impl Step for StopTxNodes { } } +pub struct RestartNonLeaderNodes { + ctx: UpgradeContext, +} + +impl RestartNonLeaderNodes { + pub fn new(ctx: UpgradeContext) -> Self { + Self { ctx } + } +} + +#[async_trait] +impl Step for RestartNonLeaderNodes { + fn name(&self) -> &str { + "RestartNonLeaderNodes" + } + + async fn build(&self) -> anyhow::Result { + build_restart_nodes_sequence_ctx( + &self.ctx, + "restart-non-leader-nodes", + "restart-non-leader-node", + self.ctx.restart_nodes(), + ) + } +} + +pub struct RestartSelectedStandby { + ctx: UpgradeContext, +} + +impl RestartSelectedStandby { + pub fn new(ctx: UpgradeContext) -> Self { + Self { ctx } + } +} + +#[async_trait] +impl Step for RestartSelectedStandby { + fn name(&self) -> &str { + "RestartSelectedStandby" + } + + async fn build(&self) -> anyhow::Result { + build_restart_nodes_sequence_ctx( + &self.ctx, + "restart-selected-standby", + "restart-selected-standby", + self.ctx.temporary_leader_nodes(), + ) + } +} + +pub struct RestartTemporaryLeaders { + ctx: UpgradeContext, +} + +impl RestartTemporaryLeaders { + pub fn new(ctx: UpgradeContext) -> Self { + Self { ctx } + } +} + +#[async_trait] +impl Step for RestartTemporaryLeaders { + fn name(&self) -> &str { + "RestartTemporaryLeaders" + } + + async fn build(&self) -> anyhow::Result { + build_restart_nodes_sequence_ctx( + &self.ctx, + "restart-temporary-leaders", + "restart-temporary-leader", + self.ctx.temporary_leader_nodes(), + ) + } +} + pub struct StopLog { ctx: UpgradeContext, } @@ -597,10 +1454,16 @@ impl Step for CleanEloqStoreData { .and_then(|cc| cc.eloq_store_reuse_local_files) .unwrap_or(false); if !should_skip_cleanup { + let first_batch_hosts = + hosts_from_host_ports(&self.ctx.first_batch_nodes()); let clean_tasks = EloqStoreDataCleanTask::build_tasks( start_cmd, &self.ctx.config, - None, + if first_batch_hosts.is_empty() { + None + } else { + Some(first_batch_hosts.as_slice()) + }, ); if !clean_tasks.is_empty() { let len = clean_tasks.len(); @@ -692,26 +1555,18 @@ impl Step for StartTx { } async fn build(&self) -> anyhow::Result { - let mut start_tx = EloqTxCtlTask::from_config( - SubCommand::Start { - cluster: self.ctx.cluster.clone(), - nodes: Vec::new(), - }, - &self.ctx.deploy, - ServerType::Tx, - ); - - if self.ctx.has_voter() { - let start_voter = EloqTxCtlTask::from_config( + let start_tx = if self.ctx.has_standby() { + build_start_node_tasks(&self.ctx, "start-tx", self.ctx.second_batch_nodes()).executable + } else { + EloqTxCtlTask::from_config( SubCommand::Start { cluster: self.ctx.cluster.clone(), nodes: Vec::new(), }, &self.ctx.deploy, - ServerType::Voter, - ); - start_tx.extend(start_voter); - } + ServerType::Tx, + ) + }; Ok(single_barrier_ctx("start-tx", start_tx)) } @@ -742,6 +1597,7 @@ impl Step for WaitCurrentMaster { let (topology_tx, _) = watch::channel(ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), }); let mut executable = IndexMap::new(); executable.insert( @@ -796,60 +1652,25 @@ impl Step for WaitTxReplicaReady { if !self.ctx.has_standby() { return Ok(TaskExecutionContext::dummy()); } - - let tx_nodes = self.ctx.tx_host_ports(); - let standby_nodes = self.ctx.standby_host_ports(); - if tx_nodes.len() != standby_nodes.len() { + let targets = self.ctx.second_batch_nodes(); + let topology = + fetch_cluster_nodes(&self.ctx, "rolling-update-wait-second-batch-topology").await?; + let managed_nodes = self.ctx.managed_tx_and_standby_set(); + let current_masters = connected_managed_nodes(&topology.masters, &managed_nodes); + let Some(source_master) = current_masters.first() else { bail!( - "tx/standby node count mismatch: tx={}, standby={}", - tx_nodes.len(), - standby_nodes.len() - ); - } - let mut executable = IndexMap::new(); - - for (source_addr, target_addr) in standby_nodes.iter().zip(tx_nodes.iter()) { - let Some((source_host, source_port_str)) = source_addr.split_once(':') else { - bail!("invalid host:port in standby list: '{source_addr}'"); - }; - let Ok(source_port) = source_port_str.parse::() else { - bail!("invalid port in standby list: '{source_addr}'"); - }; - let Some((target_host, target_port_str)) = target_addr.split_once(':') else { - bail!("invalid host:port in tx list: '{target_addr}'"); - }; - let Ok(target_port) = target_port_str.parse::() else { - bail!("invalid port in tx list: '{target_addr}'"); - }; - let task_id = TaskId { - cmd: "topology".to_string(), - task: format!("wait-tx-replica-ready-{target_port}"), - host: target_host.to_string(), - }; - executable.insert( - task_id.clone(), - TaskInstance { - task_input: HashMap::default(), - task: Box::new( - WaitReplicaReadyTask::new( - task_id, - self.ctx.redis_cluster_startup_nodes(), - source_host.to_string(), - source_port, - target_host.to_string(), - target_port, - self.ctx.redis_password.clone(), - ) - .with_service_endpoints( - self.ctx.deploy.connection.service_endpoints.clone(), - ), - ), - task_host: TaskHost::Local, - }, + "rolling update could not find current master while waiting for updated old master replicas; topology reported masters={:?}, replicas={:?}", + topology.masters, + topology.replicas ); - } - - Ok(single_barrier_ctx("wait-tx-replica-ready", executable)) + }; + build_wait_replica_ready_tasks( + &self.ctx, + "wait-tx-replica-ready", + "wait-tx-replica-ready", + source_master, + &targets, + ) } } @@ -873,62 +1694,25 @@ impl Step for WaitStandbyReplicaReady { if !self.ctx.has_standby() { return Ok(TaskExecutionContext::dummy()); } - - let tx_nodes = self.ctx.tx_host_ports(); - let standby_nodes = self.ctx.standby_host_ports(); - if tx_nodes.len() != standby_nodes.len() { + let targets = self.ctx.first_batch_nodes(); + let topology = + fetch_cluster_nodes(&self.ctx, "rolling-update-wait-first-batch-topology").await?; + let managed_nodes = self.ctx.managed_tx_and_standby_set(); + let current_masters = connected_managed_nodes(&topology.masters, &managed_nodes); + let Some(source_master) = current_masters.first() else { bail!( - "tx/standby node count mismatch: tx={}, standby={}", - tx_nodes.len(), - standby_nodes.len() + "rolling update could not find current master while waiting for updated replica; topology reported masters={:?}, replicas={:?}", + topology.masters, + topology.replicas ); - } - - let mut executable = IndexMap::new(); - - for (source_addr, target_addr) in tx_nodes.iter().zip(standby_nodes.iter()) { - let Some((source_host, source_port_str)) = source_addr.split_once(':') else { - bail!("invalid host:port in tx list: '{source_addr}'"); - }; - let Ok(source_port) = source_port_str.parse::() else { - bail!("invalid port in tx list: '{source_addr}'"); - }; - let Some((target_host, target_port_str)) = target_addr.split_once(':') else { - bail!("invalid host:port in standby list: '{target_addr}'"); - }; - let Ok(target_port) = target_port_str.parse::() else { - bail!("invalid port in standby list: '{target_addr}'"); - }; - - let task_id = TaskId { - cmd: "topology".to_string(), - task: format!("wait-standby-replica-ready-{target_port}"), - host: target_host.to_string(), - }; - executable.insert( - task_id.clone(), - TaskInstance { - task_input: HashMap::default(), - task: Box::new( - WaitReplicaReadyTask::new( - task_id, - self.ctx.redis_cluster_startup_nodes(), - source_host.to_string(), - source_port, - target_host.to_string(), - target_port, - self.ctx.redis_password.clone(), - ) - .with_service_endpoints( - self.ctx.deploy.connection.service_endpoints.clone(), - ), - ), - task_host: TaskHost::Local, - }, - ); - } - - Ok(single_barrier_ctx("wait-standby-replica-ready", executable)) + }; + build_wait_replica_ready_tasks( + &self.ctx, + "wait-standby-replica-ready", + "wait-standby-replica-ready", + source_master, + &targets, + ) } } @@ -978,15 +1762,11 @@ impl Step for StartStandby { if !self.ctx.has_standby() { return Ok(TaskExecutionContext::dummy()); } - let start = EloqTxCtlTask::from_config( - SubCommand::Start { - cluster: self.ctx.cluster.clone(), - nodes: Vec::new(), - }, - &self.ctx.deploy, - ServerType::Standby, - ); - Ok(single_barrier_ctx("start-standby", start)) + Ok(build_start_node_tasks( + &self.ctx, + "start-standby", + self.ctx.first_batch_nodes(), + )) } } @@ -1094,25 +1874,51 @@ impl Step for VerifyVersion { /// Build the list of steps for a rolling binary upgrade (`eloqctl update`). /// -/// Strategy: update standby first, then failover, then update old master. -/// This minimizes downtime because the master continues serving during -/// the standby update phase. +/// Strategy: upload the new binaries, restart one connected standby for each +/// current master, wait for that standby to reconnect, fail over to the upgraded +/// standby, then restart the remaining tx/standby nodes one by one. The selected +/// standby is not restarted again after failover because it already runs the +/// upgraded binary. Voter nodes are intentionally left running because their +/// braft readiness is not exposed by Redis CLUSTER NODES. The final leader +/// movement is left to the cluster's normal election/preferred-leader logic. pub fn build_upgrade_steps(ctx: UpgradeContext) -> Vec> { - vec![ + if !ctx.has_standby() { + let mut steps: Vec> = vec![ + Box::new(DownloadAndExtract::new(ctx.clone())), + Box::new(UploadToAllNodes::new(ctx.clone())), + Box::new(StopTxNodes::new(ctx.clone())), + ]; + if !ctx.skip_log_restart { + steps.push(Box::new(StopLog::new(ctx.clone()))); + steps.push(Box::new(StartLogAndWait::new(ctx.clone()))); + } + steps.push(Box::new(StartTx::new(ctx.clone()))); + steps.push(Box::new(VerifyVersion::new(ctx))); + return steps; + } + + let mut steps: Vec> = vec![ Box::new(DownloadAndExtract::new(ctx.clone())), - Box::new(StopStandbyOnly::new(ctx.clone())), - Box::new(UploadToStandby::new(ctx.clone())), - Box::new(StopLog::new(ctx.clone())), - Box::new(CleanEloqStoreData::new(ctx.clone())), - Box::new(StartLogAndWait::new(ctx.clone())), - Box::new(StartStandby::new(ctx.clone())), - Box::new(WaitStandbyReplicaReady::new(ctx.clone())), - Box::new(FailoverAndStopOldMaster::new(ctx.clone())), - Box::new(UploadToMaster::new(ctx.clone())), - Box::new(StartTx::new(ctx.clone())), - Box::new(WaitTxReplicaReady::new(ctx.clone())), - Box::new(VerifyVersion::new(ctx)), - ] + Box::new(UploadToAllNodes::new(ctx.clone())), + Box::new(SelectStandbyForFailover::new(ctx.clone())), + Box::new(RestartSelectedStandby::new(ctx.clone())), + Box::new(FailoverToStandby::new(ctx.clone())), + ]; + + if ctx.skip_log_restart { + info!("Skipping log service restart during rolling update (--skip-log-restart)"); + } else { + steps.push(Box::new(StopLog::new(ctx.clone()))); + } + + if !ctx.skip_log_restart { + steps.push(Box::new(StartLogAndWait::new(ctx.clone()))); + } + + steps.push(Box::new(RestartNonLeaderNodes::new(ctx.clone()))); + steps.push(Box::new(VerifyVersion::new(ctx))); + + steps } /// Build the list of steps for a rolling config restart (`eloqctl update-conf --restart`). diff --git a/src/cluster_mgr/src/cli/task/task_utils.rs b/src/cluster_mgr/src/cli/task/task_utils.rs index 3afa5690..f44138bd 100644 --- a/src/cluster_mgr/src/cli/task/task_utils.rs +++ b/src/cluster_mgr/src/cli/task/task_utils.rs @@ -57,6 +57,8 @@ pub(crate) const PROCESS_PID: &str = "_process_pid_"; pub(crate) const PROCESS_PID_LIST: &str = "_process_pid_list_"; pub(crate) const DEFAULT_ELOQ_METRICS_PORT: u16 = 18081; pub(crate) const TX_INTERNAL_PORT_DELTA: u16 = 10000; +const CHECK_PID_RETRY_ATTEMPTS: usize = 3; +const CHECK_PID_RETRY_INTERVAL: Duration = Duration::from_secs(1); pub(crate) type NodeId = u32; pub(crate) type NodeGroupId = u32; @@ -159,9 +161,33 @@ where F: Fn(String) -> Option, T: Any + Debug, { - let mut cmd_exec_rs = ssh_session - .command(find_ps_cmd.as_str(), CollectOutput) - .await?; + let mut cmd_exec_rs = None; + let mut last_err = None; + for attempt in 1..=CHECK_PID_RETRY_ATTEMPTS { + match ssh_session + .command(find_ps_cmd.as_str(), CollectOutput) + .await + { + Ok(rs) => { + cmd_exec_rs = Some(rs); + break; + } + Err(err) => { + let err_msg = err.to_string(); + error!( + "check_pid command failed on attempt {}/{}: {}. cmd={}", + attempt, CHECK_PID_RETRY_ATTEMPTS, err_msg, find_ps_cmd + ); + last_err = Some(err_msg); + if attempt < CHECK_PID_RETRY_ATTEMPTS { + sleep(CHECK_PID_RETRY_INTERVAL).await; + } + } + } + } + let mut cmd_exec_rs = cmd_exec_rs.ok_or_else(|| { + anyhow!(last_err.unwrap_or_else(|| "check_pid command failed".to_string())) + })?; let cmd_status = cmd_exec_rs.get(CMD_STATUS).ok_or_else(|| { anyhow!( "check_pid failed: CMD_STATUS key missing from command result. \ @@ -414,6 +440,7 @@ pub async fn stop_with_hot_standby( let (tx_channel, rx_standby) = watch::channel::(ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), }); let rx_tx = tx_channel.subscribe(); @@ -600,6 +627,7 @@ pub async fn stop_with_failover( let (topology_tx, failover_rx) = watch::channel::(ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), }); // Create additional receivers that will get the same data diff --git a/src/cluster_mgr/src/cli/task/topology_update_task.rs b/src/cluster_mgr/src/cli/task/topology_update_task.rs index 3b7a3847..7d468226 100644 --- a/src/cluster_mgr/src/cli/task/topology_update_task.rs +++ b/src/cluster_mgr/src/cli/task/topology_update_task.rs @@ -118,6 +118,7 @@ impl TopologyUpdateTask { let (_, rx) = watch::channel(ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), }); let task = Box::new(TopologyUpdateTask { diff --git a/src/cluster_mgr/src/cli/task/tx_conf_update_task.rs b/src/cluster_mgr/src/cli/task/tx_conf_update_task.rs index f79cf15a..e260fa7d 100644 --- a/src/cluster_mgr/src/cli/task/tx_conf_update_task.rs +++ b/src/cluster_mgr/src/cli/task/tx_conf_update_task.rs @@ -461,6 +461,7 @@ impl TaskExecutor for TxConfUpdateTask { nodes: ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), }, cluster_config: Some(content), }; diff --git a/src/cluster_mgr/src/cli/task/wait_replica_ready_task.rs b/src/cluster_mgr/src/cli/task/wait_replica_ready_task.rs index d8b55036..277537fe 100644 --- a/src/cluster_mgr/src/cli/task/wait_replica_ready_task.rs +++ b/src/cluster_mgr/src/cli/task/wait_replica_ready_task.rs @@ -8,9 +8,11 @@ use std::collections::HashMap; use std::time::Duration; use tokio::sync::watch; use tokio::time::sleep; +use tracing::info; const REPLICA_READY_RETRIES: usize = 120; const REPLICA_READY_RETRY_DELAY: Duration = Duration::from_secs(2); +const WAIT_PROGRESS_LOG_INTERVAL: usize = 5; #[derive(Clone, Debug)] pub struct WaitReplicaReadyTask { @@ -24,6 +26,41 @@ pub struct WaitReplicaReadyTask { service_endpoints: Option>, } +#[derive(Clone, Debug)] +pub struct WaitNodeReadyTask { + task_id: TaskId, + startup_nodes: Vec, + target_host: String, + target_port: u16, + required_role: Option<&'static str>, + password: Option, + service_endpoints: Option>, +} + +fn describe_nodes(nodes: &[crate::cli::task::redis_op_task::NodeInfo]) -> String { + let described = nodes + .iter() + .map(|node| { + format!( + "{}:{} ({})", + node.ip, + node.port, + if node.connected { + "connected" + } else { + "disconnected" + } + ) + }) + .collect::>(); + + if described.is_empty() { + "".to_string() + } else { + described.join(", ") + } +} + impl WaitReplicaReadyTask { pub fn new( task_id: TaskId, @@ -81,6 +118,117 @@ impl WaitReplicaReadyTask { let (tx, _rx) = watch::channel(ClusterNodes { masters: Vec::new(), replicas: Vec::new(), + voters: Vec::new(), + }); + let result = RedisOpTask::new( + task_id, + self.startup_nodes.clone(), + "cluster topology".to_string(), + tx, + self.password.clone(), + true, + ) + .with_service_endpoints(self.service_endpoints.clone()) + .execute(TaskHost::Local, HashMap::default()) + .await?; + + let values = result.ok_or_else(|| anyhow::anyhow!("missing topology task result"))?; + let output = values + .get(CMD_OUTPUT) + .cloned() + .unwrap_or_else(|| TaskArgValue::Str("missing cluster topology output".to_string())); + let status = values + .get(CMD_STATUS) + .cloned() + .unwrap_or(TaskArgValue::Number(1)); + + match (status, output) { + (TaskArgValue::Number(0), TaskArgValue::Str(json)) => { + Ok(serde_json::from_str::(&json)?) + } + (_, TaskArgValue::Str(err)) => Err(anyhow::anyhow!(err)), + _ => Err(anyhow::anyhow!("unexpected topology task output")), + } + } +} + +impl WaitNodeReadyTask { + pub fn new( + task_id: TaskId, + startup_nodes: Vec, + target_host: String, + target_port: u16, + password: Option, + ) -> Self { + Self { + task_id, + startup_nodes, + target_host, + target_port, + required_role: None, + password, + service_endpoints: None, + } + } + + pub fn with_service_endpoints( + mut self, + service_endpoints: Option>, + ) -> Self { + self.service_endpoints = service_endpoints; + self + } + + pub fn require_master(mut self) -> Self { + self.required_role = Some("master"); + self + } + + fn find_connected_target<'a>( + &self, + cluster_nodes: &'a ClusterNodes, + ) -> Option<(&'static str, &'a crate::cli::task::redis_op_task::NodeInfo)> { + cluster_nodes + .masters + .iter() + .find(|node| { + node.ip == self.target_host && node.port == self.target_port && node.connected + }) + .map(|node| ("master", node)) + .or_else(|| { + cluster_nodes + .replicas + .iter() + .find(|node| { + node.ip == self.target_host + && node.port == self.target_port + && node.connected + }) + .map(|node| ("replica", node)) + }) + .or_else(|| { + cluster_nodes + .voters + .iter() + .find(|node| { + node.ip == self.target_host + && node.port == self.target_port + && node.connected + }) + .map(|node| ("voter", node)) + }) + } + + async fn fetch_cluster_nodes(&self) -> Result { + let task_id = TaskId { + cmd: "topology".to_string(), + task: format!("{}-topology", self.task_id.task), + host: "_local".to_string(), + }; + let (tx, _rx) = watch::channel(ClusterNodes { + masters: Vec::new(), + replicas: Vec::new(), + voters: Vec::new(), }); let result = RedisOpTask::new( task_id, @@ -135,41 +283,11 @@ impl TaskExecutor for WaitReplicaReadyTask { let mut last_seen = String::from("required nodes not yet observed as connected in cluster topology"); - for _ in 0..REPLICA_READY_RETRIES { + for attempt in 1..=REPLICA_READY_RETRIES { match self.fetch_cluster_nodes().await { Ok(cluster_nodes) => { - let masters = cluster_nodes - .masters - .iter() - .map(|node| { - format!( - "{}:{} ({})", - node.ip, - node.port, - if node.connected { - "connected" - } else { - "disconnected" - } - ) - }) - .collect::>(); - let replicas = cluster_nodes - .replicas - .iter() - .map(|node| { - format!( - "{}:{} ({})", - node.ip, - node.port, - if node.connected { - "connected" - } else { - "disconnected" - } - ) - }) - .collect::>(); + let masters = describe_nodes(&cluster_nodes.masters); + let replicas = describe_nodes(&cluster_nodes.replicas); if self.find_connected_master(&cluster_nodes).is_some() && self.find_connected_target_replica(&cluster_nodes).is_some() { @@ -178,34 +296,25 @@ impl TaskExecutor for WaitReplicaReadyTask { CMD_OUTPUT.to_string(), TaskArgValue::Str(format!( "Master {source} and replica {target} are connected and ready for failover. Masters: {}. Replicas: {}", - masters.join(", "), - replicas.join(", ") + masters, + replicas )), ); return Ok(Some(task_result)); } - last_seen = if masters.is_empty() && replicas.is_empty() { - "cluster currently reports no masters or replicas".to_string() - } else { - format!( - "masters currently visible: {}; replicas currently visible: {}", - if masters.is_empty() { - "".to_string() - } else { - masters.join(", ") - }, - if replicas.is_empty() { - "".to_string() - } else { - replicas.join(", ") - } - ) - }; + last_seen = format!( + "masters currently visible: {masters}; replicas currently visible: {replicas}" + ); } Err(err) => { last_seen = err.to_string(); } } + if attempt == 1 || attempt % WAIT_PROGRESS_LOG_INTERVAL == 0 { + info!( + "Waiting for master {source} and replica {target} to become connected ({attempt}/{REPLICA_READY_RETRIES}): {last_seen}" + ); + } sleep(REPLICA_READY_RETRY_DELAY).await; } @@ -219,3 +328,88 @@ impl TaskExecutor for WaitReplicaReadyTask { Ok(Some(task_result)) } } + +#[async_trait] +impl TaskExecutor for WaitNodeReadyTask { + fn identifier(&self) -> TaskId { + self.task_id.clone() + } + + async fn execute( + &self, + _task_host: TaskHost, + _task_arg: HashMap, + ) -> Result> { + let mut task_result = HashMap::from([( + CMD.to_string(), + TaskArgValue::Str("wait node ready".to_string()), + )]); + + let target = format!("{}:{}", self.target_host, self.target_port); + let mut last_seen = + String::from("target node not yet observed as connected in cluster topology"); + + for attempt in 1..=REPLICA_READY_RETRIES { + match self.fetch_cluster_nodes().await { + Ok(cluster_nodes) => { + let has_connected_master = + cluster_nodes.masters.iter().any(|node| node.connected); + if has_connected_master { + if let Some((role, _node)) = self.find_connected_target(&cluster_nodes) { + if self.required_role.is_some_and(|required| required != role) { + last_seen = format!( + "target node {target} is connected as {role}, waiting for {}", + self.required_role.unwrap() + ); + if attempt == 1 || attempt % WAIT_PROGRESS_LOG_INTERVAL == 0 { + info!( + "Waiting for node {target} to become connected as {} ({attempt}/{REPLICA_READY_RETRIES}): {last_seen}", + self.required_role.unwrap() + ); + } + sleep(REPLICA_READY_RETRY_DELAY).await; + continue; + } + task_result.insert(CMD_STATUS.to_string(), TaskArgValue::Number(0)); + task_result.insert( + CMD_OUTPUT.to_string(), + TaskArgValue::Str(format!( + "Node {target} is connected as {role}; cluster has a connected master" + )), + ); + return Ok(Some(task_result)); + } + } + last_seen = format!( + "masters: {}; replicas: {}; voters: {}", + describe_nodes(&cluster_nodes.masters), + describe_nodes(&cluster_nodes.replicas), + describe_nodes(&cluster_nodes.voters) + ); + } + Err(err) => { + last_seen = err.to_string(); + } + } + if attempt == 1 || attempt % WAIT_PROGRESS_LOG_INTERVAL == 0 { + let role = self + .required_role + .map(|role| format!(" as {role}")) + .unwrap_or_default(); + info!( + "Waiting for node {target} to become connected{role} ({attempt}/{REPLICA_READY_RETRIES}): {last_seen}" + ); + } + sleep(REPLICA_READY_RETRY_DELAY).await; + } + + task_result.insert(CMD_STATUS.to_string(), TaskArgValue::Number(1)); + task_result.insert( + CMD_OUTPUT.to_string(), + TaskArgValue::Str(format!( + "Node {target} did not become connected in cluster topology in time: {last_seen}" + )), + ); + Ok(Some(task_result)) + } +} diff --git a/src/cluster_mgr/src/config/deployment.rs b/src/cluster_mgr/src/config/deployment.rs index e24c43cb..f5910f3d 100644 --- a/src/cluster_mgr/src/config/deployment.rs +++ b/src/cluster_mgr/src/config/deployment.rs @@ -1506,22 +1506,10 @@ impl Deployment { } if self.log_service.is_some() { - // If MINIO is used, skip setting txlog_service_list but keep other log configs - let is_minio = self - .storage_service - .as_ref() - .and_then(|s| s.rocksdb.as_ref()) - .map(|r| matches!(r, RocksDB::MINIO(_))) - .unwrap_or(false); - self.build_log_config() .into_iter() .for_each(|(key, conf_val)| { - if is_minio && key == "txlog_service_list" { - // skip - } else { - ini.set(SECTION_CLUSTER, &key, Some(conf_val)); - } + ini.set(SECTION_CLUSTER, &key, Some(conf_val)); }); }