Skip to main content

traverse_runtime/
data_store_coordinator.rs

1use std::sync::{Arc, Mutex};
2
3/// Host-facing, fenced write coordinator governed by Spec 093.
4#[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    /// Acquires the exclusive writer lease for `owner` and returns its generation.
30    ///
31    /// # Errors
32    ///
33    /// Returns [`DataStoreCoordinatorError::OwnerLocked`] when another owner holds
34    /// the lease, or [`DataStoreCoordinatorError::CoordinatorUnavailable`] when
35    /// the coordinator state is unavailable or cannot issue another generation.
36    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    /// Validates that `owner` still holds the supplied lease generation.
53    ///
54    /// # Errors
55    ///
56    /// Returns [`DataStoreCoordinatorError::OwnerLocked`] when the owner or
57    /// generation is stale, or [`DataStoreCoordinatorError::CoordinatorUnavailable`]
58    /// when the coordinator state is unavailable.
59    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    /// Releases the exclusive writer lease held by `owner` at `generation`.
72    ///
73    /// # Errors
74    ///
75    /// Returns [`DataStoreCoordinatorError::OwnerLocked`] when the owner or
76    /// generation is stale, or [`DataStoreCoordinatorError::CoordinatorUnavailable`]
77    /// when the coordinator state is unavailable.
78    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}