use super::state::VerificationV1;
use crate::repo::RepoName;
use crate::store::{
codec::CODEC_V1,
index::{IndexEntry, IndexValue},
keys,
};
use crate::{Batch, NamespaceStore, Partition, StoreError, Value};
use mkit_core::hash::Hash;
use serde::{Deserialize, Serialize};
pub const WINDOW_BYTES: u64 = 16 << 20;
pub const DEFAULT_ENTRY_CAP: u32 = 4096;
mod hex {
pub(super) mod bytes {
use serde::{Deserialize, Deserializer, Serialize, Serializer};
pub(in super::super) fn serialize<S: Serializer>(
bytes: &[u8],
s: S,
) -> Result<S::Ok, S::Error> {
mkit_core::hash::to_hex_bytes(bytes).serialize(s)
}
pub(in super::super) fn deserialize<'de, D: Deserializer<'de>>(
d: D,
) -> Result<Vec<u8>, D::Error> {
let text = String::deserialize(d)?;
if text.len() % 2 != 0 || !text.is_ascii() {
return Err(serde::de::Error::custom("malformed hex"));
}
(0..text.len())
.step_by(2)
.map(|i| u8::from_str_radix(&text[i..i + 2], 16).map_err(serde::de::Error::custom))
.collect()
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Phase {
#[default]
Decode,
ClosureResolve,
EmitIndex,
AwaitDelivery,
Extract,
Verify,
Recheck,
Watch,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Kind {
#[default]
Unknown,
Pack,
Packlist,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Outcome {
BaseMissing,
Blocked,
BaseCapped,
ClosureCapped,
ClosureMissing,
PacklistMissing,
OpenClosure,
ExternalTooDeep,
DecodeBudget,
ExtractionUnavailable,
ObjectBlocked,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct ExtractionGroupMember {
pub pack: Hash,
pub ticket: Hash,
pub bytes: u64,
pub created_at_ms: u64,
pub already_verified: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
#[allow(clippy::struct_excessive_bools)] pub struct VerifyJobV1 {
pub generation: u64,
pub gone: bool,
pub member_body_id: Option<Hash>,
#[serde(skip)]
pub members_loaded: bool,
pub extraction_head: Option<Hash>,
#[serde(default)]
pub extraction: Option<ExtractionV1>,
#[serde(default)]
pub extraction_group: Vec<ExtractionGroupMember>,
pub ticket_id: Hash,
pub created_at_ms: u64,
pub pack_len: u64,
pub phase: Phase,
pub kind: Kind,
pub version: u32,
#[serde(with = "hex::bytes")]
pub cursor: Vec<u8>,
pub etag: Option<String>,
pub entries: u64,
pub in_pack_bytes: u64,
pub external_bytes: u64,
pub windows_done: u32,
pub attempts: u32,
pub entry_cap: u32,
pub closure_cap: u32,
pub restarts: u8,
pub bad_signature: bool,
pub extract_needed: bool,
#[serde(with = "hex::bytes")]
pub scan: Vec<u8>,
pub owed: u64,
pub final_pass: bool,
#[serde(skip)]
pub satisfying: Vec<Hash>,
pub last_relay_seq: Option<u64>,
pub closure_final_at_ms: Option<u64>,
#[serde(skip)]
pub packlist: Vec<Hash>,
pub packlist_prev: Option<Hash>,
pub outcome: Option<Outcome>,
}
#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct ExtractionV1 {
pub sources: Vec<ExtractionSource>,
pub group: Hash,
pub stage: u8,
pub member: usize,
#[serde(with = "hex::bytes")]
pub scan: Vec<u8>,
pub staged_objects: u64,
pub staged_bytes: u64,
pub selected_bytes: u64,
pub object: Option<Hash>,
pub length: u64,
pub chunk: u32,
pub chunk_offset: u64,
pub written: u64,
pub cvs: u32,
pub root: Option<Hash>,
#[serde(with = "hex::bytes")]
pub session: Vec<u8>,
pub uploaded: u32,
pub relay: Option<u64>,
pub reconstruction: Option<MemberCursor>,
}
#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct MemberCursor {
pub target: Hash,
pub next: Hash,
pub preferred: Option<(Hash, u64)>,
pub level: u32,
pub local: bool,
pub ascending: bool,
pub canonical: Option<(Hash, u64, u32, Hash)>,
pub bytes: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct ExtractionSource {
pub member: ExtractionGroupMember,
pub etag: Option<String>,
pub version: u32,
pub entries: u64,
pub decoded: u64,
pub member_body_id: Option<Hash>,
}
impl VerifyJobV1 {
pub(super) fn closure_retry(&self) -> bool {
matches!(
self.outcome,
Some(Outcome::ClosureMissing | Outcome::PacklistMissing | Outcome::BaseCapped)
) && self
.extraction
.as_ref()
.is_some_and(|x| x.object.is_none() && (x.stage <= 2 || (10..=13).contains(&x.stage)))
}
#[must_use]
pub fn new(ticket_id: Hash, created_at_ms: u64, pack_len: u64, entry_cap: u32) -> Self {
Self {
ticket_id,
created_at_ms,
pack_len,
entry_cap,
closure_cap: 4,
members_loaded: true,
..Self::default()
}
}
pub fn restart(&mut self) {
let fresh = Self::new(
self.ticket_id,
self.created_at_ms,
self.pack_len,
self.entry_cap,
);
*self = Self {
restarts: self.restarts.saturating_add(1),
extraction_group: self.extraction_group.clone(),
extraction_head: self.extraction_head,
..fresh
};
}
#[must_use]
pub fn usable(&self) -> bool {
!self.gone && self.outcome.is_none() && matches!(self.phase, Phase::Recheck | Phase::Watch)
}
}
#[must_use]
pub fn encode_job(job: &VerifyJobV1) -> Value {
let mut bytes = vec![CODEC_V1];
serde_json::to_writer(&mut bytes, job).expect("job DTO serializes");
Value::new(bytes)
}
pub fn decode_job(value: &Value) -> Result<VerifyJobV1, StoreError> {
let Some((&CODEC_V1, body)) = value.as_bytes().split_first() else {
return Err(StoreError::Corrupt("bad verification job version".into()));
};
let job: VerifyJobV1 = serde_json::from_slice(body)
.map_err(|_| StoreError::Corrupt("bad verification job".into()))?;
validate_header(&job, value)?;
Ok(job)
}
pub const MAX_JOB_HEADER_BYTES: usize = 16 << 10;
fn validate_header(job: &VerifyJobV1, raw: &Value) -> Result<(), StoreError> {
let bounded = raw.as_bytes().len() <= MAX_JOB_HEADER_BYTES
&& job.cursor.len() <= 4096
&& job.scan.len() <= 324
&& job.etag.as_ref().is_none_or(|e| e.len() <= 64)
&& job.extraction_group.len() <= crate::store::outbox::MAX_TICKETS_PER_ADVANCE
&& job.extraction.as_ref().is_none_or(|x| {
job.cursor.is_empty()
&& x.scan.len() <= 324
&& x.session.len() <= 1024
&& x.cvs <= 10_000
&& x.uploaded <= 10_000
&& x.reconstruction
.as_ref()
.is_none_or(|r| r.canonical.is_none_or(|(_, n, _, _)| n <= 8 << 20))
&& x.stage <= 13
&& x.member <= x.sources.len()
&& x.sources.len() <= crate::store::outbox::MAX_TICKETS_PER_ADVANCE
&& x.sources
.iter()
.all(|s| s.etag.as_ref().is_none_or(|e| e.len() <= 64))
});
if bounded {
Ok(())
} else {
Err(StoreError::Corrupt("oversized job header".into()))
}
}
#[derive(Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct MemberLists {
satisfying: Vec<Hash>,
packlist: Vec<Hash>,
}
fn member_id(bytes: &[u8]) -> Hash {
let mut h = mkit_core::hash::Hasher::new();
h.update(b"mkit-job-members:v1");
h.update(bytes);
h.finalize()
}
fn member_key(repo: &RepoName, pack: &Hash, id: &Hash) -> crate::Key {
keys::verify_row(repo, pack, keys::VC_CANDIDATE, Some(id))
}
pub fn write_job(
mut batch: Batch,
job: &mut VerifyJobV1,
prior: Option<&Value>,
repo: &RepoName,
pack: &Hash,
) -> Result<Batch, StoreError> {
let old = prior.map(decode_job).transpose()?;
job.generation = old
.as_ref()
.map_or(0, |j| j.generation)
.checked_add(1)
.ok_or_else(|| StoreError::Corrupt("job generation overflow".into()))?;
let old_body = old.as_ref().and_then(|j| j.member_body_id);
if job.members_loaded {
if job.satisfying.len() > crate::store::index::MAX_LOOKUP_IDS
|| job.packlist.len()
> crate::store::index::MAX_LOOKUP_IDS
+ crate::store::outbox::MAX_TICKETS_PER_ADVANCE
{
return Err(StoreError::Corrupt("oversized job member lists".into()));
}
job.member_body_id = if job.satisfying.is_empty() && job.packlist.is_empty() {
None
} else {
let mut bytes = vec![CODEC_V1];
serde_json::to_writer(
&mut bytes,
&MemberLists {
satisfying: job.satisfying.clone(),
packlist: job.packlist.clone(),
},
)
.map_err(StoreError::unavailable)?;
let id = member_id(&bytes);
if Some(id) != old_body {
batch = batch.put(member_key(repo, pack, &id), Value::new(bytes));
}
Some(id)
};
} else if job.member_body_id != old_body {
return Err(StoreError::Corrupt(
"unloaded job member lists changed".into(),
));
}
let header = encode_job(job);
validate_header(job, &header)?;
Ok(batch.put(keys::verify_job(repo, pack), header))
}
pub async fn hydrate_job<S: NamespaceStore>(
store: &S,
source: &Partition,
repo: &RepoName,
pack: &Hash,
job: &mut VerifyJobV1,
) -> Result<(), StoreError> {
if !job.gone
&& let Some(id) = job.member_body_id
{
let raw = store
.get(source, &member_key(repo, pack, &id))
.await?
.ok_or_else(|| StoreError::Unavailable("job member body disappeared".into()))?;
let Some((&CODEC_V1, bytes)) = raw.as_bytes().split_first() else {
return Err(StoreError::Corrupt("bad job member body version".into()));
};
if member_id(raw.as_bytes()) != id {
return Err(StoreError::Corrupt("bad job member body digest".into()));
}
let lists: MemberLists = serde_json::from_slice(bytes)
.map_err(|_| StoreError::Corrupt("bad job member body".into()))?;
if lists.satisfying.len() > crate::store::index::MAX_LOOKUP_IDS
|| lists.packlist.len()
> crate::store::index::MAX_LOOKUP_IDS
+ crate::store::outbox::MAX_TICKETS_PER_ADVANCE
{
return Err(StoreError::Corrupt("oversized job member body".into()));
}
job.satisfying = lists.satisfying;
job.packlist = lists.packlist;
}
job.members_loaded = true;
Ok(())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct FrameRow {
pub value: IndexValue,
pub object_type: u8,
pub external: Option<Hash>,
}
pub fn encode_frame(id: &Hash, row: &FrameRow) -> Result<Value, StoreError> {
let index = crate::store::codec::encode_object_index(id, &row.value)?;
let mut bytes = vec![row.object_type, u8::from(row.external.is_some())];
if let Some(base) = &row.external {
bytes.extend_from_slice(base);
}
bytes.extend_from_slice(index.as_bytes());
Ok(Value::new(bytes))
}
pub fn decode_frame(id: &Hash, value: &Value) -> Result<FrameRow, StoreError> {
let corrupt = || StoreError::Corrupt("bad verification frame row".into());
let bytes = value.as_bytes();
let (&object_type, rest) = bytes.split_first().ok_or_else(corrupt)?;
let (external, rest) = match rest.split_first() {
Some((0, rest)) => (None, rest),
Some((1, rest)) => {
let (base, rest) = rest.split_first_chunk::<32>().ok_or_else(corrupt)?;
(Some(*base), rest)
}
_ => return Err(corrupt()),
};
Ok(FrameRow {
value: crate::store::codec::decode_object_index(id, &Value::new(rest.to_vec()))?,
object_type,
external,
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct BaseRow {
pub size: u64,
pub depth: u32,
pub entry: u64,
}
#[must_use]
pub fn encode_base(row: &BaseRow) -> Value {
let mut bytes = row.size.to_be_bytes().to_vec();
bytes.extend_from_slice(&row.depth.to_be_bytes());
bytes.extend_from_slice(&row.entry.to_be_bytes());
Value::new(bytes)
}
pub fn decode_base(value: &Value) -> Result<BaseRow, StoreError> {
let bytes: &[u8; 20] = value
.as_bytes()
.try_into()
.map_err(|_| StoreError::Corrupt("bad verification base row".into()))?;
Ok(BaseRow {
size: u64::from_be_bytes(bytes[..8].try_into().unwrap_or_default()),
depth: u32::from_be_bytes(bytes[8..12].try_into().unwrap_or_default()),
entry: u64::from_be_bytes(bytes[12..].try_into().unwrap_or_default()),
})
}
#[must_use]
pub fn index_entry(id: Hash, row: &FrameRow) -> IndexEntry {
IndexEntry {
object: id,
value: row.value,
}
}
#[must_use]
pub fn timer_reference(repo: &RepoName, pack: &Hash) -> Vec<u8> {
let mut reference = repo.as_str().as_bytes().to_vec();
reference.push(0);
reference.extend_from_slice(pack);
reference
}
#[must_use]
pub fn parse_reference(reference: &[u8]) -> Option<(RepoName, Hash)> {
let sep = reference.iter().position(|&b| b == 0)?;
let pack: Hash = reference[sep + 1..].try_into().ok()?;
let repo = RepoName::new(String::from_utf8(reference[..sep].to_vec()).ok()?).ok()?;
Some((repo, pack))
}
type Stored<T> = Option<(T, Value)>;
pub async fn read_job<S: NamespaceStore>(
store: &S,
source: &Partition,
repo: &RepoName,
pack: &Hash,
) -> Result<(Stored<VerifyJobV1>, Stored<VerificationV1>), StoreError> {
let rows = store
.get_many(
source,
&[keys::verify_job(repo, pack), keys::verification(repo, pack)],
)
.await?;
let [job, state] = <[_; 2]>::try_from(rows)
.map_err(|_| StoreError::Corrupt("short verification read".into()))?;
let job = if let Some(raw) = job {
let mut job = decode_job(&raw)?;
hydrate_job(store, source, repo, pack, &mut job).await?;
Some((job, raw))
} else {
None
};
Ok((
job,
state
.map(|raw| super::state::decode(&raw).map(|state| (state, raw)))
.transpose()?,
))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn job_frame_and_base_codecs_round_trip() {
let mut job = VerifyJobV1::new([1; 32], 5, 99, 4096);
job.cursor = vec![0xab, 0x01];
job.satisfying = vec![[2; 32]];
job.outcome = Some(Outcome::BaseMissing);
let header = decode_job(&encode_job(&job)).unwrap();
assert!(header.satisfying.is_empty());
assert!(!header.members_loaded);
assert_eq!(header.outcome, job.outcome);
assert!(decode_job(&Value::new(b"\x02{}".to_vec())).is_err());
let frame = FrameRow {
value: IndexValue {
frame_offset: 12,
frame_length: 40,
wire_type: 0x02,
decoded_size: 7,
chain_depth: 2,
delta_base: Some([3; 32]),
},
object_type: 3,
external: Some([4; 32]),
};
let id = [9; 32];
assert_eq!(
decode_frame(&id, &encode_frame(&id, &frame).unwrap()).unwrap(),
frame
);
let base = BaseRow {
size: 1,
depth: 2,
entry: 3,
};
assert_eq!(decode_base(&encode_base(&base)).unwrap(), base);
let name = RepoName::new("a").unwrap();
assert_eq!(
parse_reference(&timer_reference(&name, &[7; 32])),
Some((name, [7; 32]))
);
job.restart();
assert_eq!(
(job.restarts, job.phase, job.cursor.len()),
(1, Phase::Decode, 0)
);
}
}