traverse_runtime/
data_store_coordinator.rs1use std::sync::{Arc, Mutex};
2
3#[derive(Clone, Debug)]
5pub struct DataStoreCoordinator {
6 state: Arc<Mutex<State>>,
7}
8
9#[derive(Debug, Default)]
10struct State {
11 generation: u64,
12 owner: Option<String>,
13}
14
15#[derive(Clone, Debug, PartialEq, Eq)]
16pub enum DataStoreCoordinatorError {
17 OwnerLocked,
18 CoordinatorUnavailable,
19}
20
21impl DataStoreCoordinator {
22 #[must_use]
23 pub fn new() -> Self {
24 Self {
25 state: Arc::new(Mutex::new(State::default())),
26 }
27 }
28
29 pub fn acquire(&self, owner: &str) -> Result<u64, DataStoreCoordinatorError> {
37 let mut state = self
38 .state
39 .lock()
40 .map_err(|_| DataStoreCoordinatorError::CoordinatorUnavailable)?;
41 if state.owner.is_some() {
42 return Err(DataStoreCoordinatorError::OwnerLocked);
43 }
44 state.generation = state
45 .generation
46 .checked_add(1)
47 .ok_or(DataStoreCoordinatorError::CoordinatorUnavailable)?;
48 state.owner = Some(owner.to_owned());
49 Ok(state.generation)
50 }
51
52 pub fn validate(&self, owner: &str, generation: u64) -> Result<(), DataStoreCoordinatorError> {
60 let state = self
61 .state
62 .lock()
63 .map_err(|_| DataStoreCoordinatorError::CoordinatorUnavailable)?;
64 if state.owner.as_deref() == Some(owner) && state.generation == generation {
65 Ok(())
66 } else {
67 Err(DataStoreCoordinatorError::OwnerLocked)
68 }
69 }
70
71 pub fn release(&self, owner: &str, generation: u64) -> Result<(), DataStoreCoordinatorError> {
79 let mut state = self
80 .state
81 .lock()
82 .map_err(|_| DataStoreCoordinatorError::CoordinatorUnavailable)?;
83 if state.owner.as_deref() != Some(owner) || state.generation != generation {
84 return Err(DataStoreCoordinatorError::OwnerLocked);
85 }
86 state.owner = None;
87 Ok(())
88 }
89}
90
91impl Default for DataStoreCoordinator {
92 fn default() -> Self {
93 Self::new()
94 }
95}
96
97#[cfg(test)]
98mod tests {
99 use super::*;
100
101 #[test]
102 fn fences_stale_owner_after_takeover() {
103 let coordinator = DataStoreCoordinator::new();
104 let first = 1;
105 assert_eq!(coordinator.acquire("first"), Ok(first));
106 assert_eq!(coordinator.release("first", first), Ok(()));
107 let second = 2;
108 assert_eq!(coordinator.acquire("second"), Ok(second));
109 assert!(second > first);
110 assert_eq!(
111 coordinator.validate("first", first),
112 Err(DataStoreCoordinatorError::OwnerLocked)
113 );
114 }
115
116 #[test]
117 fn rejects_contending_and_invalid_owner_requests() {
118 let coordinator = DataStoreCoordinator::new();
119 let generation = 1;
120 assert_eq!(coordinator.acquire("owner"), Ok(generation));
121 assert_eq!(
122 coordinator.acquire("contender"),
123 Err(DataStoreCoordinatorError::OwnerLocked)
124 );
125 assert_eq!(
126 coordinator.release("contender", generation),
127 Err(DataStoreCoordinatorError::OwnerLocked)
128 );
129 }
130
131 #[test]
132 fn default_coordinator_validates_its_active_owner() {
133 let coordinator = DataStoreCoordinator::default();
134 let generation = 1;
135 assert_eq!(coordinator.acquire("owner"), Ok(generation));
136 assert_eq!(coordinator.validate("owner", generation), Ok(()));
137 }
138
139 #[test]
140 fn reports_unavailable_when_the_state_lock_is_poisoned() {
141 let coordinator = DataStoreCoordinator::new();
142 poison(&coordinator);
143
144 assert_eq!(
145 coordinator.acquire("owner"),
146 Err(DataStoreCoordinatorError::CoordinatorUnavailable)
147 );
148 assert_eq!(
149 coordinator.validate("owner", 1),
150 Err(DataStoreCoordinatorError::CoordinatorUnavailable)
151 );
152 assert_eq!(
153 coordinator.release("owner", 1),
154 Err(DataStoreCoordinatorError::CoordinatorUnavailable)
155 );
156 poison(&coordinator);
157 }
158
159 #[test]
160 fn reports_unavailable_when_generations_are_exhausted() {
161 let coordinator = DataStoreCoordinator::new();
162 set_generation(&coordinator, u64::MAX);
163
164 assert_eq!(
165 coordinator.acquire("owner"),
166 Err(DataStoreCoordinatorError::CoordinatorUnavailable)
167 );
168 }
169
170 #[test]
171 fn generation_setup_tolerates_an_unavailable_coordinator() {
172 let coordinator = DataStoreCoordinator::new();
173 poison(&coordinator);
174 set_generation(&coordinator, u64::MAX);
175
176 assert_eq!(
177 coordinator.acquire("owner"),
178 Err(DataStoreCoordinatorError::CoordinatorUnavailable)
179 );
180 }
181
182 fn set_generation(coordinator: &DataStoreCoordinator, generation: u64) {
183 let Ok(mut state) = coordinator.state.lock() else {
184 return;
185 };
186 state.generation = generation;
187 }
188
189 fn poison(coordinator: &DataStoreCoordinator) {
190 let state = Arc::clone(&coordinator.state);
191 let result = std::thread::spawn(move || {
192 let Ok(_guard) = state.lock() else {
193 return;
194 };
195 std::panic::resume_unwind(Box::new(()));
196 })
197 .join();
198 assert!(result.is_err() || coordinator.state.lock().is_err());
199 }
200}