Skip to main content

ursula_control/
state.rs

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}