Skip to main content

mkit_server/indexed/
state.rs

1//! Per-(repository, pack) verification state, outside the advance batch.
2
3use crate::repo::RepoName;
4use crate::store::{codec::CODEC_V1, keys};
5use crate::{
6    Batch, BatchOutcome, NamespaceStore, Partition, Precondition, ServerError, StoreError, Value,
7};
8use mkit_core::hash::Hash;
9use serde::{Deserialize, Serialize};
10
11/// Native inline verifier lease: long enough for ordinary packs, bounded so
12/// a crashed verifier can be retried. Operators cap concurrent decodes.
13pub const VERIFICATION_LEASE_MS: u64 = 30_000;
14
15/// Durable verification result. Only failures determined by pack content
16/// alone are `Rejected`; repository membership misses are retried.
17#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
18#[serde(tag = "state", rename_all = "snake_case", deny_unknown_fields)]
19pub enum VerificationV1 {
20    Pending {
21        lease_until_ms: u64,
22    },
23    Verified {
24        pack_len: u64,
25        verified_at_ms: u64,
26        #[serde(default, skip_serializing_if = "Option::is_none")]
27        publication: Option<Box<super::publication::resume::Progress>>,
28    },
29    Rejected {
30        code: String,
31        message: String,
32    },
33}
34
35/// Encode one state with the metadata codec version byte.
36///
37/// # Panics
38/// Serialization of this fixed integer-and-string DTO into a `Vec` cannot
39/// fail; a failure here indicates a broken serializer invariant.
40#[must_use]
41pub fn encode(value: &VerificationV1) -> Value {
42    let mut bytes = vec![CODEC_V1];
43    serde_json::to_writer(&mut bytes, value).expect("verification DTO serializes");
44    Value::new(bytes)
45}
46
47/// Decode a state, failing closed on corrupt or future values.
48pub fn decode(value: &Value) -> Result<VerificationV1, StoreError> {
49    let Some((&CODEC_V1, body)) = value.as_bytes().split_first() else {
50        return Err(StoreError::Corrupt("bad verification version".into()));
51    };
52    let decoded: VerificationV1 = serde_json::from_slice(body)
53        .map_err(|_| StoreError::Corrupt("bad verification state".into()))?;
54    if let VerificationV1::Verified {
55        publication: Some(progress),
56        ..
57    } = &decoded
58    {
59        if value.as_bytes().len() > super::publication::resume::MAX_STATE_BYTES {
60            return Err(StoreError::Corrupt(
61                "publication checkpoint too large".into(),
62            ));
63        }
64        progress.validate()?;
65    }
66    Ok(decoded)
67}
68
69/// Three-operation state CAS: `NotAfter`, prior-value guard, and put.
70pub async fn write<S: NamespaceStore>(
71    store: &S,
72    source: &Partition,
73    repo: &RepoName,
74    pack: &Hash,
75    prior: Option<&Value>,
76    state: &VerificationV1,
77    deadline_ms: u64,
78) -> Result<bool, StoreError> {
79    let key = keys::verification(repo, pack);
80    let guard = match prior {
81        Some(value) => Precondition::Equals(key.clone(), value.clone()),
82        None => Precondition::Absent(key.clone()),
83    };
84    let batch = Batch::new()
85        .require(Precondition::NotAfter(deadline_ms))
86        .require(guard)
87        .put(key, encode(state));
88    match store.apply(source, batch).await? {
89        BatchOutcome::Committed => Ok(true),
90        BatchOutcome::PreconditionFailed { .. } => Ok(false),
91        BatchOutcome::DeadlinePassed { .. } => Err(StoreError::unavailable(std::io::Error::other(
92            "verification deadline passed",
93        ))),
94    }
95}
96
97/// Release only the pending state this invocation wrote. A concurrent
98/// verifier's newer state, or a content rejection, survives the CAS loss.
99pub async fn clear_pending<S: NamespaceStore>(
100    store: &S,
101    source: &Partition,
102    repo: &RepoName,
103    pack: &Hash,
104    pending_raw: &Value,
105    deadline_ms: u64,
106) -> Result<(), StoreError> {
107    let key = keys::verification(repo, pack);
108    let batch = Batch::new()
109        .require(Precondition::NotAfter(deadline_ms))
110        .require(Precondition::Equals(key.clone(), pending_raw.clone()))
111        .delete(key);
112    let _ = store.apply(source, batch).await?;
113    Ok(())
114}
115
116/// Read the current state and its exact CAS value.
117pub async fn read<S: NamespaceStore>(
118    store: &S,
119    source: &Partition,
120    repo: &RepoName,
121    pack: &Hash,
122) -> Result<Option<(VerificationV1, Value)>, StoreError> {
123    store
124        .get(source, &keys::verification(repo, pack))
125        .await?
126        .map(|raw| Ok((decode(&raw)?, raw)))
127        .transpose()
128}
129
130/// Live lease held by another verifier, with no replay outcome.
131#[must_use]
132pub fn concurrent_pending(state: &VerificationV1, now_ms: u64) -> Option<ServerError> {
133    match state {
134        VerificationV1::Pending { lease_until_ms } if *lease_until_ms > now_ms => {
135            Some(super::pending(lease_until_ms.saturating_sub(now_ms)))
136        }
137        _ => None,
138    }
139}
140
141#[cfg(test)]
142mod tests {
143    use super::*;
144
145    #[test]
146    fn codec_golden_and_live_lease_boundary() {
147        let state = VerificationV1::Verified {
148            pack_len: 123,
149            verified_at_ms: 456,
150            publication: None,
151        };
152        let bytes = encode(&state);
153        assert_eq!(
154            bytes.as_bytes(),
155            b"\x01{\"state\":\"verified\",\"pack_len\":123,\"verified_at_ms\":456}"
156        );
157        assert_eq!(decode(&bytes).unwrap(), state);
158        assert!(decode(&Value::new(b"\x02{}".to_vec())).is_err());
159        let pending = VerificationV1::Pending {
160            lease_until_ms: 1000,
161        };
162        assert!(concurrent_pending(&pending, 999).is_some());
163        assert!(concurrent_pending(&pending, 1000).is_none());
164    }
165}