qdrant_edge/shard/
segment_manifest.rs1use std::collections::HashMap;
10use std::sync::Arc;
11
12use crate::common::save_on_disk::SaveOnDisk;
13use crate::segment::common::operation_error::{OperationError, OperationResult};
14pub use crate::segment::segment_constructor::NewSegmentToken;
16use serde::{Deserialize, Serialize};
17use uuid::Uuid;
18
19use crate::shard::segment_holder::SegmentHolder;
20
21#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
28#[serde(rename_all = "snake_case")]
29pub enum SegmentManifestState {
30 Active,
32 UnderConstruction,
34 Optimizing {
37 holder: String,
39 lease_until: u64,
41 },
42 Retiring,
44}
45
46impl SegmentManifestState {
47 pub fn is_usable(&self) -> bool {
51 match self {
52 SegmentManifestState::Active
53 | SegmentManifestState::Optimizing {
54 holder: _,
55 lease_until: _,
56 } => true,
57 SegmentManifestState::UnderConstruction | SegmentManifestState::Retiring => false,
58 }
59 }
60
61 pub fn is_optimizer_mark(&self) -> bool {
66 match self {
67 SegmentManifestState::Optimizing {
68 holder: _,
69 lease_until: _,
70 }
71 | SegmentManifestState::Retiring => true,
72 SegmentManifestState::Active | SegmentManifestState::UnderConstruction => false,
73 }
74 }
75}
76
77#[derive(Clone, Debug, Default, Eq, PartialEq, Serialize, Deserialize)]
101#[serde(transparent)]
102pub struct SegmentsManifest {
103 segments: HashMap<Uuid, SegmentManifestState>,
104}
105
106impl SegmentsManifest {
107 pub fn from_segment_holder(holder: &SegmentHolder) -> Self {
112 let segments = holder
113 .iter()
114 .map(|(_, locked_segment)| {
115 let uuid = locked_segment.get_read().read().segment_uuid();
116 (uuid, SegmentManifestState::Active)
117 })
118 .collect();
119 Self { segments }
120 }
121
122 pub fn set(&mut self, uuid: Uuid, state: SegmentManifestState) -> Option<SegmentManifestState> {
124 self.segments.insert(uuid, state)
125 }
126
127 pub fn remove(&mut self, uuid: &Uuid) -> Option<SegmentManifestState> {
129 self.segments.remove(uuid)
130 }
131
132 #[must_use]
138 pub fn preserving(mut self, previous: &SegmentsManifest) -> Self {
139 for (uuid, state) in previous.iter() {
140 if state.is_optimizer_mark() && self.segments.contains_key(uuid) {
141 self.segments.insert(*uuid, state.clone());
142 }
143 }
144 self
145 }
146
147 pub fn sync(
157 manifest: Option<&Arc<SaveOnDisk<SegmentsManifest>>>,
158 holder: &SegmentHolder,
159 extra_segment: Option<Uuid>,
160 ) -> OperationResult<()> {
161 let Some(manifest) = manifest else {
162 return Ok(());
163 };
164
165 let mut rebuilt = Self::from_segment_holder(holder);
166 if let Some(uuid) = extra_segment {
167 rebuilt.set(uuid, SegmentManifestState::Active);
168 }
169
170 manifest
173 .write_optional(|previous| {
174 let current = rebuilt.preserving(previous);
175 (*previous != current).then_some(current)
176 })
177 .map_err(|err| {
178 OperationError::service_error(format!("failed to persist segment manifest: {err}"))
179 })?;
180 Ok(())
181 }
182
183 pub fn get(&self, uuid: &Uuid) -> Option<SegmentManifestState> {
184 self.segments.get(uuid).cloned()
185 }
186
187 pub fn iter(&self) -> impl Iterator<Item = (&Uuid, &SegmentManifestState)> {
188 self.segments.iter()
189 }
190
191 pub fn len(&self) -> usize {
192 self.segments.len()
193 }
194
195 pub fn is_empty(&self) -> bool {
196 self.segments.is_empty()
197 }
198}
199
200impl FromIterator<(Uuid, SegmentManifestState)> for SegmentsManifest {
201 fn from_iter<I: IntoIterator<Item = (Uuid, SegmentManifestState)>>(iter: I) -> Self {
202 Self {
203 segments: iter.into_iter().collect(),
204 }
205 }
206}
207
208#[cfg(test)]
209mod tests {
210 use super::*;
211
212 #[test]
213 fn serializes_active_as_snake_case_uuid_map() {
214 let uuid = Uuid::parse_str("1b4e28ba-2fa1-11d2-883f-0016d3cca427").unwrap();
215 let manifest = SegmentsManifest {
216 segments: HashMap::from([(uuid, SegmentManifestState::Active)]),
217 };
218
219 let json = serde_json::to_string(&manifest).unwrap();
220 assert_eq!(json, r#"{"1b4e28ba-2fa1-11d2-883f-0016d3cca427":"active"}"#);
221
222 let parsed: SegmentsManifest = serde_json::from_str(&json).unwrap();
223 assert_eq!(parsed, manifest);
224 assert_eq!(parsed.get(&uuid), Some(SegmentManifestState::Active));
225 }
226
227 #[test]
228 fn deserializes_future_states_for_forward_compat() {
229 let active = Uuid::parse_str("1b4e28ba-2fa1-11d2-883f-0016d3cca427").unwrap();
230 let building = Uuid::parse_str("6ba7b810-9dad-11d1-80b4-00c04fd430c8").unwrap();
231 let retiring = Uuid::parse_str("6ba7b811-9dad-11d1-80b4-00c04fd430c8").unwrap();
232
233 let json = format!(
234 r#"{{"{active}":"active","{building}":"under_construction","{retiring}":"retiring"}}"#,
235 );
236 let parsed: SegmentsManifest = serde_json::from_str(&json).unwrap();
237
238 assert_eq!(parsed.get(&active), Some(SegmentManifestState::Active));
239 assert_eq!(
240 parsed.get(&building),
241 Some(SegmentManifestState::UnderConstruction),
242 );
243 assert_eq!(parsed.get(&retiring), Some(SegmentManifestState::Retiring));
244 }
245
246 #[test]
247 fn optimizing_roundtrips_as_tagged_object() {
248 let uuid = Uuid::parse_str("1b4e28ba-2fa1-11d2-883f-0016d3cca427").unwrap();
249 let manifest: SegmentsManifest = [(
250 uuid,
251 SegmentManifestState::Optimizing {
252 holder: "indexer-1".to_string(),
253 lease_until: 1_752_000_000,
254 },
255 )]
256 .into_iter()
257 .collect();
258
259 let json = serde_json::to_string(&manifest).unwrap();
260 assert_eq!(
261 json,
262 r#"{"1b4e28ba-2fa1-11d2-883f-0016d3cca427":{"optimizing":{"holder":"indexer-1","lease_until":1752000000}}}"#,
263 );
264
265 let parsed: SegmentsManifest = serde_json::from_str(&json).unwrap();
266 assert_eq!(parsed, manifest);
267 }
268
269 #[test]
272 fn unit_states_keep_their_bare_string_form() {
273 let uuid = Uuid::parse_str("6ba7b811-9dad-11d1-80b4-00c04fd430c8").unwrap();
274 let manifest: SegmentsManifest = [(uuid, SegmentManifestState::Retiring)]
275 .into_iter()
276 .collect();
277 assert_eq!(
278 serde_json::to_string(&manifest).unwrap(),
279 r#"{"6ba7b811-9dad-11d1-80b4-00c04fd430c8":"retiring"}"#,
280 );
281 }
282
283 #[test]
284 fn preserving_keeps_in_progress_marks_for_listed_segments_only() {
285 let kept = Uuid::parse_str("1b4e28ba-2fa1-11d2-883f-0016d3cca427").unwrap();
286 let gone = Uuid::parse_str("6ba7b810-9dad-11d1-80b4-00c04fd430c8").unwrap();
287 let fresh = Uuid::parse_str("6ba7b811-9dad-11d1-80b4-00c04fd430c8").unwrap();
288
289 let optimizing = SegmentManifestState::Optimizing {
290 holder: "indexer-1".to_string(),
291 lease_until: 42,
292 };
293 let previous: SegmentsManifest = [
294 (kept, optimizing.clone()),
295 (gone, SegmentManifestState::Retiring),
296 ]
297 .into_iter()
298 .collect();
299
300 let rebuilt: SegmentsManifest = [
302 (kept, SegmentManifestState::Active),
303 (fresh, SegmentManifestState::Active),
304 ]
305 .into_iter()
306 .collect();
307
308 let merged = rebuilt.preserving(&previous);
309 assert_eq!(merged.get(&kept), Some(optimizing));
310 assert_eq!(merged.get(&fresh), Some(SegmentManifestState::Active));
311 assert_eq!(merged.get(&gone), None, "marks drop with the segment");
312 }
313}