Skip to main content

canwu_sim/runtime/
page_store.rs

1//! Content-addressed state pages used by the Format 8 checkpoint boundary.
2//!
3//! Pages are deliberately domain-neutral: the runtime commits the canonical
4//! bytes and a host decides where (and whether) to physically store them. A
5//! page ID is the hash of the canonical uncompressed bytes, so compression or
6//! blob placement cannot change semantic identity.
7
8use 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;
15/// Hard predecode cap for one initial or incremental Format-8 page graph.
16/// A one-million-entry non-collision Patricia map can contain `2N - 1`
17/// logical node pages; this leaves bounded room for its manifest and compact
18/// decision buckets without making the count unbounded.
19pub 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/// Persistable host-side mark/sweep interlock for content-addressed state
112/// pages. Preparing, verifying, or durably enqueueing a root protects every
113/// declared reachable page across process restart. A committed root takes over
114/// that lease atomically before the transient handle may disappear.
115#[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}