1use crate::pipeline::{D34Shards, ShardMap, SinglePartition};
4use crate::repo::RepoId;
5use crate::rt::BoxFuture;
6use crate::store::outbox::OutboxBuilder;
7use crate::store::publication::{self, Advance, Clearance, Publication, Witness};
8use crate::store::{
9 Batch, BlobKey, Key, NamespaceStore, Partition, Precondition, StoreError, Value, keys,
10};
11use crate::timers::{DueTimer, Fired, TimerCtx, TimerHandler, TimerKind, registry::kinds};
12use std::collections::BTreeMap;
13
14pub const MAX_RECHECK_CALLS: u32 = 128;
18const REMOTE_PAGE_KEYS: usize = 8;
19const LOCAL_PAGE_KEYS: usize = 256;
20
21pub struct PublicationRecheck<T> {
23 pub target: T,
25}
26
27#[derive(Debug)]
29pub struct BudgetedRecheck<T> {
30 recheck: PublicationRecheck<T>,
31 budget: crate::purge::SliceBudget,
32}
33
34impl<T> PublicationRecheck<T> {
35 pub fn new(target: T) -> Self {
37 Self { target }
38 }
39
40 #[must_use]
43 pub fn with_alarm_budget(self, budget: crate::purge::SliceBudget) -> BudgetedRecheck<T> {
44 BudgetedRecheck {
45 recheck: self,
46 budget,
47 }
48 }
49}
50
51#[derive(Debug, Default, PartialEq, Eq)]
54struct Progress {
55 binding: mkit_core::hash::Hash,
56 position: u32,
57}
58impl Progress {
59 fn decode(value: &Value) -> Result<Self, StoreError> {
60 let bytes = value.as_bytes();
61 if bytes.len() != 37 || bytes[0] != 1 {
62 return Err(StoreError::Corrupt(
63 "invalid publication recheck cursor".into(),
64 ));
65 }
66 let binding = bytes[1..33]
67 .try_into()
68 .map_err(|_| StoreError::Corrupt("invalid recheck binding".into()))?;
69 let position = u32::from_le_bytes(
70 bytes[33..]
71 .try_into()
72 .map_err(|_| StoreError::Corrupt("invalid recheck position".into()))?,
73 );
74 if position > 8192 {
75 return Err(StoreError::Corrupt(
76 "publication recheck position exceeds bound".into(),
77 ));
78 }
79 Ok(Self { binding, position })
80 }
81 fn encode(&self) -> Value {
82 let mut bytes = Vec::with_capacity(37);
83 bytes.push(1);
84 bytes.extend_from_slice(&self.binding);
85 bytes.extend_from_slice(&self.position.to_le_bytes());
86 Value::new(bytes)
87 }
88 fn bind(&mut self, advance: &Value, state: &Publication) {
89 let mut digest = mkit_core::hash::Hasher::new();
90 digest.update(&state.generation.to_le_bytes());
91 digest.update(&state.boundary.to_le_bytes());
92 digest.update(advance.as_bytes());
93 let binding = digest.finalize();
94 if self.binding != binding {
95 self.binding = binding;
96 self.position = 0;
97 }
98 }
99}
100
101pub(crate) fn initial_value() -> Value {
103 Progress::default().encode()
104}
105
106fn dependency_groups(
107 source: &Partition,
108 shards: &dyn ShardMap,
109 repo: &RepoId,
110 advance: &Advance,
111) -> BTreeMap<Partition, Vec<Key>> {
112 let mut needed = advance
113 .dependencies
114 .iter()
115 .filter(|id| !advance.additions.contains(id))
116 .chain(advance.external_bases.iter())
117 .copied()
118 .collect::<Vec<_>>();
119 needed.sort_unstable();
120 needed.dedup();
121 let mut groups: BTreeMap<Partition, Vec<Key>> = BTreeMap::new();
122 for pack in needed {
123 let p = shards.membership(repo, &BlobKey::pack(pack));
124 let key = if p == *source {
125 keys::membership(&repo.name, &pack)
126 } else {
127 keys::published_member(&repo.name, &pack)
128 };
129 groups.entry(p).or_default().push(key);
130 }
131 groups
132}
133
134impl<T> core::fmt::Debug for PublicationRecheck<T> {
135 fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
136 f.debug_struct("PublicationRecheck").finish_non_exhaustive()
137 }
138}
139
140pub async fn dependencies<S: NamespaceStore, T: NamespaceStore>(
144 local: &S,
145 target: &T,
146 source: &Partition,
147 shards: &dyn ShardMap,
148 repo: &RepoId,
149 advance: &Advance,
150) -> Result<bool, StoreError> {
151 let groups = dependency_groups(source, shards, repo, advance);
152 for (p, keys) in groups {
153 for page in keys.chunks(8) {
154 let rows = if p == *source {
155 local.get_many(&p, page).await?
156 } else {
157 target.get_many(&p, page).await?
158 };
159 if rows.len() != page.len() {
160 return Err(StoreError::Corrupt(
161 "short publication dependency read".into(),
162 ));
163 }
164 for raw in rows {
165 let Some(raw) = raw else { return Ok(false) };
166 if !Witness::decode(&raw)?.visible(false, advance.generation) {
167 return Ok(false);
168 }
169 }
170 }
171 }
172 Ok(true)
173}
174
175async fn visible_rows<S: NamespaceStore>(
176 store: &S,
177 partition: &Partition,
178 keys: &[Key],
179) -> Result<Vec<Option<Value>>, StoreError> {
180 let rows = store.get_many(partition, keys).await?;
181 if rows.len() != keys.len() {
182 return Err(StoreError::Corrupt(
183 "short publication dependency read".into(),
184 ));
185 }
186 for raw in rows.iter().flatten() {
189 Witness::decode(raw)?;
190 }
191 Ok(rows)
192}
193
194#[allow(clippy::too_many_arguments)]
195async fn resume_dependencies<S: NamespaceStore, T: NamespaceStore>(
196 local: &S,
197 target: &T,
198 source: &Partition,
199 shards: &dyn ShardMap,
200 repo: &RepoId,
201 advance: &Advance,
202 progress: &mut Progress,
203 alarm_budget: Option<&crate::purge::SliceBudget>,
204) -> Result<bool, StoreError> {
205 let groups = dependency_groups(source, shards, repo, advance);
206 let needed = groups
207 .iter()
208 .filter(|(p, _)| *p != source)
209 .map(|(_, keys)| keys.len())
210 .sum::<usize>();
211 let position = usize::try_from(progress.position)
212 .map_err(|_| StoreError::Corrupt("invalid recheck position".into()))?;
213 if position > needed {
214 return Err(StoreError::Corrupt(
215 "publication recheck position exceeds dependencies".into(),
216 ));
217 }
218 if let Some(keys) = groups.get(source) {
221 for page in keys.chunks(LOCAL_PAGE_KEYS) {
222 let rows = visible_rows(local, source, page).await?;
223 for raw in rows {
224 if raw
225 .as_ref()
226 .map(Witness::decode)
227 .transpose()?
228 .is_none_or(|w| !w.visible(false, advance.generation))
229 {
230 return Ok(false);
231 }
232 }
233 }
234 }
235 let mut offset = 0;
236 let mut calls = 0;
237 for (partition, keys) in groups.iter().filter(|(p, _)| *p != source) {
238 let skip = position.saturating_sub(offset).min(keys.len());
239 offset += keys.len();
240 for page in keys[skip..].chunks(REMOTE_PAGE_KEYS) {
241 if calls == MAX_RECHECK_CALLS || alarm_budget.is_some_and(|budget| !budget.charge(1)) {
242 return Ok(false);
243 }
244 calls += 1;
245 let rows = visible_rows(target, partition, page).await?;
246 for raw in rows {
247 if raw
248 .as_ref()
249 .map(Witness::decode)
250 .transpose()?
251 .is_none_or(|w| !w.visible(false, advance.generation))
252 {
253 return Ok(false);
254 }
255 progress.position += 1;
256 }
257 }
258 }
259 Ok(true)
260}
261
262fn location(
263 partition: &Partition,
264 key: &Key,
265) -> Result<(RepoId, String, u64, &'static dyn ShardMap), StoreError> {
266 let Some(keys::ParsedKey::Advance {
267 repo,
268 name,
269 sequence,
270 }) = keys::parse(key)
271 else {
272 return Err(StoreError::Corrupt(
273 "invalid publication timer reference".into(),
274 ));
275 };
276 let ns = match partition {
277 Partition::Namespace(ns) | Partition::Ref { ns, .. } => ns.clone(),
278 _ => {
279 return Err(StoreError::Corrupt(
280 "publication timer on wrong partition".into(),
281 ));
282 }
283 };
284 let repo = RepoId {
285 namespace: ns,
286 name: repo,
287 };
288 let shards: &dyn ShardMap = if matches!(partition, Partition::Namespace(_)) {
289 &SinglePartition
290 } else {
291 &D34Shards
292 };
293 if shards.ref_shard(&repo, &name) != *partition {
294 return Err(StoreError::Corrupt("misrouted publication timer".into()));
295 }
296 Ok((repo, name, sequence, shards))
297}
298
299impl<S: NamespaceStore, T: NamespaceStore> TimerHandler<S> for PublicationRecheck<T> {
300 fn kind(&self) -> TimerKind {
301 kinds::PUBLICATION_RECHECK
302 }
303 fn max_per_tick(&self) -> Option<u32> {
304 Some(1)
305 }
306 fn fire<'a>(
307 &'a self,
308 ctx: &'a TimerCtx<'a, S>,
309 timer: &'a DueTimer,
310 ) -> BoxFuture<'a, Result<Fired, StoreError>> {
311 self.fire_inner(ctx, timer, None)
312 }
313}
314
315impl<S: NamespaceStore, T: NamespaceStore> TimerHandler<S> for BudgetedRecheck<T> {
316 fn kind(&self) -> TimerKind {
317 kinds::PUBLICATION_RECHECK
318 }
319 fn max_per_tick(&self) -> Option<u32> {
320 Some(1)
321 }
322 fn fire<'a>(
323 &'a self,
324 ctx: &'a TimerCtx<'a, S>,
325 timer: &'a DueTimer,
326 ) -> BoxFuture<'a, Result<Fired, StoreError>> {
327 self.recheck.fire_inner(ctx, timer, Some(&self.budget))
328 }
329}
330
331impl<T: NamespaceStore> PublicationRecheck<T> {
332 #[allow(clippy::too_many_lines)] fn fire_inner<'a, S: NamespaceStore>(
334 &'a self,
335 ctx: &'a TimerCtx<'a, S>,
336 timer: &'a DueTimer,
337 alarm_budget: Option<&'a crate::purge::SliceBudget>,
338 ) -> BoxFuture<'a, Result<Fired, StoreError>> {
339 Box::pin(async move {
340 if matches!(
341 keys::parse(&Key::new(timer.reference.clone())),
342 Some(keys::ParsedKey::Verification { .. })
343 ) {
344 return crate::indexed::publication::resume::fire(
345 ctx,
346 &self.target,
347 timer,
348 alarm_budget,
349 )
350 .await;
351 }
352 let mut progress = Progress::decode(&timer.value)?;
353 let key = Key::new(timer.reference.clone());
354 let (repo, name, sequence, shards) = location(ctx.partition, &key)?;
355 let wanted = [
356 key.clone(),
357 keys::publication(&repo.name, &name),
358 keys::outbox_sequence(),
359 keys::outcome_backlog(),
360 ];
361 let rows = ctx.store.get_many(ctx.partition, &wanted).await?;
362 if rows.len() != wanted.len() {
363 return Err(StoreError::Corrupt("short publication recheck read".into()));
364 }
365 let raw = rows[0]
366 .as_ref()
367 .ok_or_else(|| StoreError::Corrupt("missing retained publication work".into()))?;
368 let mut changed = Advance::decode(raw)?;
369 let state_raw = rows[1]
370 .as_ref()
371 .ok_or_else(|| StoreError::Corrupt("missing publication state".into()))?;
372 let state = Publication::decode(Some(state_raw))?;
373 if changed.sequence != sequence {
374 return Err(StoreError::Corrupt(
375 "publication timer sequence mismatch".into(),
376 ));
377 }
378 progress.bind(raw, &state);
379 let mut batch =
380 Batch::new().require(Precondition::NotAfter(ctx.now_ms.saturating_add(10_000)));
381 let eligible = changed.generation == state.generation
382 && (changed.state == Clearance::Pending || changed.state.publishable())
383 && changed.obligations.iter().all(|o| o.state.publishable());
384 let complete = if eligible {
385 resume_dependencies(
386 ctx.store,
387 &self.target,
388 ctx.partition,
389 shards,
390 &repo,
391 &changed,
392 &mut progress,
393 alarm_budget,
394 )
395 .await?
396 } else {
397 progress.position = 0;
400 false
401 };
402 if !complete || changed.state.publishable() {
403 batch.preconditions.extend([
404 Precondition::Equals(key, raw.clone()),
405 Precondition::Equals(wanted[1].clone(), state_raw.clone()),
406 ]);
407 if !complete {
408 return Ok(Fired::Reschedule {
409 due_at_ms: ctx.now_ms.saturating_add(publication::RECHECK_MS),
410 value: progress.encode(),
411 batch,
412 });
413 }
414 return Ok(Fired::Done(batch));
415 }
416 changed.state = Clearance::Cleared;
417 let eligible = publication::prefix(
418 ctx.store,
419 ctx.partition,
420 &repo.name,
421 &name,
422 &state,
423 &changed,
424 )
425 .await?;
426 let mut outbox = OutboxBuilder::new(rows[2].as_ref(), rows[3].as_ref())?;
427 publication::clear(
428 &repo,
429 &name,
430 ctx.partition,
431 shards,
432 state_raw,
433 raw,
434 &changed,
435 eligible,
436 &mut batch.preconditions,
437 &mut batch.writes,
438 &mut outbox,
439 )?;
440 outbox.relay_at(ctx.now_ms);
441 outbox.try_finish(&mut batch.preconditions, &mut batch.writes)?;
442 Ok(Fired::Done(batch))
443 })
444 }
445}
446
447#[cfg(test)]
448mod tests;