mkit_server/indexed/
state.rs1use 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
11pub const VERIFICATION_LEASE_MS: u64 = 30_000;
14
15#[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#[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
47pub 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
69pub 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
97pub 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
116pub 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#[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}