1use 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
17pub const MAX_UNPUBLISHED_ADVANCES: u64 = 64;
19pub const MAX_ADVANCE_ITEMS: usize = 4096;
21pub const RECHECK_MS: u64 = 5_000;
23
24#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
26#[serde(deny_unknown_fields)]
27pub struct Pair {
28 pub head: Option<Hash>,
30 pub packmap: Option<Hash>,
32}
33
34#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
36#[serde(rename_all = "snake_case")]
37pub enum Clearance {
38 Pending,
40 Cleared,
42 Held,
44 Hit,
46 Resolved,
48}
49impl Clearance {
50 #[must_use]
52 pub fn publishable(self) -> bool {
53 matches!(self, Self::Cleared | Self::Resolved)
54 }
55}
56
57#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
59#[serde(deny_unknown_fields)]
60pub struct Obligation {
61 pub id: Hash,
63 pub state: Clearance,
65}
66
67#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
69#[serde(deny_unknown_fields)]
70pub struct Advance {
71 pub sequence: u64,
73 pub generation: u64,
75 pub value: Pair,
77 pub additions: Vec<Hash>,
79 pub dependencies: Vec<Hash>,
81 pub external_bases: Vec<Hash>,
83 pub obligations: Vec<Obligation>,
85 pub state: Clearance,
87 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 pub fn encode(&self) -> Result<Value, StoreError> {
115 self.validate()?;
116 encode(self)
117 }
118 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#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
129#[serde(deny_unknown_fields)]
130pub struct Publication {
131 pub sequence: u64,
133 pub published: u64,
135 pub boundary: u64,
137 pub generation: u64,
139 pub value: Pair,
141}
142impl Publication {
143 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 pub fn encode(&self) -> Result<Value, StoreError> {
153 let value = encode(self)?;
154 Self::decode(Some(&value))?;
155 Ok(value)
156 }
157}
158
159#[derive(Debug, Clone, Copy, PartialEq, Eq)]
161pub struct Witness {
162 pub generation: u64,
164 pub sequence: u64,
166 pub published: bool,
168 pub held: bool,
170}
171impl Witness {
172 #[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 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 #[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#[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#[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#[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 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 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 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 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
384pub 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
399pub 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#[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;