Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

6 changes: 6 additions & 0 deletions src/cluster_mgr/src/cli/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
1 change: 1 addition & 0 deletions src/cluster_mgr/src/cli/task/db_update_task.rs
Original file line number Diff line number Diff line change
Expand Up @@ -281,6 +281,7 @@ impl TaskExecutor for DbDeploymentUpdateTask {
nodes: ClusterNodes {
masters: Vec::new(),
replicas: Vec::new(),
voters: Vec::new(),
},
cluster_config: Some(content),
}
Expand Down
28 changes: 23 additions & 5 deletions src/cluster_mgr/src/cli/task/eloq_tx_ctl_task.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand Down Expand Up @@ -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<String> = Vec::new();
match server_type {
"txservice" => {
Expand Down Expand Up @@ -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::<u16>().ok())
Expand Down Expand Up @@ -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),
Expand Down
71 changes: 64 additions & 7 deletions src/cluster_mgr/src/cli/task/failover_op_task.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ pub struct FailoverOpTask {
receiver: watch::Receiver<ClusterNodes>,
password: Option<String>,
service_endpoints: Option<HashMap<String, ServiceEndpoint>>,
require_explicit_target: bool,
}

impl FailoverOpTask {
Expand All @@ -43,6 +44,7 @@ impl FailoverOpTask {
receiver,
password,
service_endpoints: None,
require_explicit_target: false,
}
}

Expand All @@ -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));
}

Expand Down
2 changes: 2 additions & 0 deletions src/cluster_mgr/src/cli/task/group/db_cluster_ctrl_group.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand Down Expand Up @@ -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));
Expand Down
2 changes: 2 additions & 0 deletions src/cluster_mgr/src/cli/task/group/failover_group.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,13 +52,15 @@ fn failover_task_group(
let (pre_failover_tx, pre_failover_rx) = watch::channel::<ClusterNodes>(ClusterNodes {
masters: Vec::new(),
replicas: Vec::new(),
voters: Vec::new(),
});

// Create another channel for the post-failover task
// This time, we'll keep a reference to both sender and receiver to ensure the channel remains open
let (post_failover_tx, post_failover_rx) = watch::channel::<ClusterNodes>(ClusterNodes {
masters: Vec::new(),
replicas: Vec::new(),
voters: Vec::new(),
});

// Keep the receiver alive to prevent the channel from closing
Expand Down
1 change: 1 addition & 0 deletions src/cluster_mgr/src/cli/task/group/launch_group.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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());

Expand Down
2 changes: 2 additions & 0 deletions src/cluster_mgr/src/cli/task/group/scale_group.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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());

Expand Down
1 change: 1 addition & 0 deletions src/cluster_mgr/src/cli/task/group/scale_log_group.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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());
Expand Down
1 change: 1 addition & 0 deletions src/cluster_mgr/src/cli/task/group/update_config_group.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Loading
Loading