Skip to main content

mkit_server/timers/
publication_recheck.rs

1//! Durable dependency rechecks. One retained timer per blocked advance avoids
2//! an unbounded reverse-dependency fanout when membership publishes across refs.
3use 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
14/// Maximum routed witness reads per fire. Together with relay (512),
15/// verification (256), outcomes (64), rollup (32) and backup (1), this uses
16/// at most 993 of the existing 1,000-operation Paid alarm allowance.
17pub const MAX_RECHECK_CALLS: u32 = 128;
18const REMOTE_PAGE_KEYS: usize = 8;
19const LOCAL_PAGE_KEYS: usize = 256;
20
21/// One retained timer resumes a canonical dependency walk across alarms.
22pub struct PublicationRecheck<T> {
23    /// Routed metadata client for dependency witnesses outside the source shard.
24    pub target: T,
25}
26
27/// A recheck sharing its caller's whole-alarm allowance.
28#[derive(Debug)]
29pub struct BudgetedRecheck<T> {
30    recheck: PublicationRecheck<T>,
31    budget: crate::purge::SliceBudget,
32}
33
34impl<T> PublicationRecheck<T> {
35    /// Construct a recheck with the fixed routed-call cap.
36    pub fn new(target: T) -> Self {
37        Self { target }
38    }
39
40    /// Reserve shared operations before routed reads. The target must not
41    /// charge the same allowance again; source-local SQL calls are separate.
42    #[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// The current timer codec is fixed-width: version, bound row digest, then
52// the next routed witness ordinal in little endian. No pre-launch migration.
53#[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
101/// Initial value written atomically with a retained pending advance.
102pub(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
140/// Verify published dependencies against current versioned witnesses. Own
141/// additions satisfy closure, but never external delta-base dependencies.
142/// Missing/reordered projections reduce visibility. Corruption fails closed.
143pub 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    // Decode before advancing the cursor; missing/nonvisible rows are handled
187    // individually below, so no later row in a page skips an outstanding one.
188    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    // Unlike launch pm projections, source-local membership can be replaced
219    // by pending additions. Never cache its visibility across fires.
220    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)] // Dispatch and the existing guarded publication settlement.
333    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                // A hold/hit or outstanding obligation never becomes cleared by
398                // the timer. If it is later released, check every witness afresh.
399                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;