Skip to main content

mkit_server/store/
publication.rs

1//! `RefShard` authority for paired publication and retained inspection obligations.
2//!
3//! Sequence, value, membership and relay work join the caller's guarded apply.
4//! Ref/`RepoIndex` projections are witnesses, never evidence of live membership.
5use std::collections::BTreeSet;
6
7use mkit_core::hash::Hash;
8use serde::{Deserialize, Serialize};
9
10use super::outbox::{OutboxBuilder, guard};
11use super::{
12    BlobKey, Key, NamespaceStore, Partition, Precondition, StoreError, Value, Write, codec, keys,
13};
14use crate::pipeline::ShardMap;
15use crate::repo::{RepoId, RepoName};
16
17/// Bound per-ref outstanding values and the one-round prefix read.
18pub const MAX_UNPUBLISHED_ADVANCES: u64 = 64;
19/// Maximum durable dependency/obligation entries per advance.
20pub const MAX_ADVANCE_ITEMS: usize = 4096;
21/// Durable blocked advances retry without client traffic every five seconds.
22pub const RECHECK_MS: u64 = 5_000;
23
24/// The branch head and packmap, or one non-branch target.
25#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
26#[serde(deny_unknown_fields)]
27pub struct Pair {
28    /// Head, or the other ref's target.
29    pub head: Option<Hash>,
30    /// Packmap of a branch; absent for other refs.
31    pub packmap: Option<Hash>,
32}
33
34/// The clearance states from SPEC-SERVER ยง10.2.
35#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
36#[serde(rename_all = "snake_case")]
37pub enum Clearance {
38    /// Inspection or membership remains outstanding.
39    Pending,
40    /// All obligations and dependencies permit publication.
41    Cleared,
42    /// Quarantine requires release or re-inspection.
43    Held,
44    /// Rejection requires takedown completion.
45    Hit,
46    /// Takedown and remaining obligations are complete.
47    Resolved,
48}
49impl Clearance {
50    /// Whether this state permits membership and prefix publication.
51    #[must_use]
52    pub fn publishable(self) -> bool {
53        matches!(self, Self::Cleared | Self::Resolved)
54    }
55}
56
57/// Stable inspection identity and its independently retained result.
58#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
59#[serde(deny_unknown_fields)]
60pub struct Obligation {
61    /// Stable identity; superseding an inspection uses a different id.
62    pub id: Hash,
63    /// The inspection result, independently of membership dependencies.
64    pub state: Clearance,
65}
66
67/// One retained advance; this is also the inspector's durable handoff.
68#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
69#[serde(deny_unknown_fields)]
70pub struct Advance {
71    /// Never-reused sequence number.
72    pub sequence: u64,
73    /// Repository membership incarnation, supplied by lifecycle enforcement.
74    pub generation: u64,
75    /// Resulting paired live values.
76    pub value: Pair,
77    /// Packs added by this advance, including MKPL nodes.
78    pub additions: Vec<Hash>,
79    /// Packmap/closure dependencies; own additions may satisfy these.
80    pub dependencies: Vec<Hash>,
81    /// External delta source packs; own additions never satisfy these.
82    pub external_bases: Vec<Hash>,
83    /// Individually retained inspector obligations.
84    pub obligations: Vec<Obligation>,
85    /// Aggregate clearance, including dependencies and flags.
86    pub state: Clearance,
87    /// Logical operation correlation; never a credential.
88    pub operation: Hash,
89}
90impl Advance {
91    fn validate(&self) -> Result<(), StoreError> {
92        if self.sequence == 0
93            || self.additions.len() > super::outbox::MAX_TICKETS_PER_ADVANCE
94            || self.dependencies.len() > MAX_ADVANCE_ITEMS
95            || self.external_bases.len() > MAX_ADVANCE_ITEMS
96            || self.obligations.len() > MAX_ADVANCE_ITEMS
97            || !unique(&self.additions)
98            || !unique(&self.dependencies)
99            || !unique(&self.external_bases)
100        {
101            return Err(StoreError::Invalid("invalid publication advance".into()));
102        }
103        let ids: BTreeSet<_> = self.obligations.iter().map(|o| o.id).collect();
104        if ids.len() != self.obligations.len()
105            || self.state.publishable() && self.obligations.iter().any(|o| !o.state.publishable())
106        {
107            return Err(StoreError::Invalid(
108                "invalid publication obligations".into(),
109            ));
110        }
111        Ok(())
112    }
113    /// Encode a bounded v1 retained value.
114    pub fn encode(&self) -> Result<Value, StoreError> {
115        self.validate()?;
116        encode(self)
117    }
118    /// Decode and validate the entire record before using any field.
119    pub fn decode(value: &Value) -> Result<Self, StoreError> {
120        let row: Self = decode(value)?;
121        row.validate()
122            .map_err(|e| StoreError::Corrupt(e.to_string().into()))?;
123        Ok(row)
124    }
125}
126
127/// Persistent sequence, contiguous pointer and deletion boundary.
128#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
129#[serde(deny_unknown_fields)]
130pub struct Publication {
131    /// Last successfully applied advance; never deleted or reset.
132    pub sequence: u64,
133    /// Greatest cleared-or-resolved prefix after the deletion boundary.
134    pub published: u64,
135    /// Older verdicts cannot change ref values at or below this boundary.
136    pub boundary: u64,
137    /// Membership incarnation; lifecycle invalidation increments this.
138    pub generation: u64,
139    /// Paired published values, absent at deletion.
140    pub value: Pair,
141}
142impl Publication {
143    /// Decode persistent state; absence is the only initial state.
144    pub fn decode(value: Option<&Value>) -> Result<Self, StoreError> {
145        let row: Self = value.map(decode).transpose()?.unwrap_or_default();
146        if row.boundary > row.published || row.published > row.sequence {
147            return Err(StoreError::Corrupt("invalid publication prefix".into()));
148        }
149        Ok(row)
150    }
151    /// Encode persistent state.
152    pub fn encode(&self) -> Result<Value, StoreError> {
153        let value = encode(self)?;
154        Self::decode(Some(&value))?;
155        Ok(value)
156    }
157}
158
159/// A versioned local membership or published `RepoIndex` witness.
160#[derive(Debug, Clone, Copy, PartialEq, Eq)]
161pub struct Witness {
162    /// Membership generation, invalidated by repository deletion.
163    pub generation: u64,
164    /// Advance that established clearance; zero is immediate unticketed membership.
165    pub sequence: u64,
166    /// Independently cleared membership, even when the ref prefix is behind.
167    pub published: bool,
168    /// Serving stop; applies to writers too.
169    pub held: bool,
170}
171impl Witness {
172    /// Fixed-width versioned encoding; no unchecked optional fields.
173    #[must_use]
174    pub fn encode(self) -> Value {
175        let mut bytes = vec![1, u8::from(self.published), u8::from(self.held)];
176        bytes.extend_from_slice(&self.generation.to_be_bytes());
177        bytes.extend_from_slice(&self.sequence.to_be_bytes());
178        Value::new(bytes)
179    }
180    /// Decode a witness. Empty membership is the explicitly immediate upload form.
181    pub fn decode(value: &Value) -> Result<Self, StoreError> {
182        let bytes = value.as_bytes();
183        if bytes.is_empty() {
184            return Ok(Self {
185                generation: 0,
186                sequence: 0,
187                published: true,
188                held: false,
189            });
190        }
191        if bytes.len() != 19 || bytes[0] != 1 || bytes[1] > 1 || bytes[2] > 1 {
192            return Err(StoreError::Corrupt("invalid clearance witness".into()));
193        }
194        let number = |offset| -> Result<u64, StoreError> {
195            bytes
196                .get(offset..offset + 8)
197                .and_then(|b| b.try_into().ok())
198                .map(u64::from_be_bytes)
199                .ok_or_else(|| StoreError::Corrupt("short clearance witness".into()))
200        };
201        Ok(Self {
202            generation: number(3)?,
203            sequence: number(11)?,
204            published: bytes[1] != 0,
205            held: bytes[2] != 0,
206        })
207    }
208    /// Whether the caller may use this membership in its chosen view.
209    #[must_use]
210    pub fn visible(self, writer: bool, generation: u64) -> bool {
211        !self.held && self.generation == generation && (writer || self.published)
212    }
213}
214
215fn unique(ids: &[Hash]) -> bool {
216    ids.iter().collect::<BTreeSet<_>>().len() == ids.len()
217}
218fn encode<T: Serialize>(row: &T) -> Result<Value, StoreError> {
219    let mut bytes = vec![1];
220    bytes.extend(
221        serde_json::to_vec(row).map_err(|_| StoreError::Invalid("publication encoding".into()))?,
222    );
223    if bytes.len() > super::MAX_VALUE_BYTES {
224        return Err(StoreError::Invalid(
225            "publication value exceeds limit".into(),
226        ));
227    }
228    Ok(Value::new(bytes))
229}
230fn decode<T: serde::de::DeserializeOwned>(value: &Value) -> Result<T, StoreError> {
231    let bytes = value.as_bytes();
232    if bytes.first() != Some(&1) || bytes.len() > super::MAX_VALUE_BYTES {
233        return Err(StoreError::Corrupt(
234            "invalid publication version or size".into(),
235        ));
236    }
237    serde_json::from_slice(&bytes[1..])
238        .map_err(|_| StoreError::Corrupt("invalid publication record".into()))
239}
240
241/// Canonical shared sequence name for a branch pair, unchanged for another ref.
242#[must_use]
243pub fn sequence_ref(name: &str) -> String {
244    mkit_attest::grant::packmap_head(name).unwrap_or_else(|| name.to_owned())
245}
246
247/// Resulting branch pair or standalone ref names in stable order.
248#[must_use]
249pub fn value_refs(name: &str, value: &Pair) -> Vec<(String, Option<Hash>)> {
250    let mut refs = vec![(name.to_owned(), value.head)];
251    if let Some(packmap) = mkit_attest::grant::head_packmap(name) {
252        refs.push((packmap, value.packmap));
253    }
254    refs
255}
256
257/// Append one successful ref write. The caller commits this fragment with live refs.
258/// A deletion is an immediate boundary; older retained obligations remain intact.
259#[allow(clippy::too_many_arguments)]
260pub fn append(
261    repo: &RepoId,
262    name: &str,
263    source: &Partition,
264    shards: &dyn ShardMap,
265    prior: Option<&Value>,
266    mut advance: Advance,
267    deleted: bool,
268    pre: &mut Vec<Precondition>,
269    writes: &mut Vec<Write>,
270    outbox: &mut OutboxBuilder,
271) -> Result<Publication, StoreError> {
272    let name = sequence_ref(name);
273    let mut state = Publication::decode(prior)?;
274    if !deleted && state.sequence - state.published >= MAX_UNPUBLISHED_ADVANCES {
275        return Err(StoreError::unavailable("publication backlog full"));
276    }
277    state.sequence = state
278        .sequence
279        .checked_add(1)
280        .ok_or_else(|| StoreError::Corrupt("advance sequence overflow".into()))?;
281    advance.sequence = state.sequence;
282    advance.generation = state.generation;
283    advance.validate()?;
284    if !advance.state.publishable() {
285        // One timer per advance, installed in the same transaction as its retained
286        // obligations. Periodic rechecks cover cross-ref publication and delayed
287        // projections without an unbounded reverse dependency fanout.
288        writes.push(Write::Put(
289            keys::timer(
290                0,
291                crate::timers::registry::kinds::PUBLICATION_RECHECK.get(),
292                keys::advance(&repo.name, &name, advance.sequence).as_bytes(),
293            ),
294            crate::timers::publication_recheck::initial_value(),
295        ));
296    }
297    let key = keys::publication(&repo.name, &name);
298    pre.push(guard(key.clone(), prior));
299    if deleted {
300        state.boundary = state.sequence;
301        state.published = state.sequence;
302        // A partial deletion removes its component immediately, but must not
303        // publish the surviving live component from an uncleared advance.
304        state.value = Pair {
305            head: advance.value.head.and(state.value.head),
306            packmap: advance.value.packmap.and(state.value.packmap),
307        };
308        project_refs(repo, &name, source, shards, &state.value, writes, outbox);
309    } else if advance.state.publishable() && state.published + 1 == state.sequence {
310        state.published = state.sequence;
311        state.value = advance.value.clone();
312        project_refs(repo, &name, source, shards, &state.value, writes, outbox);
313    }
314    // Completed obligation-free values need no retained work once the prefix
315    // includes them. Keep blocked intermediates and every inspection obligation.
316    if advance.sequence > state.published || !advance.obligations.is_empty() {
317        writes.push(Write::Put(
318            keys::advance(&repo.name, &name, advance.sequence),
319            advance.encode()?,
320        ));
321    }
322    project_members(repo, source, shards, &advance, writes, outbox);
323    writes.push(Write::Put(key, state.encode()?));
324    Ok(state)
325}
326
327fn project_refs(
328    repo: &RepoId,
329    name: &str,
330    source: &Partition,
331    shards: &dyn ShardMap,
332    value: &Pair,
333    writes: &mut Vec<Write>,
334    outbox: &mut OutboxBuilder,
335) {
336    for (name, id) in value_refs(name, value) {
337        let key = keys::published_ref(&repo.name, &name);
338        writes.push(id.map_or_else(
339            || Write::Delete(key.clone()),
340            |id| Write::Put(key.clone(), codec::encode_ref_id(&id)),
341        ));
342        let target = shards.ref_index(repo, &name);
343        if target != *source {
344            let key = keys::published_index(&repo.name, &name);
345            match id {
346                Some(id) => outbox.relay(&target, vec![(key, codec::encode_ref_id(&id))]),
347                None => outbox.relay_delete(&target, vec![key]),
348            }
349        }
350    }
351}
352fn project_members(
353    repo: &RepoId,
354    source: &Partition,
355    shards: &dyn ShardMap,
356    advance: &Advance,
357    writes: &mut Vec<Write>,
358    outbox: &mut OutboxBuilder,
359) {
360    for pack in &advance.additions {
361        let witness = Witness {
362            generation: advance.generation,
363            sequence: advance.sequence,
364            published: advance.state.publishable(),
365            held: matches!(advance.state, Clearance::Held | Clearance::Hit),
366        };
367        let key = keys::membership(&repo.name, pack);
368        // Replace the live membership fragment rather than spending another per-ticket put.
369        writes.retain(|w| !matches!(w, Write::Put(k, _) | Write::Delete(k) if k == &key));
370        writes.push(Write::Put(key.clone(), witness.encode()));
371        let target = shards.membership(repo, &BlobKey::pack(*pack));
372        if target != *source {
373            outbox.relay(&target, vec![(key, witness.encode())]);
374            if witness.published {
375                outbox.relay(
376                    &target,
377                    vec![(keys::published_member(&repo.name, pack), witness.encode())],
378                );
379            }
380        }
381    }
382}
383
384/// Read authoritative sequence state without consulting a projection.
385pub async fn read<S: NamespaceStore>(
386    store: &S,
387    source: &Partition,
388    repo: &RepoName,
389    name: &str,
390) -> Result<Publication, StoreError> {
391    Publication::decode(
392        store
393            .get(source, &keys::publication(repo, &sequence_ref(name)))
394            .await?
395            .as_ref(),
396    )
397}
398
399/// Bounded read of every value that can advance the pointer. No projection is used.
400pub async fn prefix<S: NamespaceStore>(
401    store: &S,
402    source: &Partition,
403    repo: &RepoName,
404    name: &str,
405    state: &Publication,
406    changed: &Advance,
407) -> Result<(u64, Pair), StoreError> {
408    if state.sequence - state.published > MAX_UNPUBLISHED_ADVANCES {
409        return Err(StoreError::Corrupt(
410            "publication prefix exceeds bound".into(),
411        ));
412    }
413    if state.published == state.sequence {
414        return Ok((state.published, state.value.clone()));
415    }
416    let name = sequence_ref(name);
417    let wanted: Vec<Key> = (state.published + 1..=state.sequence)
418        .map(|sequence| keys::advance(repo, &name, sequence))
419        .collect();
420    let rows = store.get_many(source, &wanted).await?;
421    if rows.len() != wanted.len() {
422        return Err(StoreError::Corrupt("short publication prefix read".into()));
423    }
424    let mut result = (state.published, state.value.clone());
425    for (sequence, raw) in (state.published + 1..=state.sequence).zip(rows) {
426        let current = if changed.sequence == sequence {
427            changed.clone()
428        } else {
429            Advance::decode(
430                raw.as_ref()
431                    .ok_or_else(|| StoreError::Corrupt("missing retained advance".into()))?,
432            )?
433        };
434        if current.sequence != sequence || current.generation != state.generation {
435            return Err(StoreError::Corrupt(
436                "retained advance binding mismatch".into(),
437            ));
438        }
439        if !current.state.publishable() {
440            break;
441        }
442        result = (sequence, current.value);
443    }
444    Ok(result)
445}
446
447/// Guarded clearance fragment. A caller has verified obligations, flags and membership
448/// dependencies against current witnesses before calling this function. Resolved
449/// advances must already name their completed takedown replacements.
450#[allow(clippy::too_many_arguments)]
451pub fn clear(
452    repo: &RepoId,
453    name: &str,
454    source: &Partition,
455    shards: &dyn ShardMap,
456    state_raw: &Value,
457    advance_raw: &Value,
458    changed: &Advance,
459    eligible: (u64, Pair),
460    pre: &mut Vec<Precondition>,
461    writes: &mut Vec<Write>,
462    outbox: &mut OutboxBuilder,
463) -> Result<(), StoreError> {
464    let name = sequence_ref(name);
465    let mut state = Publication::decode(Some(state_raw))?;
466    let old = Advance::decode(advance_raw)?;
467    changed.validate()?;
468    if changed.sequence != old.sequence
469        || changed.sequence > state.sequence
470        || changed.generation != old.generation
471        || changed.generation != state.generation
472        || eligible.0 < state.published
473        || eligible.0 > state.sequence
474        || old.state == Clearance::Hit && changed.state == Clearance::Cleared
475    {
476        return Err(StoreError::Invalid(
477            "invalid publication clearance transition".into(),
478        ));
479    }
480    let key = keys::advance(&repo.name, &name, changed.sequence);
481    pre.push(guard(key.clone(), Some(advance_raw)));
482    writes.push(Write::Put(key, changed.encode()?));
483    pre.push(guard(keys::publication(&repo.name, &name), Some(state_raw)));
484    project_members(repo, source, shards, changed, writes, outbox);
485    if eligible.0 > state.published {
486        state.published = eligible.0;
487        state.value = eligible.1;
488        project_refs(repo, &name, source, shards, &state.value, writes, outbox);
489    }
490    writes.push(Write::Put(
491        keys::publication(&repo.name, &name),
492        state.encode()?,
493    ));
494    Ok(())
495}
496
497#[cfg(test)]
498#[path = "publication_tests.rs"]
499mod tests;