Skip to main content

qdrant_edge/shard/
segment_manifest.rs

1//! The segment manifest (`segments_manifest.json`, sitting next to the `segments/` directory): a
2//! small file listing the shard's segments and
3//! their state, so out-of-process readers (e.g. a read-only follower, possibly over object storage)
4//! can discover segments without scanning the filesystem.
5//!
6//! Defines the on-disk structure and a helper to build it from a [`SegmentHolder`]. The *logic* of
7//! when to write it lives in the owner of the segments (e.g. `LocalShard`).
8
9use std::collections::HashMap;
10use std::sync::Arc;
11
12use crate::common::save_on_disk::SaveOnDisk;
13use crate::segment::common::operation_error::{OperationError, OperationResult};
14// Re-exported from `segment` (where `build_segment` mints it) since the manifest is what consumes it.
15pub use crate::segment::segment_constructor::NewSegmentToken;
16use serde::{Deserialize, Serialize};
17use uuid::Uuid;
18
19use crate::shard::segment_holder::SegmentHolder;
20
21/// State of a segment in the manifest.
22///
23/// The shard itself only writes [`Active`](SegmentManifestState::Active); the in-progress states
24/// come from out-of-process rewriters (the serverless indexer). A data-carrying state serializes
25/// as a tagged object, which readers predating it cannot parse — write it only once every reader
26/// understands it.
27#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
28#[serde(rename_all = "snake_case")]
29pub enum SegmentManifestState {
30    /// Live segment, part of the shard's data and safe to read.
31    Active,
32    /// Reserved for future use: a segment being built that is not yet ready to read.
33    UnderConstruction,
34    /// Being rebuilt by an out-of-process optimizer; still live and readable. The claim is a
35    /// lease: past `lease_until` another optimizer may take the segment over.
36    Optimizing {
37        /// Identity of the optimizer run holding the claim (e.g. pod name).
38        holder: String,
39        /// Unix seconds after which the claim is stale.
40        lease_until: u64,
41    },
42    /// A superseded segment that is still readable but pending removal.
43    Retiring,
44}
45
46impl SegmentManifestState {
47    /// Whether a reader may serve the segment: live ([`Active`](Self::Active)), or
48    /// [`Optimizing`](Self::Optimizing) — being rebuilt out-of-process, but live (and still
49    /// receiving deletes) until the swap.
50    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    /// Whether the state is a mark owned by an out-of-process optimizer — a rebuild claim
62    /// ([`Optimizing`](Self::Optimizing)) or a pending removal ([`Retiring`](Self::Retiring)) —
63    /// rather than something derivable from the segment holder. A holder rebuild must not erase
64    /// such marks (see [`SegmentsManifest::preserving`]).
65    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/// Contents of `segments_manifest.json`: a flat map of segment UUID to its state, e.g.
78/// `{ "1b4e28ba-...": "active", "6ba7b810-...": "active" }`.
79///
80/// # Consistency assumptions (important for readers)
81///
82/// The manifest is *not* an exact, instantaneous snapshot of the on-disk segment set. To guarantee
83/// that no data is ever lost across the non-atomic create/delete transitions, it deliberately errs
84/// on the side of listing more segments than strictly necessary:
85///
86/// - **May list a segment that is not yet fully created.** A new segment is registered here as soon
87///   as it exists on disk — potentially *before* its version file is written (i.e. before it is
88///   finalized). This guarantees a real, writable segment is never missing from the manifest. A
89///   reader must therefore tolerate a listed segment that is incomplete / not yet loadable and skip
90///   it (a segment without a saved version is not ready to read).
91/// - **May list a segment that has already been (or is about to be) deleted.** Removal from the
92///   manifest is intentionally not strict: a superseded segment may linger in the manifest after its
93///   data has been dropped, or be dropped from disk slightly after it leaves the manifest. A reader
94///   must tolerate a listed segment that no longer exists on disk and skip it.
95///
96/// In other words: the manifest is a *superset-biased* view. Over-listing (duplicate/extra segments)
97/// is always safe and is resolved by point versioning during reads; under-listing a live segment is
98/// the only thing that would lose data, and the writer ordering prevents it. Readers must handle
99/// both of the above cases gracefully rather than assuming every listed segment is present and ready.
100#[derive(Clone, Debug, Default, Eq, PartialEq, Serialize, Deserialize)]
101#[serde(transparent)]
102pub struct SegmentsManifest {
103    segments: HashMap<Uuid, SegmentManifestState>,
104}
105
106impl SegmentsManifest {
107    /// Build a manifest from the current live segments in the holder, marking every segment
108    /// [`Active`](SegmentManifestState::Active).
109    ///
110    /// Includes segments temporarily wrapped as proxies during optimization — they are still live.
111    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    /// Set `uuid`'s state, returning the previous one.
123    pub fn set(&mut self, uuid: Uuid, state: SegmentManifestState) -> Option<SegmentManifestState> {
124        self.segments.insert(uuid, state)
125    }
126
127    /// Remove `uuid`'s entry, returning its state.
128    pub fn remove(&mut self, uuid: &Uuid) -> Option<SegmentManifestState> {
129        self.segments.remove(uuid)
130    }
131
132    /// Merge rule for persisting a rebuild: a rebuild from the segment holder marks every live
133    /// segment `Active`, which would erase an out-of-process optimizer's marks. Re-apply
134    /// `previous`'s optimizer marks (`Optimizing`/`Retiring`) to segments this manifest still
135    /// lists; marks for segments gone from the holder drop with them. Staleness is governed by
136    /// the lease, not by rebuilds.
137    #[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    /// Rebuild the manifest from `holder` and persist it if it changed.
148    ///
149    /// No-op when `manifest` is `None` (the `write_segment_manifest` flag is off). Errors
150    /// propagate so callers that gate destructive work (e.g. deleting superseded segments) on a
151    /// fresh manifest can abort instead of risking a stale manifest.
152    ///
153    /// Idempotent and cheap: it only writes when the live segment set differs from what is
154    /// persisted, so it is safe to call on every relevant segment-set transition. Usually invoked
155    /// via [`SegmentHolder::sync_segment_manifest`], which owns the manifest.
156    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        // Merge under the write lock, so a concurrent rebuild cannot slip between a
171        // read and a write and erase preserved marks.
172        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    /// Byte-compat guard: today's writers produce bare strings for the unit states; adding
270    /// the data-carrying variant must not change how those serialize.
271    #[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        // A holder rebuild lists every live segment as Active; `gone` left the holder.
301        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}