Skip to main content

mkit_server/takedown/
discovery.rs

1//! Checkpointed holder discovery over configured or provable namespace roots.
2use crate::pipeline::{MAX_APPLY_WINDOW, ShardMap};
3use crate::store::{BlobKey, Cursor, StoreError, codec, keys, watermark};
4use crate::{
5    Addressing, Batch, Key, NamespaceKey, NamespaceStore, Partition, RepoId, RepoName, Value,
6};
7use mkit_core::hash::Hash;
8use serde::{Deserialize, Serialize};
9use watermark::namespace_relay_watermark_step;
10type RecoveryMarkers = (Option<Vec<u8>>, Option<Vec<u8>>);
11
12/// Durable checkpoint; commit it atomically with the returned context batch.
13#[derive(Debug, Clone, Default, Serialize, Deserialize)]
14#[serde(deny_unknown_fields)]
15pub struct DiscoveryState {
16    version: u8,
17    namespaces: Vec<String>,
18    exhaustive: bool,
19    addressing_single: bool,
20    single_repo: Option<(String, String)>,
21    safety_cut: u64,
22    namespace: usize,
23    phase: u8,
24    registry_cursor: Option<Vec<u8>>,
25    shard_cursor: Option<Vec<u8>>,
26    object_cursor: Option<Vec<u8>>,
27    repo: Option<String>,
28    watermark: Option<Vec<u8>>,
29    generation: Option<RecoveryMarkers>,
30    binding: Option<(Hash, Hash, bool)>,
31}
32fn bad() -> StoreError {
33    StoreError::Corrupt("invalid takedown discovery checkpoint/source".into())
34}
35impl DiscoveryState {
36    /// Bind namespace configuration, the named source and the safety cut.
37    pub fn new(
38        addressing: &Addressing,
39        named_namespace: &NamespaceKey,
40        created: u64,
41        margin: u64,
42    ) -> Result<Self, StoreError> {
43        let exhaustive = !matches!(addressing, Addressing::Multi(multi) if matches!(multi.namespace_policy, crate::policy::NamespacePolicy::Any { .. }));
44        let namespaces = if exhaustive {
45            super::configured_namespaces(addressing)
46        } else {
47            vec![named_namespace.clone()]
48        };
49        if namespaces.is_empty() || margin == 0 {
50            return Err(bad());
51        }
52        let mut roots: Vec<_> = namespaces
53            .into_iter()
54            .map(|ns| ns.as_str().to_owned())
55            .collect();
56        roots.sort();
57        roots.dedup();
58        if exhaustive && !roots.iter().any(|root| root == named_namespace.as_str()) {
59            return Err(bad());
60        }
61        let single_repo = match addressing {
62            Addressing::Single { repo } => Some((
63                repo.namespace.as_str().to_owned(),
64                repo.name.as_str().to_owned(),
65            )),
66            Addressing::Multi(_) => None,
67        };
68        Ok(Self {
69            version: 1,
70            namespaces: roots,
71            exhaustive,
72            addressing_single: single_repo.is_some(),
73            single_repo,
74            safety_cut: created
75                .checked_add(u64::try_from(MAX_APPLY_WINDOW.as_millis()).map_err(|_| bad())?)
76                .and_then(|v| v.checked_add(margin))
77                .ok_or_else(bad)?,
78            ..Self::default()
79        })
80    }
81    /// Whether the accepted namespace universe is exhaustive.
82    #[must_use]
83    pub fn exhaustive(&self) -> bool {
84        self.exhaustive
85    }
86    /// Whether the current namespace candidate traversal finished.
87    #[must_use]
88    pub fn traversed(&self) -> bool {
89        self.namespace == self.namespaces.len()
90    }
91    /// Resume one streamed known-holder namespace under the original cut and binding.
92    /// The caller owns its durable namespace-page cursor; repeats are idempotent.
93    pub fn next_candidate(&mut self, namespace: &NamespaceKey) -> Result<(), StoreError> {
94        if self.exhaustive || !self.traversed() {
95            return Err(bad());
96        }
97        self.namespaces = vec![namespace.as_str().to_owned()];
98        self.namespace = 0;
99        Ok(())
100    }
101}
102async fn member_context<S: NamespaceStore>(
103    store: &S,
104    shards: &dyn ShardMap,
105    repo: &RepoId,
106    action: &Hash,
107    object: &Hash,
108    pack_id: &Hash,
109) -> Result<Option<(Key, Value)>, StoreError> {
110    if let Some(member) = store
111        .get(
112            &shards.membership(repo, &BlobKey::pack(*pack_id)),
113            &keys::membership(&repo.name, pack_id),
114        )
115        .await?
116    {
117        if !member.as_bytes().is_empty() {
118            return Err(bad());
119        }
120        let identity = format!("{}\0{}\0", repo.namespace.as_str(), repo.name.as_str());
121        let key = Key::new(
122            [
123                b"b\0\xffdiscovery-context\0".as_slice(),
124                action,
125                object,
126                identity.as_bytes(),
127                pack_id,
128            ]
129            .concat(),
130        );
131        let value = serde_json::to_vec(&serde_json::json!({"version":1,"pack":pack_id,"known_signers":[],
132                        "signer_metadata":"unavailable_in_existing_source","context_complete":false})).map_err(|_| bad())?;
133        return Ok(Some((key, Value::new(value))));
134    }
135    Ok(None)
136}
137
138/// One bounded step; a successful read is not durable until the caller commits.
139#[derive(Debug)]
140pub struct DiscoveryStep {
141    /// Checkpoint to commit alongside `batch` under the workflow's CAS guard.
142    pub state: DiscoveryState,
143    /// Context writes for `root`; this function performs no storage mutations.
144    pub batch: Batch,
145    /// Repository with confirmed membership in this bounded step.
146    pub repository: Option<RepoId>,
147    /// All exhaustive roots passed watermarks and applicable configured-repository
148    /// or registry/active-shard enumeration.
149    /// Any always returns false; this never claims verified preservation completion.
150    pub complete: bool,
151    /// This candidate traversal finished; Any remains non-exhaustive and incomplete.
152    pub traversed: bool,
153}
154/// One watermark page, repository candidate or object-index row per call.
155/// Use the alarm's budgeted/local-aware store; never explicitly route a self call.
156#[allow(
157    clippy::too_many_lines,
158    clippy::too_many_arguments,
159    reason = "Keep the mutually exclusive bounded checkpoint transitions together."
160)]
161pub async fn step<S: NamespaceStore>(
162    store: &S,
163    shards: &dyn ShardMap,
164    _root: &Partition,
165    action: &Hash,
166    object: &Hash,
167    is_pack: bool,
168    now: u64,
169    mut state: DiscoveryState,
170) -> Result<DiscoveryStep, StoreError> {
171    if state.version != 1
172        || state.namespaces.is_empty()
173        || state.namespace > state.namespaces.len()
174        || state.phase > 3
175        || state.addressing_single != state.single_repo.is_some()
176        || (state.phase > 0 && (state.generation.is_none() || state.binding.is_none()))
177        || state
178            .binding
179            .is_some_and(|b| b != (*action, *object, is_pack))
180    {
181        return Err(bad());
182    }
183    state.binding = Some((*action, *object, is_pack));
184    let mut batch = Batch::new();
185    let mut repository = None;
186    if state.namespace < state.namespaces.len() && now > state.safety_cut {
187        let ns = NamespaceKey::from_stored(state.namespaces[state.namespace].clone());
188        let coordinator = shards.coordinator(&ns);
189        if state.addressing_single
190            && state.single_repo.as_ref().is_none_or(|(namespace, _)| {
191                state.namespaces.len() != 1 || *namespace != ns.as_str()
192            })
193        {
194            return Err(bad());
195        }
196        let generation = watermark::check_recovery(store, &coordinator, None)
197            .await
198            .map_err(|_| bad())?;
199        let generation = (
200            generation.0.map(|v| v.as_bytes().to_vec()),
201            generation.1.map(|v| v.as_bytes().to_vec()),
202        );
203        if state
204            .generation
205            .as_ref()
206            .is_some_and(|old| old != &generation)
207        {
208            state.phase = 0;
209            state.registry_cursor = None;
210            state.shard_cursor = None;
211            state.object_cursor = None;
212            state.repo = None;
213            state.watermark = None;
214            state.generation = None;
215            return Ok(DiscoveryStep {
216                state,
217                batch,
218                repository,
219                complete: false,
220                traversed: false,
221            });
222        }
223        state.generation = Some(generation);
224        if let Some(name) = state.repo.clone() {
225            let repo = RepoId {
226                namespace: ns,
227                name: RepoName::new(name).map_err(|_| bad())?,
228            };
229            if is_pack {
230                if let Some((key, value)) =
231                    member_context(store, shards, &repo, action, object, object).await?
232                {
233                    batch = batch.put(key, value);
234                    repository = Some(repo.clone());
235                }
236                state.repo = None;
237            } else {
238                let (start, end) = keys::object_index_range(&repo.name, object);
239                let cursor = state.object_cursor.clone().map(Cursor::new);
240                let page = store
241                    .scan(
242                        &shards.object_index(&repo, object),
243                        &start,
244                        &end,
245                        cursor.as_ref(),
246                        1,
247                    )
248                    .await?;
249                for (key, value) in page.entries {
250                    let Some(keys::ParsedKey::ObjectIndex {
251                        repo: indexed_repo,
252                        object: indexed_object,
253                        pack_id,
254                    }) = keys::parse(&key)
255                    else {
256                        return Err(bad());
257                    };
258                    if indexed_repo != repo.name || indexed_object != *object {
259                        return Err(bad());
260                    }
261                    codec::decode_object_index(object, &value)?;
262                    if let Some((key, value)) =
263                        member_context(store, shards, &repo, action, object, &pack_id).await?
264                    {
265                        batch = batch.put(key, value);
266                        repository = Some(repo.clone());
267                    }
268                }
269                state.object_cursor = page.next.map(|c| c.as_bytes().to_vec());
270                if state.object_cursor.is_none() {
271                    state.repo = None;
272                }
273            }
274        } else if state.phase == 0 {
275            let value = if matches!(coordinator, Partition::Namespace(_)) {
276                Some(crate::relay::relay_watermark(store, &coordinator, now).await?)
277            } else {
278                let checkpoint = state
279                    .watermark
280                    .as_deref()
281                    .map(watermark::WatermarkCheckpoint::decode)
282                    .transpose()?;
283                match namespace_relay_watermark_step(store, &coordinator, now, checkpoint, 1)
284                    .await
285                    .map_err(|_| bad())?
286                {
287                    watermark::WatermarkStep::Pending(next) => {
288                        state.watermark = Some(next.encode());
289                        None
290                    }
291                    watermark::WatermarkStep::Complete(value) => {
292                        state.watermark = None;
293                        Some(value)
294                    }
295                }
296            };
297            if value.is_some_and(|v| v > state.safety_cut) {
298                state.phase = if state.addressing_single { 3 } else { 1 };
299                if state.addressing_single {
300                    state.repo = state.single_repo.as_ref().map(|(_, name)| name.clone());
301                }
302            }
303        } else if state.phase == 1 {
304            let (start, end) = keys::class_range(keys::TAG_REPO_REGISTRY);
305            let cursor = state.registry_cursor.clone().map(Cursor::new);
306            let page = store
307                .scan(&coordinator, &start, &end, cursor.as_ref(), 1)
308                .await?;
309            for (key, value) in page.entries {
310                let Some(keys::ParsedKey::RepoRecord(repo)) = keys::parse(&key) else {
311                    return Err(bad());
312                };
313                codec::decode_repo_record(&value)?;
314                state.repo = Some(repo.as_str().to_owned());
315            }
316            state.registry_cursor = page.next.map(|c| c.as_bytes().to_vec());
317            if state.registry_cursor.is_none() {
318                state.phase = 2;
319            }
320        } else if state.phase == 2 {
321            if matches!(coordinator, Partition::Namespace(_)) {
322                state.phase = 3;
323            } else {
324                let cursor = state.shard_cursor.clone().map(Cursor::new);
325                let page = watermark::active_shards(store, &coordinator, cursor.as_ref(), 1)
326                    .await
327                    .map_err(|_| bad())?;
328                for partition in page.shards {
329                    let Partition::Ref { repo, .. } = partition else {
330                        return Err(bad());
331                    };
332                    state.repo = Some(repo.as_str().to_owned());
333                }
334                state.shard_cursor = page.next.map(|c| c.as_bytes().to_vec());
335                if state.shard_cursor.is_none() {
336                    state.phase = 3;
337                }
338            }
339        } else {
340            state.namespace += 1;
341            state.phase = 0;
342            state.generation = None;
343            state.registry_cursor = None;
344            state.shard_cursor = None;
345        }
346    }
347    let traversed = state.traversed();
348    let complete = traversed && state.exhaustive();
349    Ok(DiscoveryStep {
350        state,
351        batch,
352        repository,
353        complete,
354        traversed,
355    })
356}
357
358#[cfg(all(test, feature = "memory"))]
359#[path = "discovery_tests.rs"]
360mod tests;