1use super::{ArchiveStoreOutcome, CanwuError, ErrorCode, canonical_byte_hash};
9use serde::{Deserialize, Serialize};
10use std::collections::{BTreeMap, BTreeSet};
11
12pub const STATE_PAGE_FORMAT_VERSION: u32 = 1;
13pub const STATE_PAGE_CODEC: &str = "raw-canonical-v1";
14pub const MAX_STATE_PAGE_BYTES: usize = 4 * 1024 * 1024;
15pub const MAX_STATE_DELTA_PAGES: usize = 4_194_304;
20pub const STATE_PAGE_RETENTION_FORMAT_VERSION: u32 = 1;
21
22#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
23#[serde(deny_unknown_fields)]
24pub struct StatePageBlob {
25 pub format_version: u32,
26 pub page_id: String,
27 pub codec: String,
28 pub decoded_bytes: u64,
29 pub bytes: Vec<u8>,
30}
31
32impl StatePageBlob {
33 pub fn new(bytes: Vec<u8>) -> Result<Self, CanwuError> {
34 if bytes.is_empty() || bytes.len() > MAX_STATE_PAGE_BYTES {
35 return Err(page_error("state page bytes exceed the bounded page limit"));
36 }
37 let decoded_bytes = u64::try_from(bytes.len())
38 .map_err(|_| page_error("state page byte count is not representable"))?;
39 let page_id = state_page_id(&bytes);
40 Ok(Self {
41 format_version: STATE_PAGE_FORMAT_VERSION,
42 page_id,
43 codec: STATE_PAGE_CODEC.to_owned(),
44 decoded_bytes,
45 bytes,
46 })
47 }
48
49 pub fn validate(&self) -> Result<(), CanwuError> {
50 if self.format_version != STATE_PAGE_FORMAT_VERSION {
51 return Err(page_error(format!(
52 "state page format {} is unsupported; expected {STATE_PAGE_FORMAT_VERSION}",
53 self.format_version
54 )));
55 }
56 if self.codec != STATE_PAGE_CODEC {
57 return Err(page_error(
58 "state page codec is not the canonical runtime codec",
59 ));
60 }
61 if self.bytes.is_empty() || self.bytes.len() > MAX_STATE_PAGE_BYTES {
62 return Err(page_error("state page bytes exceed the bounded page limit"));
63 }
64 if self.decoded_bytes != self.bytes.len() as u64 {
65 return Err(page_error("state page decoded byte count is inconsistent"));
66 }
67 if self.page_id != state_page_id(&self.bytes) {
68 return Err(page_error(
69 "state page ID does not match its canonical bytes",
70 ));
71 }
72 Ok(())
73 }
74}
75
76#[must_use]
77pub fn state_page_id(bytes: &[u8]) -> String {
78 canonical_byte_hash("canwu.state-page.v1", bytes)
79}
80
81pub trait StatePageProvider {
82 fn load_state_page(&self, page_id: &str) -> Result<Option<StatePageBlob>, CanwuError>;
83}
84
85pub trait StatePageStore: StatePageProvider {
86 fn store_state_page(&self, page: &StatePageBlob) -> Result<ArchiveStoreOutcome, CanwuError>;
87}
88
89#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
90#[serde(rename_all = "snake_case")]
91pub enum StatePageRetentionPhase {
92 Prepared,
93 Verified,
94 DurableIngress,
95 Committed,
96 Abandoned,
97}
98
99#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
100#[serde(deny_unknown_fields)]
101pub struct StatePageRetentionHandle {
102 pub format_version: u32,
103 pub handle_id: String,
104 pub source_root: String,
105 pub target_root: String,
106 pub page_ids: BTreeSet<String>,
107 pub prepared_epoch: u64,
108 pub phase: StatePageRetentionPhase,
109}
110
111#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
116#[serde(deny_unknown_fields)]
117pub struct StatePageRetentionLedger {
118 pub format_version: u32,
119 pub gc_epoch: u64,
120 pub handles: BTreeMap<String, StatePageRetentionHandle>,
121 pub committed_roots: BTreeMap<String, BTreeSet<String>>,
122}
123
124impl Default for StatePageRetentionLedger {
125 fn default() -> Self {
126 Self {
127 format_version: STATE_PAGE_RETENTION_FORMAT_VERSION,
128 gc_epoch: 0,
129 handles: BTreeMap::new(),
130 committed_roots: BTreeMap::new(),
131 }
132 }
133}
134
135impl StatePageRetentionLedger {
136 pub fn prepare(
137 &mut self,
138 delta: &PreparedStateDelta,
139 provider: &dyn StatePageProvider,
140 ) -> Result<String, CanwuError> {
141 self.validate()?;
142 delta.validate()?;
143 let reachable_page_ids = state_page_closure(delta, provider)?;
144 for page_id in &reachable_page_ids {
145 validate_hash(page_id, "retained state page ID")?;
146 }
147 let handle_id = canonical_byte_hash(
148 "canwu.state-page-retention-handle.v1",
149 &serde_json::to_vec(&(
150 &delta.token_hash,
151 &delta.source_root,
152 &delta.target_root,
153 &reachable_page_ids,
154 self.gc_epoch,
155 ))
156 .map_err(|error| page_error(format!("cannot encode retention handle: {error}")))?,
157 );
158 let handle = StatePageRetentionHandle {
159 format_version: STATE_PAGE_RETENTION_FORMAT_VERSION,
160 handle_id: handle_id.clone(),
161 source_root: delta.source_root.clone(),
162 target_root: delta.target_root.clone(),
163 page_ids: reachable_page_ids,
164 prepared_epoch: self.gc_epoch,
165 phase: StatePageRetentionPhase::Prepared,
166 };
167 if let Some(existing) = self.handles.get(&handle_id) {
168 if existing != &handle {
169 return Err(page_error(
170 "state page retention handle collides with different content",
171 ));
172 }
173 return Ok(handle_id);
174 }
175 self.handles.insert(handle_id.clone(), handle);
176 Ok(handle_id)
177 }
178
179 pub fn verify(
180 &mut self,
181 handle_id: &str,
182 provider: &dyn StatePageProvider,
183 ) -> Result<(), CanwuError> {
184 let handle = self
185 .handles
186 .get(handle_id)
187 .cloned()
188 .ok_or_else(|| page_error("state page retention handle is unknown"))?;
189 if !matches!(
190 handle.phase,
191 StatePageRetentionPhase::Prepared | StatePageRetentionPhase::Verified
192 ) {
193 return Err(page_error(
194 "state page retention handle cannot enter verified state",
195 ));
196 }
197 for page_id in &handle.page_ids {
198 let page = provider.load_state_page(page_id)?.ok_or_else(|| {
199 CanwuError::new(
200 ErrorCode::StatePageUnavailable,
201 "retained state page is unavailable",
202 )
203 })?;
204 page.validate()?;
205 if page.page_id != *page_id {
206 return Err(page_error("retained state page identity changed"));
207 }
208 }
209 let observed =
210 state_page_closure_for_root(&handle.target_root, provider, &BTreeMap::new())?;
211 if observed != handle.page_ids {
212 return Err(page_error(
213 "retained state-page closure changed after preparation",
214 ));
215 }
216 self.handles
217 .get_mut(handle_id)
218 .ok_or_else(|| page_error("state page retention handle disappeared during verify"))?
219 .phase = StatePageRetentionPhase::Verified;
220 Ok(())
221 }
222
223 pub fn mark_durable_ingress(&mut self, handle_id: &str) -> Result<(), CanwuError> {
224 self.transition(
225 handle_id,
226 StatePageRetentionPhase::Verified,
227 StatePageRetentionPhase::DurableIngress,
228 )
229 }
230
231 pub fn commit(&mut self, handle_id: &str) -> Result<(), CanwuError> {
232 let handle = self
233 .handles
234 .get(handle_id)
235 .cloned()
236 .ok_or_else(|| page_error("state page retention handle is unknown"))?;
237 if handle.phase != StatePageRetentionPhase::DurableIngress
238 && handle.phase != StatePageRetentionPhase::Committed
239 {
240 return Err(page_error(
241 "only durable maintenance ingress may commit a state-page root",
242 ));
243 }
244 if let Some(existing) = self.committed_roots.get(&handle.target_root)
245 && existing != &handle.page_ids
246 {
247 return Err(page_error(
248 "committed state root is bound to different reachable pages",
249 ));
250 }
251 self.committed_roots
252 .insert(handle.target_root.clone(), handle.page_ids.clone());
253 self.handles
254 .get_mut(handle_id)
255 .ok_or_else(|| page_error("state page retention handle disappeared during commit"))?
256 .phase = StatePageRetentionPhase::Committed;
257 Ok(())
258 }
259
260 pub fn abandon(&mut self, handle_id: &str) -> Result<(), CanwuError> {
261 let handle = self
262 .handles
263 .get_mut(handle_id)
264 .ok_or_else(|| page_error("state page retention handle is unknown"))?;
265 if matches!(
266 handle.phase,
267 StatePageRetentionPhase::DurableIngress | StatePageRetentionPhase::Committed
268 ) {
269 return Err(page_error(
270 "durable or committed state-page retention cannot be abandoned",
271 ));
272 }
273 handle.phase = StatePageRetentionPhase::Abandoned;
274 Ok(())
275 }
276
277 pub fn release_committed_root(&mut self, root: &str) -> Result<(), CanwuError> {
278 validate_hash(root, "released state root")?;
279 self.committed_roots.remove(root);
280 self.handles.retain(|_, handle| {
281 !(handle.phase == StatePageRetentionPhase::Committed && handle.target_root == root)
282 });
283 self.validate()?;
284 Ok(())
285 }
286
287 pub fn begin_gc_epoch(&mut self) -> Result<u64, CanwuError> {
288 self.gc_epoch = self
289 .gc_epoch
290 .checked_add(1)
291 .ok_or_else(|| page_error("state-page GC epoch is exhausted"))?;
292 Ok(self.gc_epoch)
293 }
294
295 #[must_use]
296 pub fn reachable_page_ids(&self) -> BTreeSet<String> {
297 let mut reachable = self
298 .committed_roots
299 .values()
300 .flat_map(|pages| pages.iter().cloned())
301 .collect::<BTreeSet<_>>();
302 for handle in self.handles.values().filter(|handle| {
303 !matches!(
304 handle.phase,
305 StatePageRetentionPhase::Abandoned | StatePageRetentionPhase::Committed
306 )
307 }) {
308 reachable.extend(handle.page_ids.iter().cloned());
309 }
310 reachable
311 }
312
313 #[must_use]
314 pub fn sweep_candidates(
315 &self,
316 all_page_ids: impl IntoIterator<Item = String>,
317 ) -> BTreeSet<String> {
318 let reachable = self.reachable_page_ids();
319 all_page_ids
320 .into_iter()
321 .filter(|page_id| !reachable.contains(page_id))
322 .collect()
323 }
324
325 pub fn validate(&self) -> Result<(), CanwuError> {
326 if self.format_version != STATE_PAGE_RETENTION_FORMAT_VERSION {
327 return Err(page_error("unsupported state-page retention format"));
328 }
329 for (handle_id, handle) in &self.handles {
330 validate_hash(handle_id, "state-page retention handle ID")?;
331 validate_hash(&handle.source_root, "retention source root")?;
332 validate_hash(&handle.target_root, "retention target root")?;
333 if handle.format_version != STATE_PAGE_RETENTION_FORMAT_VERSION
334 || handle.handle_id != *handle_id
335 || handle.page_ids.is_empty()
336 || handle.prepared_epoch > self.gc_epoch
337 {
338 return Err(page_error("state-page retention handle is inconsistent"));
339 }
340 for page_id in &handle.page_ids {
341 validate_hash(page_id, "retained state page ID")?;
342 }
343 if handle.phase == StatePageRetentionPhase::Committed
344 && self.committed_roots.get(&handle.target_root) != Some(&handle.page_ids)
345 {
346 return Err(page_error(
347 "committed retention handle did not transfer its lease",
348 ));
349 }
350 }
351 for (root, pages) in &self.committed_roots {
352 validate_hash(root, "committed state root")?;
353 if pages.is_empty() {
354 return Err(page_error("committed state root has no reachable pages"));
355 }
356 for page_id in pages {
357 validate_hash(page_id, "committed state page ID")?;
358 }
359 }
360 Ok(())
361 }
362
363 fn transition(
364 &mut self,
365 handle_id: &str,
366 expected: StatePageRetentionPhase,
367 next: StatePageRetentionPhase,
368 ) -> Result<(), CanwuError> {
369 let handle = self
370 .handles
371 .get_mut(handle_id)
372 .ok_or_else(|| page_error("state page retention handle is unknown"))?;
373 if handle.phase == next {
374 return Ok(());
375 }
376 if handle.phase != expected {
377 return Err(page_error("state page retention transition is invalid"));
378 }
379 handle.phase = next;
380 Ok(())
381 }
382}
383
384fn state_page_closure(
385 delta: &PreparedStateDelta,
386 provider: &dyn StatePageProvider,
387) -> Result<BTreeSet<String>, CanwuError> {
388 let new_pages = delta
389 .new_pages
390 .iter()
391 .map(|page| (page.page_id.clone(), page.clone()))
392 .collect::<BTreeMap<_, _>>();
393 let reachable = state_page_closure_for_root(&delta.target_root, provider, &new_pages)?;
394 if delta
395 .new_pages
396 .iter()
397 .any(|page| !reachable.contains(&page.page_id))
398 {
399 return Err(page_error(
400 "prepared state delta contains a page outside the target-root closure",
401 ));
402 }
403 Ok(reachable)
404}
405
406fn state_page_closure_for_root(
407 root: &str,
408 provider: &dyn StatePageProvider,
409 new_pages: &BTreeMap<String, StatePageBlob>,
410) -> Result<BTreeSet<String>, CanwuError> {
411 let mut reachable = BTreeSet::new();
412 let mut pending = vec![root.to_owned()];
413 while let Some(page_id) = pending.pop() {
414 if !reachable.insert(page_id.clone()) {
415 continue;
416 }
417 if reachable.len() > MAX_STATE_DELTA_PAGES {
418 return Err(page_error("state-page closure exceeds the hard page limit"));
419 }
420 let page = match new_pages.get(&page_id).cloned() {
421 Some(page) => page,
422 None => provider.load_state_page(&page_id)?.ok_or_else(|| {
423 CanwuError::new(
424 ErrorCode::StatePageUnavailable,
425 "state-page closure references an unavailable page",
426 )
427 })?,
428 };
429 page.validate()?;
430 if page.page_id != page_id {
431 return Err(page_error(
432 "state-page closure provider changed page identity",
433 ));
434 }
435 let value = serde_json::from_slice::<serde_json::Value>(&page.bytes).map_err(|error| {
436 page_error(format!(
437 "state page cannot be decoded for reachability: {error}"
438 ))
439 })?;
440 pending.extend(state_page_children(&value)?);
441 }
442 Ok(reachable)
443}
444
445fn state_page_children(value: &serde_json::Value) -> Result<Vec<String>, CanwuError> {
446 let Some(object) = value.as_object() else {
447 return Ok(Vec::new());
448 };
449 let mut children = Vec::new();
450 if object.contains_key("checkpoint_without_paged_state")
451 && object.contains_key("domain_records")
452 {
453 let records = object
454 .get("domain_records")
455 .and_then(serde_json::Value::as_object)
456 .ok_or_else(|| page_error("paged checkpoint domain roots are malformed"))?;
457 for field in [
458 "primary",
459 "reverse_references",
460 "successor_of",
461 "predecessors_of",
462 ] {
463 if let Some(page_id) = records.get(field).and_then(serde_json::Value::as_str) {
464 children.push(page_id.to_owned());
465 }
466 }
467 if let Some(page_id) = object
468 .get("decision_manifest_page_id")
469 .and_then(serde_json::Value::as_str)
470 {
471 children.push(page_id.to_owned());
472 }
473 } else if let Some(page_id) = object
474 .get("hot_page_id")
475 .and_then(serde_json::Value::as_str)
476 {
477 children.push(page_id.to_owned());
478 if let Some(pages) = object
479 .get("archive_directory_page_ids")
480 .and_then(serde_json::Value::as_array)
481 {
482 children.extend(
483 pages
484 .iter()
485 .filter_map(serde_json::Value::as_str)
486 .map(str::to_owned),
487 );
488 }
489 } else if let Some(pages) = object
490 .get("archive_bucket_pages")
491 .and_then(serde_json::Value::as_array)
492 {
493 children.extend(pages.iter().filter_map(|entry| {
494 entry
495 .as_array()
496 .and_then(|pair| pair.get(1))
497 .and_then(serde_json::Value::as_str)
498 .map(str::to_owned)
499 }));
500 } else if object.get("node").and_then(serde_json::Value::as_str) == Some("branch") {
501 for field in ["left_page", "right_page"] {
502 let page_id = object
503 .get(field)
504 .and_then(serde_json::Value::as_str)
505 .ok_or_else(|| page_error("Patricia branch page is missing a child"))?;
506 children.push(page_id.to_owned());
507 }
508 }
509 for page_id in &children {
510 validate_hash(page_id, "reachable state page ID")?;
511 }
512 Ok(children)
513}
514
515#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
516#[serde(deny_unknown_fields)]
517pub struct PreparedStateDelta {
518 pub format_version: u32,
519 pub source_root: String,
520 pub target_root: String,
521 pub new_pages: Vec<StatePageBlob>,
522 pub token_hash: String,
523}
524
525impl PreparedStateDelta {
526 pub fn validate(&self) -> Result<(), CanwuError> {
527 if self.format_version != STATE_PAGE_FORMAT_VERSION {
528 return Err(page_error(
529 "prepared state delta uses an unsupported format",
530 ));
531 }
532 validate_hash(&self.source_root, "state delta source root")?;
533 validate_hash(&self.target_root, "state delta target root")?;
534 validate_hash(&self.token_hash, "state delta token")?;
535 if self.new_pages.len() > MAX_STATE_DELTA_PAGES {
536 return Err(page_error("prepared state delta contains too many pages"));
537 }
538 let mut ids = BTreeSet::new();
539 for page in &self.new_pages {
540 page.validate()?;
541 if !ids.insert(page.page_id.clone()) {
542 return Err(page_error("prepared state delta contains duplicate pages"));
543 }
544 }
545 let expected =
546 canonical_hash_for_delta(&self.source_root, &self.target_root, &self.new_pages);
547 if self.token_hash != expected {
548 return Err(page_error("prepared state delta token is inconsistent"));
549 }
550 Ok(())
551 }
552}
553
554pub fn prepare_state_delta(
555 source_root: &str,
556 target_root: &str,
557 pages: Vec<StatePageBlob>,
558) -> Result<PreparedStateDelta, CanwuError> {
559 validate_hash(source_root, "state delta source root")?;
560 validate_hash(target_root, "state delta target root")?;
561 let prepared = PreparedStateDelta {
562 format_version: STATE_PAGE_FORMAT_VERSION,
563 source_root: source_root.to_owned(),
564 target_root: target_root.to_owned(),
565 new_pages: pages,
566 token_hash: String::new(),
567 };
568 let token_hash = canonical_hash_for_delta(
569 &prepared.source_root,
570 &prepared.target_root,
571 &prepared.new_pages,
572 );
573 let prepared = PreparedStateDelta {
574 token_hash,
575 ..prepared
576 };
577 prepared.validate()?;
578 Ok(prepared)
579}
580
581pub fn verify_state_delta(
582 prepared: &PreparedStateDelta,
583 provider: &dyn StatePageProvider,
584) -> Result<(), CanwuError> {
585 prepared.validate()?;
586 for page in &prepared.new_pages {
587 let loaded = provider.load_state_page(&page.page_id)?.ok_or_else(|| {
588 CanwuError::new(ErrorCode::StatePageUnavailable, "state page is unavailable")
589 })?;
590 loaded.validate()?;
591 if loaded != *page {
592 return Err(page_error(
593 "provider returned bytes different from the prepared page",
594 ));
595 }
596 }
597 Ok(())
598}
599
600fn canonical_hash_for_delta(
601 source_root: &str,
602 target_root: &str,
603 pages: &[StatePageBlob],
604) -> String {
605 let mut bytes = Vec::new();
606 bytes.extend_from_slice(source_root.as_bytes());
607 bytes.push(0);
608 bytes.extend_from_slice(target_root.as_bytes());
609 bytes.push(0);
610 for page in pages {
611 bytes.extend_from_slice(page.page_id.as_bytes());
612 bytes.push(0);
613 }
614 canonical_byte_hash("canwu.state-delta.v1", &bytes)
615}
616
617fn validate_hash(value: &str, label: &str) -> Result<(), CanwuError> {
618 if value.len() != 64
619 || value
620 .bytes()
621 .any(|byte| !byte.is_ascii_hexdigit() || byte.is_ascii_uppercase())
622 {
623 return Err(page_error(format!(
624 "{label} must be a lower-case 32-byte hash"
625 )));
626 }
627 Ok(())
628}
629
630fn page_error(message: impl Into<String>) -> CanwuError {
631 CanwuError::new(ErrorCode::InvalidArchive, message)
632}