1use 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#[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 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 #[must_use]
83 pub fn exhaustive(&self) -> bool {
84 self.exhaustive
85 }
86 #[must_use]
88 pub fn traversed(&self) -> bool {
89 self.namespace == self.namespaces.len()
90 }
91 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#[derive(Debug)]
140pub struct DiscoveryStep {
141 pub state: DiscoveryState,
143 pub batch: Batch,
145 pub repository: Option<RepoId>,
147 pub complete: bool,
151 pub traversed: bool,
153}
154#[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;