1use std::collections::BTreeMap;
2use std::collections::BTreeSet;
3
4use serde::Deserialize;
5use serde::Serialize;
6use ursula_shard::RaftGroupId;
7
8use crate::command::ControlCommand;
9use crate::command::ControlResponse;
10use crate::model::ClusterNode;
11use crate::model::DataGroupPlacement;
12use crate::model::GroupMigration;
13use crate::model::LearnerStatus;
14use crate::model::MetaConfig;
15use crate::model::MigrationPhase;
16use crate::model::NodeId;
17use crate::model::NodeState;
18use crate::view::GroupPlacementView;
19use crate::view::PlacementNode;
20
21#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
22pub struct ControlPlaneState {
23 pub nodes: BTreeMap<NodeId, ClusterNode>,
24 pub placements: BTreeMap<RaftGroupId, DataGroupPlacement>,
25 pub migrations: BTreeMap<u64, GroupMigration>,
26 pub active_migration: Option<u64>,
27 pub next_migration_id: u64,
28 pub config: MetaConfig,
29}
30
31impl Default for ControlPlaneState {
32 fn default() -> Self {
33 Self::new(MetaConfig::default())
34 }
35}
36
37impl ControlPlaneState {
38 pub fn new(config: MetaConfig) -> Self {
39 Self {
40 nodes: BTreeMap::new(),
41 placements: BTreeMap::new(),
42 migrations: BTreeMap::new(),
43 active_migration: None,
44 next_migration_id: 1,
45 config,
46 }
47 }
48
49 pub fn apply(&mut self, command: ControlCommand) -> ControlResponse {
50 match command {
51 ControlCommand::RegisterNode {
52 node_id,
53 client_url,
54 cluster_url,
55 labels,
56 now_ms,
57 } => self.register_node(node_id, client_url, cluster_url, labels, now_ms),
58 ControlCommand::SetNodeState {
59 node_id,
60 state,
61 now_ms,
62 } => self.set_node_state(node_id, state, now_ms),
63 ControlCommand::SeedPlacement {
64 raft_group_id,
65 voters,
66 now_ms,
67 } => self.seed_placement(raft_group_id, voters, now_ms),
68 ControlCommand::CommitPlacement {
69 raft_group_id,
70 voters,
71 learners,
72 draining,
73 now_ms,
74 } => self.commit_placement(raft_group_id, voters, learners, draining, now_ms),
75 ControlCommand::BeginMigration {
76 raft_group_id,
77 target_voters,
78 retain_removed,
79 now_ms,
80 } => self.begin_migration(raft_group_id, target_voters, retain_removed, now_ms),
81 ControlCommand::AdvanceMigration {
82 migration_id,
83 phase,
84 now_ms,
85 } => self.advance_migration(migration_id, phase, now_ms),
86 ControlCommand::SetLearnerStatus {
87 migration_id,
88 node_id,
89 status,
90 now_ms,
91 } => self.set_learner_status(migration_id, node_id, status, now_ms),
92 ControlCommand::RecordMigrationError {
93 migration_id,
94 error,
95 now_ms,
96 } => self.record_migration_error(migration_id, error, now_ms),
97 ControlCommand::FinishMigration {
98 migration_id,
99 success,
100 now_ms,
101 } => self.finish_migration(migration_id, success, now_ms),
102 ControlCommand::EvictLearner {
103 raft_group_id,
104 node_id,
105 now_ms,
106 } => self.evict_learner(raft_group_id, node_id, now_ms),
107 }
108 }
109
110 pub fn active_migration(&self) -> Option<&GroupMigration> {
111 self.active_migration
112 .and_then(|id| self.migrations.get(&id))
113 }
114
115 pub fn placement_view(&self, raft_group_id: RaftGroupId) -> Option<GroupPlacementView> {
116 let placement = self.placements.get(&raft_group_id)?;
117 let node_ids = placement
118 .voters
119 .iter()
120 .chain(placement.learners.iter())
121 .chain(placement.draining.iter());
122 let nodes = node_ids
123 .filter_map(|node_id| {
124 self.nodes.get(node_id).map(|node| {
125 (*node_id, PlacementNode {
126 node_id: *node_id,
127 client_url: node.client_url.clone(),
128 cluster_url: node.cluster_url.clone(),
129 state: node.state,
130 })
131 })
132 })
133 .collect();
134
135 Some(GroupPlacementView {
136 raft_group_id,
137 voters: placement.voters.clone(),
138 learners: placement.learners.clone(),
139 draining: placement.draining.clone(),
140 epoch: placement.epoch,
141 nodes,
142 })
143 }
144
145 fn register_node(
146 &mut self,
147 node_id: NodeId,
148 client_url: String,
149 cluster_url: String,
150 labels: BTreeMap<String, String>,
151 now_ms: u64,
152 ) -> ControlResponse {
153 let client_url = normalize_url(client_url);
154 let cluster_url = normalize_url(cluster_url);
155 if client_url.is_empty() {
156 return reject("client_url must not be empty".to_owned());
157 }
158 if cluster_url.is_empty() {
159 return reject("cluster_url must not be empty".to_owned());
160 }
161
162 let (registered_at_ms, state) =
163 self.nodes
164 .get(&node_id)
165 .map_or((now_ms, NodeState::Active), |node| {
166 let state = if node.state == NodeState::Removed {
167 NodeState::Active
168 } else {
169 node.state
170 };
171 (node.registered_at_ms, state)
172 });
173 self.nodes.insert(node_id, ClusterNode {
174 node_id,
175 client_url,
176 cluster_url,
177 state,
178 registered_at_ms,
179 updated_at_ms: now_ms,
180 labels,
181 });
182 ControlResponse::Ok
183 }
184
185 fn set_node_state(
186 &mut self,
187 node_id: NodeId,
188 state: NodeState,
189 now_ms: u64,
190 ) -> ControlResponse {
191 let Some(node) = self.nodes.get_mut(&node_id) else {
192 return reject(format!("node {node_id} is not registered"));
193 };
194 node.state = state;
195 node.updated_at_ms = now_ms;
196 ControlResponse::Ok
197 }
198
199 fn seed_placement(
200 &mut self,
201 raft_group_id: RaftGroupId,
202 voters: BTreeSet<NodeId>,
203 now_ms: u64,
204 ) -> ControlResponse {
205 if voters.is_empty() {
206 return reject("placement voters must not be empty".to_owned());
207 }
208
209 self.placements.insert(raft_group_id, DataGroupPlacement {
210 raft_group_id,
211 voters,
212 learners: BTreeSet::new(),
213 draining: BTreeSet::new(),
214 epoch: 0,
215 updated_at_ms: now_ms,
216 });
217 ControlResponse::Ok
218 }
219
220 fn commit_placement(
221 &mut self,
222 raft_group_id: RaftGroupId,
223 voters: BTreeSet<NodeId>,
224 learners: BTreeSet<NodeId>,
225 draining: BTreeSet<NodeId>,
226 now_ms: u64,
227 ) -> ControlResponse {
228 if voters.is_empty() {
229 return reject("placement voters must not be empty".to_owned());
230 }
231 if let Some(response) = self.validate_placement_nodes(&voters, &learners, &draining) {
232 return response;
233 }
234
235 let placement = self
236 .placements
237 .entry(raft_group_id)
238 .or_insert_with(|| DataGroupPlacement::empty(raft_group_id));
239 if placement.voters != voters {
240 placement.epoch = placement.epoch.saturating_add(1);
241 }
242 placement.voters = voters;
243 placement.learners = learners;
244 placement.draining = draining;
245 placement.updated_at_ms = now_ms;
246 ControlResponse::Ok
247 }
248
249 fn validate_placement_nodes(
250 &self,
251 voters: &BTreeSet<NodeId>,
252 learners: &BTreeSet<NodeId>,
253 draining: &BTreeSet<NodeId>,
254 ) -> Option<ControlResponse> {
255 if let Some(node_id) = voters.intersection(learners).next() {
256 return Some(reject(format!(
257 "node {node_id} cannot be both voter and learner"
258 )));
259 }
260 if let Some(response) = self.validate_registered_nodes("voter", voters, true) {
261 return Some(response);
262 }
263 if let Some(response) = self.validate_registered_nodes("learner", learners, false) {
264 return Some(response);
265 }
266 self.validate_registered_nodes("draining", draining, false)
267 }
268
269 fn validate_registered_nodes(
270 &self,
271 role: &str,
272 node_ids: &BTreeSet<NodeId>,
273 require_migration_eligible: bool,
274 ) -> Option<ControlResponse> {
275 for node_id in node_ids {
276 let Some(node) = self.nodes.get(node_id) else {
277 return Some(reject(format!("{role} node {node_id} is not registered")));
278 };
279 if require_migration_eligible && !node.state.is_migration_eligible() {
280 return Some(reject(format!(
281 "{role} node {node_id} is not migration eligible: {:?}",
282 node.state
283 )));
284 }
285 }
286 None
287 }
288
289 fn begin_migration(
290 &mut self,
291 raft_group_id: RaftGroupId,
292 target_voters: BTreeSet<NodeId>,
293 retain_removed: bool,
294 now_ms: u64,
295 ) -> ControlResponse {
296 if let Some(active) = self.active_migration {
297 return reject(format!("migration {active} is already running"));
298 }
299 if target_voters.is_empty() {
300 return reject("target voters must not be empty".to_owned());
301 }
302 for node_id in &target_voters {
303 let Some(node) = self.nodes.get(node_id) else {
304 return reject(format!("node {node_id} is not registered"));
305 };
306 if !node.state.is_migration_eligible() {
307 return reject(format!(
308 "node {node_id} is not migration eligible: {:?}",
309 node.state
310 ));
311 }
312 }
313
314 let Some(placement) = self.placements.get(&raft_group_id) else {
315 return reject(format!("group {} has no placement", raft_group_id.0));
316 };
317 let from_voters = placement.voters.clone();
318 let added_nodes = target_voters
319 .difference(&from_voters)
320 .copied()
321 .collect::<BTreeSet<_>>();
322 let removed_voters = from_voters
323 .difference(&target_voters)
324 .copied()
325 .collect::<BTreeSet<_>>();
326 let per_node_learner_status = added_nodes
327 .iter()
328 .copied()
329 .map(|node_id| (node_id, LearnerStatus::Pending))
330 .collect();
331
332 let migration_id = self.next_migration_id.max(1);
333 self.next_migration_id = migration_id.saturating_add(1);
334 self.migrations.insert(migration_id, GroupMigration {
335 migration_id,
336 raft_group_id,
337 from_voters,
338 target_voters,
339 added_nodes,
340 removed_voters,
341 retain_removed,
342 phase: MigrationPhase::Validating,
343 per_node_learner_status,
344 last_error: None,
345 retry_count: 0,
346 created_at_ms: now_ms,
347 updated_at_ms: now_ms,
348 });
349 self.active_migration = Some(migration_id);
350
351 ControlResponse::MigrationStarted { migration_id }
352 }
353
354 fn advance_migration(
355 &mut self,
356 migration_id: u64,
357 phase: MigrationPhase,
358 now_ms: u64,
359 ) -> ControlResponse {
360 if !phase.is_running() {
361 return reject(format!(
362 "migration {migration_id} must finish through FinishMigration"
363 ));
364 }
365 let Some(migration) = self.migrations.get_mut(&migration_id) else {
366 return reject(format!("migration {migration_id} does not exist"));
367 };
368 if migration.phase == phase {
369 return ControlResponse::Ok;
370 }
371 if !migration.phase.can_advance_to(phase) {
372 return reject(format!(
373 "migration {migration_id} cannot advance from {:?} to {:?}",
374 migration.phase, phase
375 ));
376 }
377 migration.phase = phase;
378 migration.updated_at_ms = now_ms;
379 ControlResponse::Ok
380 }
381
382 fn set_learner_status(
383 &mut self,
384 migration_id: u64,
385 node_id: NodeId,
386 status: LearnerStatus,
387 now_ms: u64,
388 ) -> ControlResponse {
389 let Some(migration) = self.migrations.get_mut(&migration_id) else {
390 return reject(format!("migration {migration_id} does not exist"));
391 };
392 if !migration.is_running() {
393 return reject(format!("migration {migration_id} is not running"));
394 }
395 if !migration.added_nodes.contains(&node_id) {
396 return reject(format!(
397 "node {node_id} is not an added learner for migration {migration_id}"
398 ));
399 }
400 migration.per_node_learner_status.insert(node_id, status);
401 migration.updated_at_ms = now_ms;
402 ControlResponse::Ok
403 }
404
405 fn record_migration_error(
406 &mut self,
407 migration_id: u64,
408 error: String,
409 now_ms: u64,
410 ) -> ControlResponse {
411 let Some(migration) = self.migrations.get_mut(&migration_id) else {
412 return reject(format!("migration {migration_id} does not exist"));
413 };
414 migration.last_error = Some(error);
415 migration.retry_count = migration.retry_count.saturating_add(1);
416 migration.updated_at_ms = now_ms;
417 ControlResponse::Ok
418 }
419
420 fn finish_migration(
421 &mut self,
422 migration_id: u64,
423 success: bool,
424 now_ms: u64,
425 ) -> ControlResponse {
426 let Some(migration) = self.migrations.get_mut(&migration_id) else {
427 return reject(format!("migration {migration_id} does not exist"));
428 };
429 if self.active_migration != Some(migration_id) {
430 return reject(format!("migration {migration_id} is not active"));
431 }
432 if !migration.is_running() {
433 return reject(format!("migration {migration_id} is not running"));
434 }
435 migration.phase = if success {
436 MigrationPhase::Succeeded
437 } else {
438 MigrationPhase::Failed
439 };
440 migration.updated_at_ms = now_ms;
441 self.active_migration = None;
442 ControlResponse::Ok
443 }
444
445 fn evict_learner(
446 &mut self,
447 raft_group_id: RaftGroupId,
448 node_id: NodeId,
449 now_ms: u64,
450 ) -> ControlResponse {
451 let Some(placement) = self.placements.get_mut(&raft_group_id) else {
452 return reject(format!("group {} has no placement", raft_group_id.0));
453 };
454 if placement.voters.contains(&node_id) {
455 return reject(format!(
456 "node {node_id} is a voter of group {} and cannot be evicted as a learner",
457 raft_group_id.0
458 ));
459 }
460 placement.learners.remove(&node_id);
461 placement.draining.remove(&node_id);
462 placement.updated_at_ms = now_ms;
463 ControlResponse::Ok
464 }
465}
466
467fn normalize_url(value: String) -> String {
468 value.trim().trim_end_matches('/').to_owned()
469}
470
471fn reject(reason: String) -> ControlResponse {
472 ControlResponse::Rejected { reason }
473}