1use super::state::VerificationV1;
10use crate::repo::RepoName;
11use crate::store::{
12 codec::CODEC_V1,
13 index::{IndexEntry, IndexValue},
14 keys,
15};
16use crate::{Batch, NamespaceStore, Partition, StoreError, Value};
17use mkit_core::hash::Hash;
18use serde::{Deserialize, Serialize};
19
20pub const WINDOW_BYTES: u64 = 16 << 20;
22pub const DEFAULT_ENTRY_CAP: u32 = 4096;
24
25mod hex {
27 pub(super) mod bytes {
28 use serde::{Deserialize, Deserializer, Serialize, Serializer};
29 pub(in super::super) fn serialize<S: Serializer>(
30 bytes: &[u8],
31 s: S,
32 ) -> Result<S::Ok, S::Error> {
33 mkit_core::hash::to_hex_bytes(bytes).serialize(s)
34 }
35 pub(in super::super) fn deserialize<'de, D: Deserializer<'de>>(
36 d: D,
37 ) -> Result<Vec<u8>, D::Error> {
38 let text = String::deserialize(d)?;
39 if text.len() % 2 != 0 || !text.is_ascii() {
40 return Err(serde::de::Error::custom("malformed hex"));
41 }
42 (0..text.len())
43 .step_by(2)
44 .map(|i| u8::from_str_radix(&text[i..i + 2], 16).map_err(serde::de::Error::custom))
45 .collect()
46 }
47 }
48}
49
50#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
53#[serde(rename_all = "snake_case")]
54pub enum Phase {
55 #[default]
57 Decode,
58 ClosureResolve,
60 EmitIndex,
62 AwaitDelivery,
64 Extract,
66 Verify,
68 Recheck,
70 Watch,
72}
73
74#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
76#[serde(rename_all = "snake_case")]
77pub enum Kind {
78 #[default]
80 Unknown,
81 Pack,
83 Packlist,
85}
86
87#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
90#[serde(rename_all = "snake_case")]
91pub enum Outcome {
92 BaseMissing,
94 Blocked,
96 BaseCapped,
98 ClosureCapped,
100 ClosureMissing,
102 PacklistMissing,
104 OpenClosure,
106 ExternalTooDeep,
108 DecodeBudget,
110 ExtractionUnavailable,
112 ObjectBlocked,
114}
115
116#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
118#[serde(deny_unknown_fields)]
119pub struct ExtractionGroupMember {
120 pub pack: Hash,
122 pub ticket: Hash,
124 pub bytes: u64,
126 pub created_at_ms: u64,
128 pub already_verified: bool,
131}
132
133#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
135#[serde(deny_unknown_fields)]
136#[allow(clippy::struct_excessive_bools)] pub struct VerifyJobV1 {
138 pub generation: u64,
140 pub gone: bool,
142 pub member_body_id: Option<Hash>,
144 #[serde(skip)]
146 pub members_loaded: bool,
147 pub extraction_head: Option<Hash>,
149 #[serde(default)]
151 pub extraction: Option<ExtractionV1>,
152 #[serde(default)]
155 pub extraction_group: Vec<ExtractionGroupMember>,
156 pub ticket_id: Hash,
158 pub created_at_ms: u64,
160 pub pack_len: u64,
162 pub phase: Phase,
164 pub kind: Kind,
166 pub version: u32,
168 #[serde(with = "hex::bytes")]
170 pub cursor: Vec<u8>,
171 pub etag: Option<String>,
173 pub entries: u64,
175 pub in_pack_bytes: u64,
177 pub external_bytes: u64,
179 pub windows_done: u32,
181 pub attempts: u32,
183 pub entry_cap: u32,
185 pub closure_cap: u32,
187 pub restarts: u8,
189 pub bad_signature: bool,
191 pub extract_needed: bool,
193 #[serde(with = "hex::bytes")]
195 pub scan: Vec<u8>,
196 pub owed: u64,
198 pub final_pass: bool,
200 #[serde(skip)]
202 pub satisfying: Vec<Hash>,
203 pub last_relay_seq: Option<u64>,
205 pub closure_final_at_ms: Option<u64>,
207 #[serde(skip)]
209 pub packlist: Vec<Hash>,
210 pub packlist_prev: Option<Hash>,
212 pub outcome: Option<Outcome>,
214}
215
216#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
218#[serde(deny_unknown_fields)]
219pub struct ExtractionV1 {
220 pub sources: Vec<ExtractionSource>,
222 pub group: Hash,
224 pub stage: u8,
226 pub member: usize,
228 #[serde(with = "hex::bytes")]
230 pub scan: Vec<u8>,
231 pub staged_objects: u64,
233 pub staged_bytes: u64,
235 pub selected_bytes: u64,
237 pub object: Option<Hash>,
239 pub length: u64,
241 pub chunk: u32,
243 pub chunk_offset: u64,
245 pub written: u64,
247 pub cvs: u32,
249 pub root: Option<Hash>,
251 #[serde(with = "hex::bytes")]
253 pub session: Vec<u8>,
254 pub uploaded: u32,
256 pub relay: Option<u64>,
258 pub reconstruction: Option<MemberCursor>,
260}
261
262#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
264#[serde(deny_unknown_fields)]
265pub struct MemberCursor {
266 pub target: Hash,
268 pub next: Hash,
270 pub preferred: Option<(Hash, u64)>,
272 pub level: u32,
274 pub local: bool,
276 pub ascending: bool,
278 pub canonical: Option<(Hash, u64, u32, Hash)>,
280 pub bytes: u64,
282}
283
284#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
286#[serde(deny_unknown_fields)]
287pub struct ExtractionSource {
288 pub member: ExtractionGroupMember,
290 pub etag: Option<String>,
292 pub version: u32,
294 pub entries: u64,
296 pub decoded: u64,
298 pub member_body_id: Option<Hash>,
300}
301
302impl VerifyJobV1 {
303 pub(super) fn closure_retry(&self) -> bool {
304 matches!(
305 self.outcome,
306 Some(Outcome::ClosureMissing | Outcome::PacklistMissing | Outcome::BaseCapped)
307 ) && self
308 .extraction
309 .as_ref()
310 .is_some_and(|x| x.object.is_none() && (x.stage <= 2 || (10..=13).contains(&x.stage)))
311 }
312
313 #[must_use]
315 pub fn new(ticket_id: Hash, created_at_ms: u64, pack_len: u64, entry_cap: u32) -> Self {
316 Self {
317 ticket_id,
318 created_at_ms,
319 pack_len,
320 entry_cap,
321 closure_cap: 4,
322 members_loaded: true,
323 ..Self::default()
324 }
325 }
326
327 pub fn restart(&mut self) {
329 let fresh = Self::new(
330 self.ticket_id,
331 self.created_at_ms,
332 self.pack_len,
333 self.entry_cap,
334 );
335 *self = Self {
336 restarts: self.restarts.saturating_add(1),
337 extraction_group: self.extraction_group.clone(),
338 extraction_head: self.extraction_head,
339 ..fresh
340 };
341 }
342
343 #[must_use]
345 pub fn usable(&self) -> bool {
346 !self.gone && self.outcome.is_none() && matches!(self.phase, Phase::Recheck | Phase::Watch)
347 }
348}
349
350#[must_use]
355pub fn encode_job(job: &VerifyJobV1) -> Value {
356 let mut bytes = vec![CODEC_V1];
357 serde_json::to_writer(&mut bytes, job).expect("job DTO serializes");
358 Value::new(bytes)
359}
360
361pub fn decode_job(value: &Value) -> Result<VerifyJobV1, StoreError> {
363 let Some((&CODEC_V1, body)) = value.as_bytes().split_first() else {
364 return Err(StoreError::Corrupt("bad verification job version".into()));
365 };
366 let job: VerifyJobV1 = serde_json::from_slice(body)
367 .map_err(|_| StoreError::Corrupt("bad verification job".into()))?;
368 validate_header(&job, value)?;
369 Ok(job)
370}
371
372pub const MAX_JOB_HEADER_BYTES: usize = 16 << 10;
375
376fn validate_header(job: &VerifyJobV1, raw: &Value) -> Result<(), StoreError> {
377 let bounded = raw.as_bytes().len() <= MAX_JOB_HEADER_BYTES
378 && job.cursor.len() <= 4096
379 && job.scan.len() <= 324
380 && job.etag.as_ref().is_none_or(|e| e.len() <= 64)
381 && job.extraction_group.len() <= crate::store::outbox::MAX_TICKETS_PER_ADVANCE
382 && job.extraction.as_ref().is_none_or(|x| {
383 job.cursor.is_empty()
384 && x.scan.len() <= 324
385 && x.session.len() <= 1024
386 && x.cvs <= 10_000
387 && x.uploaded <= 10_000
388 && x.reconstruction
389 .as_ref()
390 .is_none_or(|r| r.canonical.is_none_or(|(_, n, _, _)| n <= 8 << 20))
391 && x.stage <= 13
392 && x.member <= x.sources.len()
393 && x.sources.len() <= crate::store::outbox::MAX_TICKETS_PER_ADVANCE
394 && x.sources
395 .iter()
396 .all(|s| s.etag.as_ref().is_none_or(|e| e.len() <= 64))
397 });
398 if bounded {
399 Ok(())
400 } else {
401 Err(StoreError::Corrupt("oversized job header".into()))
402 }
403}
404
405#[derive(Serialize, Deserialize)]
406#[serde(deny_unknown_fields)]
407struct MemberLists {
408 satisfying: Vec<Hash>,
409 packlist: Vec<Hash>,
410}
411
412fn member_id(bytes: &[u8]) -> Hash {
413 let mut h = mkit_core::hash::Hasher::new();
414 h.update(b"mkit-job-members:v1");
415 h.update(bytes);
416 h.finalize()
417}
418fn member_key(repo: &RepoName, pack: &Hash, id: &Hash) -> crate::Key {
419 keys::verify_row(repo, pack, keys::VC_CANDIDATE, Some(id))
420}
421
422pub fn write_job(
425 mut batch: Batch,
426 job: &mut VerifyJobV1,
427 prior: Option<&Value>,
428 repo: &RepoName,
429 pack: &Hash,
430) -> Result<Batch, StoreError> {
431 let old = prior.map(decode_job).transpose()?;
432 job.generation = old
433 .as_ref()
434 .map_or(0, |j| j.generation)
435 .checked_add(1)
436 .ok_or_else(|| StoreError::Corrupt("job generation overflow".into()))?;
437 let old_body = old.as_ref().and_then(|j| j.member_body_id);
438 if job.members_loaded {
439 if job.satisfying.len() > crate::store::index::MAX_LOOKUP_IDS
440 || job.packlist.len()
441 > crate::store::index::MAX_LOOKUP_IDS
442 + crate::store::outbox::MAX_TICKETS_PER_ADVANCE
443 {
444 return Err(StoreError::Corrupt("oversized job member lists".into()));
445 }
446 job.member_body_id = if job.satisfying.is_empty() && job.packlist.is_empty() {
447 None
448 } else {
449 let mut bytes = vec![CODEC_V1];
450 serde_json::to_writer(
451 &mut bytes,
452 &MemberLists {
453 satisfying: job.satisfying.clone(),
454 packlist: job.packlist.clone(),
455 },
456 )
457 .map_err(StoreError::unavailable)?;
458 let id = member_id(&bytes);
459 if Some(id) != old_body {
460 batch = batch.put(member_key(repo, pack, &id), Value::new(bytes));
461 }
462 Some(id)
463 };
464 } else if job.member_body_id != old_body {
465 return Err(StoreError::Corrupt(
466 "unloaded job member lists changed".into(),
467 ));
468 }
469 let header = encode_job(job);
470 validate_header(job, &header)?;
471 Ok(batch.put(keys::verify_job(repo, pack), header))
472}
473
474pub async fn hydrate_job<S: NamespaceStore>(
477 store: &S,
478 source: &Partition,
479 repo: &RepoName,
480 pack: &Hash,
481 job: &mut VerifyJobV1,
482) -> Result<(), StoreError> {
483 if !job.gone
484 && let Some(id) = job.member_body_id
485 {
486 let raw = store
487 .get(source, &member_key(repo, pack, &id))
488 .await?
489 .ok_or_else(|| StoreError::Unavailable("job member body disappeared".into()))?;
490 let Some((&CODEC_V1, bytes)) = raw.as_bytes().split_first() else {
491 return Err(StoreError::Corrupt("bad job member body version".into()));
492 };
493 if member_id(raw.as_bytes()) != id {
494 return Err(StoreError::Corrupt("bad job member body digest".into()));
495 }
496 let lists: MemberLists = serde_json::from_slice(bytes)
497 .map_err(|_| StoreError::Corrupt("bad job member body".into()))?;
498 if lists.satisfying.len() > crate::store::index::MAX_LOOKUP_IDS
499 || lists.packlist.len()
500 > crate::store::index::MAX_LOOKUP_IDS
501 + crate::store::outbox::MAX_TICKETS_PER_ADVANCE
502 {
503 return Err(StoreError::Corrupt("oversized job member body".into()));
504 }
505 job.satisfying = lists.satisfying;
506 job.packlist = lists.packlist;
507 }
508 job.members_loaded = true;
509 Ok(())
510}
511
512#[derive(Debug, Clone, Copy, PartialEq, Eq)]
514pub struct FrameRow {
515 pub value: IndexValue,
517 pub object_type: u8,
519 pub external: Option<Hash>,
521}
522
523pub fn encode_frame(id: &Hash, row: &FrameRow) -> Result<Value, StoreError> {
525 let index = crate::store::codec::encode_object_index(id, &row.value)?;
526 let mut bytes = vec![row.object_type, u8::from(row.external.is_some())];
527 if let Some(base) = &row.external {
528 bytes.extend_from_slice(base);
529 }
530 bytes.extend_from_slice(index.as_bytes());
531 Ok(Value::new(bytes))
532}
533
534pub fn decode_frame(id: &Hash, value: &Value) -> Result<FrameRow, StoreError> {
536 let corrupt = || StoreError::Corrupt("bad verification frame row".into());
537 let bytes = value.as_bytes();
538 let (&object_type, rest) = bytes.split_first().ok_or_else(corrupt)?;
539 let (external, rest) = match rest.split_first() {
540 Some((0, rest)) => (None, rest),
541 Some((1, rest)) => {
542 let (base, rest) = rest.split_first_chunk::<32>().ok_or_else(corrupt)?;
543 (Some(*base), rest)
544 }
545 _ => return Err(corrupt()),
546 };
547 Ok(FrameRow {
548 value: crate::store::codec::decode_object_index(id, &Value::new(rest.to_vec()))?,
549 object_type,
550 external,
551 })
552}
553
554#[derive(Debug, Clone, Copy, PartialEq, Eq)]
557pub struct BaseRow {
558 pub size: u64,
560 pub depth: u32,
562 pub entry: u64,
564}
565
566#[must_use]
568pub fn encode_base(row: &BaseRow) -> Value {
569 let mut bytes = row.size.to_be_bytes().to_vec();
570 bytes.extend_from_slice(&row.depth.to_be_bytes());
571 bytes.extend_from_slice(&row.entry.to_be_bytes());
572 Value::new(bytes)
573}
574
575pub fn decode_base(value: &Value) -> Result<BaseRow, StoreError> {
577 let bytes: &[u8; 20] = value
578 .as_bytes()
579 .try_into()
580 .map_err(|_| StoreError::Corrupt("bad verification base row".into()))?;
581 Ok(BaseRow {
582 size: u64::from_be_bytes(bytes[..8].try_into().unwrap_or_default()),
583 depth: u32::from_be_bytes(bytes[8..12].try_into().unwrap_or_default()),
584 entry: u64::from_be_bytes(bytes[12..].try_into().unwrap_or_default()),
585 })
586}
587
588#[must_use]
590pub fn index_entry(id: Hash, row: &FrameRow) -> IndexEntry {
591 IndexEntry {
592 object: id,
593 value: row.value,
594 }
595}
596
597#[must_use]
599pub fn timer_reference(repo: &RepoName, pack: &Hash) -> Vec<u8> {
600 let mut reference = repo.as_str().as_bytes().to_vec();
601 reference.push(0);
602 reference.extend_from_slice(pack);
603 reference
604}
605
606#[must_use]
608pub fn parse_reference(reference: &[u8]) -> Option<(RepoName, Hash)> {
609 let sep = reference.iter().position(|&b| b == 0)?;
610 let pack: Hash = reference[sep + 1..].try_into().ok()?;
611 let repo = RepoName::new(String::from_utf8(reference[..sep].to_vec()).ok()?).ok()?;
612 Some((repo, pack))
613}
614
615type Stored<T> = Option<(T, Value)>;
616
617pub async fn read_job<S: NamespaceStore>(
619 store: &S,
620 source: &Partition,
621 repo: &RepoName,
622 pack: &Hash,
623) -> Result<(Stored<VerifyJobV1>, Stored<VerificationV1>), StoreError> {
624 let rows = store
625 .get_many(
626 source,
627 &[keys::verify_job(repo, pack), keys::verification(repo, pack)],
628 )
629 .await?;
630 let [job, state] = <[_; 2]>::try_from(rows)
631 .map_err(|_| StoreError::Corrupt("short verification read".into()))?;
632 let job = if let Some(raw) = job {
633 let mut job = decode_job(&raw)?;
634 hydrate_job(store, source, repo, pack, &mut job).await?;
635 Some((job, raw))
636 } else {
637 None
638 };
639 Ok((
640 job,
641 state
642 .map(|raw| super::state::decode(&raw).map(|state| (state, raw)))
643 .transpose()?,
644 ))
645}
646
647#[cfg(test)]
648mod tests {
649 use super::*;
650
651 #[test]
652 fn job_frame_and_base_codecs_round_trip() {
653 let mut job = VerifyJobV1::new([1; 32], 5, 99, 4096);
654 job.cursor = vec![0xab, 0x01];
655 job.satisfying = vec![[2; 32]];
656 job.outcome = Some(Outcome::BaseMissing);
657 let header = decode_job(&encode_job(&job)).unwrap();
658 assert!(header.satisfying.is_empty());
659 assert!(!header.members_loaded);
660 assert_eq!(header.outcome, job.outcome);
661 assert!(decode_job(&Value::new(b"\x02{}".to_vec())).is_err());
662 let frame = FrameRow {
663 value: IndexValue {
664 frame_offset: 12,
665 frame_length: 40,
666 wire_type: 0x02,
667 decoded_size: 7,
668 chain_depth: 2,
669 delta_base: Some([3; 32]),
670 },
671 object_type: 3,
672 external: Some([4; 32]),
673 };
674 let id = [9; 32];
675 assert_eq!(
676 decode_frame(&id, &encode_frame(&id, &frame).unwrap()).unwrap(),
677 frame
678 );
679 let base = BaseRow {
680 size: 1,
681 depth: 2,
682 entry: 3,
683 };
684 assert_eq!(decode_base(&encode_base(&base)).unwrap(), base);
685 let name = RepoName::new("a").unwrap();
686 assert_eq!(
687 parse_reference(&timer_reference(&name, &[7; 32])),
688 Some((name, [7; 32]))
689 );
690 job.restart();
691 assert_eq!(
692 (job.restarts, job.phase, job.cursor.len()),
693 (1, Phase::Decode, 0)
694 );
695 }
696}