1use crate::model::{canonicalize_options, require_text};
2use crate::{
3 DecisionAction, DecisionAttemptOutcome, DecisionAttemptRecord, DecisionControllerBinding,
4 DecisionError, DecisionErrorCode, DecisionMutation, DecisionOutcome, DecisionPolicy,
5 DecisionPolicyIdentity, DecisionTicket, DecisionTicketState, DecisionTrace, PolicyDecision,
6};
7use canwu_core::{CommandRequestId, DecisionTicketId, DecisionTraceId};
8use canwu_time::SimTime;
9use im::{OrdMap, OrdSet};
10use serde::{Deserialize, Serialize};
11use std::cell::RefCell;
12use std::collections::BTreeMap as StdBTreeMap;
13use std::fmt::Write as _;
14use std::mem::size_of;
15use std::ops::Index;
16use std::sync::Arc;
17
18mod persistent_decision_map_serde {
19 use im::OrdMap;
20 use serde::ser::SerializeMap as _;
21 use serde::{Deserialize, Deserializer, Serialize, Serializer};
22 use std::collections::BTreeMap;
23 use std::sync::Arc;
24
25 pub fn serialize<S, K, V>(entries: &OrdMap<K, Arc<V>>, serializer: S) -> Result<S::Ok, S::Error>
26 where
27 S: Serializer,
28 K: Clone + Ord + Serialize,
29 V: Serialize,
30 {
31 let mut map = serializer.serialize_map(Some(entries.len()))?;
32 for (ordinal, entry) in entries {
33 map.serialize_entry(ordinal, entry.as_ref())?;
34 }
35 map.end()
36 }
37
38 pub fn deserialize<'de, D, K, V>(deserializer: D) -> Result<OrdMap<K, Arc<V>>, D::Error>
39 where
40 D: Deserializer<'de>,
41 K: Clone + Deserialize<'de> + Ord,
42 V: Deserialize<'de>,
43 {
44 Ok(BTreeMap::<K, V>::deserialize(deserializer)?
45 .into_iter()
46 .map(|(ordinal, entry)| (ordinal, Arc::new(entry)))
47 .collect())
48 }
49}
50
51mod decision_archive_receipts_serde {
52 use super::{
53 CompactDecisionArchiveReceipt, DecisionArchivePageKey, DecisionArchiveReceipt,
54 DecisionHistoryKey,
55 };
56 use im::OrdMap;
57 use serde::de::Error as _;
58 use serde::{Deserialize, Deserializer, Serialize, Serializer};
59
60 pub fn serialize<S>(
61 buckets: &OrdMap<
62 DecisionArchivePageKey,
63 OrdMap<DecisionHistoryKey, CompactDecisionArchiveReceipt>,
64 >,
65 serializer: S,
66 ) -> Result<S::Ok, S::Error>
67 where
68 S: Serializer,
69 {
70 buckets
71 .values()
72 .flat_map(OrdMap::iter)
73 .map(|(key, receipt)| receipt.to_receipt(key))
74 .collect::<Vec<_>>()
75 .serialize(serializer)
76 }
77
78 pub fn deserialize<'de, D>(
79 deserializer: D,
80 ) -> Result<
81 OrdMap<DecisionArchivePageKey, OrdMap<DecisionHistoryKey, CompactDecisionArchiveReceipt>>,
82 D::Error,
83 >
84 where
85 D: Deserializer<'de>,
86 {
87 let entries = Vec::<DecisionArchiveReceipt>::deserialize(deserializer)?;
88 let mut buckets = OrdMap::<
89 DecisionArchivePageKey,
90 OrdMap<DecisionHistoryKey, CompactDecisionArchiveReceipt>,
91 >::new();
92 for receipt in entries {
93 let key = receipt.key.clone();
94 let bucket = super::decision_history_page_key(&key).map_err(D::Error::custom)?;
95 let value =
96 CompactDecisionArchiveReceipt::from_receipt(&receipt).map_err(D::Error::custom)?;
97 if buckets
98 .entry(bucket)
99 .or_default()
100 .insert(key, value)
101 .is_some()
102 {
103 return Err(D::Error::custom(
104 "decision archive receipt map contains a duplicate key",
105 ));
106 }
107 }
108 Ok(buckets)
109 }
110}
111
112mod decision_archive_page_directory_serde {
113 use super::DecisionArchivePageKey;
114 use im::OrdMap;
115 use serde::de::Error as _;
116 use serde::{Deserialize, Deserializer, Serialize, Serializer};
117
118 pub fn serialize<S>(
119 pages: &OrdMap<DecisionArchivePageKey, String>,
120 serializer: S,
121 ) -> Result<S::Ok, S::Error>
122 where
123 S: Serializer,
124 {
125 pages
126 .iter()
127 .map(|(key, id)| (*key, id))
128 .collect::<Vec<_>>()
129 .serialize(serializer)
130 }
131
132 pub fn deserialize<'de, D>(
133 deserializer: D,
134 ) -> Result<OrdMap<DecisionArchivePageKey, String>, D::Error>
135 where
136 D: Deserializer<'de>,
137 {
138 let entries = Vec::<(DecisionArchivePageKey, String)>::deserialize(deserializer)?;
139 let mut pages = OrdMap::new();
140 let mut previous = None;
141 for (key, id) in entries {
142 if previous.is_some_and(|previous| previous >= key) || pages.insert(key, id).is_some() {
143 return Err(D::Error::custom(
144 "decision archive page directory is not strictly ordered",
145 ));
146 }
147 previous = Some(key);
148 }
149 Ok(pages)
150 }
151}
152
153#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
157#[serde(rename_all = "snake_case", tag = "kind", content = "id")]
158pub enum DecisionHistoryKey {
159 Ticket(canwu_core::DecisionTicketId),
160 Attempt(canwu_core::DecisionRequestId),
161 Trace(canwu_core::DecisionTraceId),
162}
163
164#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
165#[serde(rename_all = "snake_case", tag = "location")]
166pub enum DecisionHistoryLocation {
167 Hot,
168 Archived {
169 locator: String,
170 },
171 Unresolved {
174 bucket: u16,
175 segment: u8,
176 },
177 Absent,
178}
179
180pub const DECISION_ARCHIVE_FORMAT_VERSION: u32 = 1;
181pub const MAX_DECISION_ARCHIVE_BATCH_ENTRIES: usize = 4_096;
182pub const MAX_DECISION_HISTORY_PAGE_SIZE: usize = 512;
183pub const MAX_DECISION_HISTORY_PAGE_BYTES: u64 = 16 * 1024 * 1024;
184pub const DECISION_HISTORY_BUCKET_BITS: u8 = 12;
185pub const DECISION_HISTORY_BUCKET_COUNT: u16 = 1 << DECISION_HISTORY_BUCKET_BITS;
186pub const DECISION_ARCHIVE_BUCKET_PAGE_FORMAT_VERSION: u32 = 1;
187pub const MAX_DECISION_ARCHIVE_BUCKET_PAGE_ENTRIES: usize = 64;
188pub const MAX_DECISION_ARCHIVE_BUCKET_PAGE_BYTES: usize = 1024 * 1024;
189
190#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
191#[serde(deny_unknown_fields)]
192pub struct DecisionArchivePageKey {
193 pub bucket: u16,
194 pub segment: u8,
195}
196
197#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
198#[serde(deny_unknown_fields)]
199pub struct DecisionHistoryQueryBudget {
200 pub max_results: usize,
201 pub max_provider_calls: usize,
202 pub max_decoded_bytes: u64,
203}
204
205impl Default for DecisionHistoryQueryBudget {
206 fn default() -> Self {
207 Self {
208 max_results: 128,
209 max_provider_calls: 128,
210 max_decoded_bytes: 4 * 1024 * 1024,
211 }
212 }
213}
214
215#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
216#[serde(deny_unknown_fields)]
217pub struct DecisionHistoryCursor {
218 pub archive_root: String,
219 pub bucket: u16,
220 pub segment: u8,
221 pub after: DecisionHistoryKey,
222}
223
224#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
225#[serde(deny_unknown_fields)]
226pub struct DecisionHistoryPage {
227 pub archive_root: String,
228 pub records: Vec<DecisionArchiveRecord>,
229 #[serde(default, skip_serializing_if = "Option::is_none")]
230 pub next_cursor: Option<DecisionHistoryCursor>,
231 pub provider_calls: u64,
232 pub decoded_bytes: u64,
233}
234
235#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
236#[serde(deny_unknown_fields)]
237pub struct DecisionArchiveReachability {
238 pub bucket_page_ids: OrdSet<String>,
239 pub blob_locators: OrdSet<String>,
240}
241
242#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
243#[serde(deny_unknown_fields)]
244pub struct DecisionArchiveBucketPage {
245 pub format_version: u32,
246 pub bucket: u16,
247 pub segment: u8,
248 pub receipts: Vec<DecisionArchiveReceipt>,
249}
250
251impl DecisionArchiveBucketPage {
252 pub fn validate(&self) -> Result<(), DecisionError> {
253 if self.format_version != DECISION_ARCHIVE_BUCKET_PAGE_FORMAT_VERSION
254 || self.bucket >= DECISION_HISTORY_BUCKET_COUNT
255 || self.receipts.is_empty()
256 || self.receipts.len() > MAX_DECISION_ARCHIVE_BUCKET_PAGE_ENTRIES
257 || self
258 .receipts
259 .windows(2)
260 .any(|pair| pair[0].key >= pair[1].key)
261 || self.receipts.iter().any(|receipt| {
262 decision_history_page_key(&receipt.key).ok()
263 != Some(DecisionArchivePageKey {
264 bucket: self.bucket,
265 segment: self.segment,
266 })
267 || CompactDecisionArchiveReceipt::from_receipt(receipt).is_err()
268 })
269 {
270 return Err(archive_error(
271 "decision archive bucket page is malformed or non-canonical",
272 ));
273 }
274 let encoded = serde_json::to_vec(self).map_err(|error| {
275 archive_error(format!(
276 "cannot encode decision archive bucket page: {error}"
277 ))
278 })?;
279 if encoded.len() > MAX_DECISION_ARCHIVE_BUCKET_PAGE_BYTES {
280 return Err(archive_error(
281 "decision archive bucket page exceeds the hard byte limit",
282 ));
283 }
284 Ok(())
285 }
286
287 pub fn state_page_id(&self) -> Result<String, DecisionError> {
288 self.validate()?;
289 let bytes = serde_json::to_vec(self).map_err(|error| {
290 archive_error(format!(
291 "cannot encode decision archive bucket page: {error}"
292 ))
293 })?;
294 let mut hasher = blake3::Hasher::new();
295 hasher.update(b"canwu.state-page.v1");
296 hasher.update(&[0]);
297 hasher.update(&bytes);
298 Ok(hasher.finalize().to_hex().to_string())
299 }
300}
301
302#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
303#[serde(deny_unknown_fields)]
304pub struct DecisionLocatorScaleMetrics {
305 pub entries: u64,
306 pub locator_pages: u64,
307 pub max_page_entries: u64,
308 pub max_page_encoded_bytes: u64,
309 pub archive_batches: u64,
310 pub exact_restart_queries: u64,
311 pub reachable_blob_locators: u64,
312 pub estimated_resident_structural_bytes: u64,
313 pub root_hash: String,
314}
315
316#[doc(hidden)]
317pub struct DecisionLocatorScaleFixture {
318 pub state: DecisionState,
319 pub archive_blobs: Vec<DecisionArchiveBlob>,
320 pub metrics: DecisionLocatorScaleMetrics,
321}
322
323#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
324#[serde(deny_unknown_fields)]
325pub struct TraceLocatorScaleMetrics {
326 pub hot_trace_entries: u64,
327 pub indexed_lookup_samples: u64,
328 pub archive_commit_entries: u64,
329 pub target_archived: bool,
330}
331
332#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
333#[serde(tag = "record", rename_all = "snake_case")]
334pub enum DecisionArchiveRecord {
335 Ticket { ticket: DecisionTicket },
336 Attempt { attempt: DecisionAttemptRecord },
337 Trace { trace: DecisionTrace },
338}
339
340impl DecisionArchiveRecord {
341 #[must_use]
342 pub fn key(&self) -> DecisionHistoryKey {
343 match self {
344 Self::Ticket { ticket } => DecisionHistoryKey::Ticket(ticket.id),
345 Self::Attempt { attempt } => DecisionHistoryKey::Attempt(attempt.request_id),
346 Self::Trace { trace } => DecisionHistoryKey::Trace(trace.id),
347 }
348 }
349}
350
351#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
352#[serde(deny_unknown_fields)]
353pub struct DecisionArchiveBlob {
354 pub format_version: u32,
355 pub key: DecisionHistoryKey,
356 pub record: DecisionArchiveRecord,
357}
358
359impl DecisionArchiveBlob {
360 pub fn validate(&self) -> Result<(), DecisionError> {
361 if self.format_version != DECISION_ARCHIVE_FORMAT_VERSION || self.key != self.record.key() {
362 return Err(DecisionError::new(
363 DecisionErrorCode::InvalidDecision,
364 "decision archive blob format or typed identity is invalid",
365 ));
366 }
367 Ok(())
368 }
369
370 pub fn content_id(&self) -> Result<String, DecisionError> {
371 self.validate()?;
372 decision_archive_hash("canwu.decision.archive-blob.v1", self)
373 }
374}
375
376#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
377#[serde(deny_unknown_fields)]
378pub struct DecisionArchiveReceipt {
379 pub format_version: u32,
380 pub key: DecisionHistoryKey,
381 pub locator: String,
382 pub content_id: String,
383 pub encoded_bytes: u64,
384}
385
386#[derive(Clone, Debug, Eq, PartialEq)]
387struct CompactDecisionArchiveReceipt {
388 content_id: [u8; 32],
389 encoded_bytes: u64,
390}
391
392impl CompactDecisionArchiveReceipt {
393 fn from_receipt(receipt: &DecisionArchiveReceipt) -> Result<Self, DecisionError> {
394 if receipt.format_version != DECISION_ARCHIVE_FORMAT_VERSION
395 || receipt.locator != receipt.content_id
396 || receipt.encoded_bytes == 0
397 {
398 return Err(archive_error("decision archive receipt is not canonical"));
399 }
400 Ok(Self {
401 content_id: decode_archive_hash(&receipt.content_id)?,
402 encoded_bytes: receipt.encoded_bytes,
403 })
404 }
405
406 fn to_receipt(&self, key: &DecisionHistoryKey) -> DecisionArchiveReceipt {
407 let content_id = encode_archive_hash(self.content_id);
408 DecisionArchiveReceipt {
409 format_version: DECISION_ARCHIVE_FORMAT_VERSION,
410 key: key.clone(),
411 locator: content_id.clone(),
412 content_id,
413 encoded_bytes: self.encoded_bytes,
414 }
415 }
416}
417
418#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
419#[serde(deny_unknown_fields)]
420pub struct PreparedDecisionArchive {
421 pub source_root: String,
422 pub token: String,
423 pub blobs: Vec<DecisionArchiveBlob>,
424 pub receipts: Vec<DecisionArchiveReceipt>,
425}
426
427#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
431#[serde(deny_unknown_fields)]
432pub struct VerifiedDecisionArchiveCommit {
433 format_version: u32,
434 source_root: String,
435 token: String,
436 receipts: Vec<DecisionArchiveReceipt>,
437 page_replacements: Vec<VerifiedDecisionArchivePageReplacement>,
438}
439
440#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
441#[serde(deny_unknown_fields)]
442struct VerifiedDecisionArchivePageReplacement {
443 page_key: DecisionArchivePageKey,
444 previous_page_id: Option<String>,
445 page: DecisionArchiveBucketPage,
446}
447
448impl VerifiedDecisionArchiveCommit {
449 #[must_use]
450 pub fn token(&self) -> &str {
451 &self.token
452 }
453
454 #[must_use]
455 pub fn source_root(&self) -> &str {
456 &self.source_root
457 }
458
459 #[must_use]
460 pub fn has_current_nonempty_shape(&self) -> bool {
461 self.format_version == DECISION_ARCHIVE_FORMAT_VERSION
462 && !self.receipts.is_empty()
463 && !self.page_replacements.is_empty()
464 }
465
466 pub fn archive_locators(&self) -> impl Iterator<Item = &str> {
467 self.receipts.iter().map(|receipt| receipt.locator.as_str())
468 }
469}
470
471#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
472#[serde(rename_all = "snake_case")]
473pub enum DecisionArchiveStoreOutcome {
474 Stored,
475 AlreadyStored,
476}
477
478pub trait DecisionArchiveProvider {
479 fn load_decision_archive(
480 &self,
481 locator: &str,
482 ) -> Result<Option<DecisionArchiveBlob>, DecisionError>;
483
484 fn load_decision_archive_bucket_page(
488 &self,
489 _page_id: &str,
490 ) -> Result<Option<DecisionArchiveBucketPage>, DecisionError> {
491 Ok(None)
492 }
493}
494
495pub trait DecisionArchiveStore: DecisionArchiveProvider {
496 fn store_decision_archive(
497 &self,
498 blob: &DecisionArchiveBlob,
499 ) -> Result<DecisionArchiveStoreOutcome, DecisionError>;
500}
501
502fn archive_error(message: impl Into<String>) -> DecisionError {
503 DecisionError::new(DecisionErrorCode::InvalidDecision, message)
504}
505
506fn history_unavailable(message: impl Into<String>) -> DecisionError {
507 DecisionError::new(DecisionErrorCode::DecisionHistoryUnavailable, message)
508}
509
510fn history_budget_error(message: impl Into<String>) -> DecisionError {
511 DecisionError::new(DecisionErrorCode::QueryBudgetExceeded, message)
512}
513
514fn decision_archive_hash(domain: &str, value: &impl Serialize) -> Result<String, DecisionError> {
515 let encoded = serde_json::to_vec(&(domain, value)).map_err(|error| {
516 archive_error(format!(
517 "cannot encode canonical decision archive content: {error}"
518 ))
519 })?;
520 Ok(blake3::hash(&encoded).to_hex().to_string())
521}
522
523fn canonical_archive_hash(value: &str) -> bool {
524 value.len() == 64
525 && value
526 .bytes()
527 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
528}
529
530fn decode_archive_hash(value: &str) -> Result<[u8; 32], DecisionError> {
531 if !canonical_archive_hash(value) {
532 return Err(archive_error(
533 "decision archive hash must be lower-case 32-byte hexadecimal",
534 ));
535 }
536 let mut decoded = [0_u8; 32];
537 let (pairs, remainder) = value.as_bytes().as_chunks::<2>();
538 debug_assert!(remainder.is_empty());
539 for (index, pair) in pairs.iter().enumerate() {
540 let digit = |byte| match byte {
541 b'0'..=b'9' => Some(byte - b'0'),
542 b'a'..=b'f' => Some(byte - b'a' + 10),
543 _ => None,
544 };
545 decoded[index] = digit(pair[0])
546 .zip(digit(pair[1]))
547 .map(|(high, low)| (high << 4) | low)
548 .ok_or_else(|| archive_error("decision archive hash contains invalid hexadecimal"))?;
549 }
550 Ok(decoded)
551}
552
553fn encode_archive_hash(value: [u8; 32]) -> String {
554 let mut encoded = String::with_capacity(64);
555 for byte in value {
556 write!(&mut encoded, "{byte:02x}").expect("writing to a string cannot fail");
557 }
558 encoded
559}
560
561#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
562pub struct DecisionHotState {
563 pub ticket_count: u64,
564 pub attempt_count: u64,
565 pub trace_count: u64,
566}
567
568#[derive(Clone, Debug, Default, Eq, PartialEq)]
569struct DecisionHotHistoryAccumulator {
570 count: u64,
571 xor: [u8; 32],
572 sum: [u8; 32],
573}
574
575impl DecisionHotHistoryAccumulator {
576 fn insert(
577 &mut self,
578 key: &DecisionHistoryKey,
579 record: &DecisionArchiveRecord,
580 ) -> Result<(), DecisionError> {
581 let digest = decision_hot_leaf_hash(key, record)?;
582 self.count = self
583 .count
584 .checked_add(1)
585 .ok_or_else(|| archive_error("decision hot-history count overflowed"))?;
586 for (target, byte) in self.xor.iter_mut().zip(digest) {
587 *target ^= byte;
588 }
589 add_digest_mod_256(&mut self.sum, digest);
590 Ok(())
591 }
592
593 fn remove(
594 &mut self,
595 key: &DecisionHistoryKey,
596 record: &DecisionArchiveRecord,
597 ) -> Result<(), DecisionError> {
598 let digest = decision_hot_leaf_hash(key, record)?;
599 self.count = self
600 .count
601 .checked_sub(1)
602 .ok_or_else(|| archive_error("decision hot-history count underflowed"))?;
603 for (target, byte) in self.xor.iter_mut().zip(digest) {
604 *target ^= byte;
605 }
606 subtract_digest_mod_256(&mut self.sum, digest);
607 Ok(())
608 }
609}
610
611#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
617#[serde(bound(serialize = "T: Serialize", deserialize = "T: Deserialize<'de>"))]
618pub struct PersistentDecisionLog<T: Clone> {
619 next_ordinal: u64,
620 #[serde(with = "persistent_decision_map_serde")]
621 entries: OrdMap<u64, Arc<T>>,
622}
623
624impl<T: Clone> Default for PersistentDecisionLog<T> {
625 fn default() -> Self {
626 Self {
627 next_ordinal: 1,
628 entries: OrdMap::new(),
629 }
630 }
631}
632
633impl<T: Clone> PersistentDecisionLog<T> {
634 #[must_use]
635 pub fn is_empty(&self) -> bool {
636 self.entries.is_empty()
637 }
638
639 fn is_default_empty(&self) -> bool {
640 self.entries.is_empty() && self.next_ordinal == 1
641 }
642
643 #[must_use]
644 pub fn len(&self) -> usize {
645 self.entries.len()
646 }
647
648 #[must_use]
649 pub fn get(&self, index: usize) -> Option<&T> {
650 self.entries.values().nth(index).map(Arc::as_ref)
651 }
652
653 #[must_use]
654 pub fn last(&self) -> Option<&T> {
655 self.entries.get_max().map(|(_, value)| value.as_ref())
656 }
657
658 #[must_use]
659 pub fn iter(&self) -> impl DoubleEndedIterator<Item = &T> + ExactSizeIterator {
660 self.entries.values().map(Arc::as_ref)
661 }
662
663 pub fn push(&mut self, value: T) -> u64 {
669 let ordinal = self.next_ordinal;
670 self.next_ordinal = self
671 .next_ordinal
672 .checked_add(1)
673 .expect("decision log ordinal range is exhausted");
674 self.entries.insert(ordinal, Arc::new(value));
675 ordinal
676 }
677
678 #[must_use]
679 pub fn get_ordinal(&self, ordinal: u64) -> Option<&T> {
680 self.entries.get(&ordinal).map(Arc::as_ref)
681 }
682
683 pub fn remove_ordinal(&mut self, ordinal: u64) -> Option<T> {
684 self.entries.remove(&ordinal).map(Arc::unwrap_or_clone)
685 }
686
687 #[must_use]
688 pub fn ordinals(&self) -> impl DoubleEndedIterator<Item = u64> + ExactSizeIterator + '_ {
689 self.entries.keys().copied()
690 }
691
692 #[must_use]
693 pub const fn next_ordinal(&self) -> u64 {
694 self.next_ordinal
695 }
696}
697
698impl<T: Clone> Index<usize> for PersistentDecisionLog<T> {
699 type Output = T;
700
701 fn index(&self, index: usize) -> &Self::Output {
702 self.get(index)
703 .expect("persistent decision log index is out of bounds")
704 }
705}
706
707pub struct PersistentDecisionLogIter<'a, T: Clone> {
708 inner: im::ordmap::Values<'a, u64, Arc<T>>,
709}
710
711impl<'a, T: Clone> Iterator for PersistentDecisionLogIter<'a, T> {
712 type Item = &'a T;
713
714 fn next(&mut self) -> Option<Self::Item> {
715 self.inner.next().map(Arc::as_ref)
716 }
717
718 fn size_hint(&self) -> (usize, Option<usize>) {
719 self.inner.size_hint()
720 }
721}
722
723impl<T: Clone> DoubleEndedIterator for PersistentDecisionLogIter<'_, T> {
724 fn next_back(&mut self) -> Option<Self::Item> {
725 self.inner.next_back().map(Arc::as_ref)
726 }
727}
728
729impl<T: Clone> ExactSizeIterator for PersistentDecisionLogIter<'_, T> {}
730
731impl<'a, T: Clone> IntoIterator for &'a PersistentDecisionLog<T> {
732 type Item = &'a T;
733 type IntoIter = PersistentDecisionLogIter<'a, T>;
734
735 fn into_iter(self) -> Self::IntoIter {
736 PersistentDecisionLogIter {
737 inner: self.entries.values(),
738 }
739 }
740}
741
742#[derive(Clone, Debug, Default, Eq, PartialEq, Serialize)]
743pub struct DecisionState {
744 #[serde(
745 default,
746 skip_serializing_if = "OrdMap::is_empty",
747 with = "persistent_decision_map_serde"
748 )]
749 pub controllers: OrdMap<String, Arc<DecisionControllerBinding>>,
750 #[serde(
751 default,
752 skip_serializing_if = "OrdMap::is_empty",
753 with = "persistent_decision_map_serde"
754 )]
755 pub tickets: OrdMap<DecisionTicketId, Arc<DecisionTicket>>,
756 #[serde(
757 default,
758 skip_serializing_if = "PersistentDecisionLog::is_default_empty"
759 )]
760 pub traces: PersistentDecisionLog<DecisionTrace>,
761 #[serde(
762 default,
763 skip_serializing_if = "PersistentDecisionLog::is_default_empty"
764 )]
765 attempts: PersistentDecisionLog<DecisionAttemptRecord>,
766 #[serde(skip)]
767 attempts_by_request: OrdMap<canwu_core::DecisionRequestId, u64>,
768 #[serde(skip)]
769 deadline_index: OrdMap<SimTime, OrdSet<DecisionTicketId>>,
770 #[serde(
771 default,
772 skip_serializing_if = "OrdMap::is_empty",
773 with = "decision_archive_receipts_serde"
774 )]
775 archive_receipt_buckets:
776 OrdMap<DecisionArchivePageKey, OrdMap<DecisionHistoryKey, CompactDecisionArchiveReceipt>>,
777 #[serde(
778 default,
779 skip_serializing_if = "OrdMap::is_empty",
780 with = "decision_archive_page_directory_serde"
781 )]
782 archive_bucket_page_ids: OrdMap<DecisionArchivePageKey, String>,
783 #[serde(default, skip_serializing_if = "is_zero_u64")]
784 archive_receipt_count: u64,
785 #[serde(skip)]
786 hot_history_accumulator: DecisionHotHistoryAccumulator,
787}
788
789#[allow(clippy::trivially_copy_pass_by_ref)]
790const fn is_zero_u64(value: &u64) -> bool {
791 *value == 0
792}
793
794impl<'de> Deserialize<'de> for DecisionState {
795 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
796 where
797 D: serde::Deserializer<'de>,
798 {
799 #[derive(Deserialize)]
800 struct PersistedDecisionState {
801 #[serde(default, with = "persistent_decision_map_serde")]
802 controllers: OrdMap<String, Arc<DecisionControllerBinding>>,
803 #[serde(default, with = "persistent_decision_map_serde")]
804 tickets: OrdMap<canwu_core::DecisionTicketId, Arc<DecisionTicket>>,
805 #[serde(default)]
806 traces: PersistentDecisionLog<DecisionTrace>,
807 #[serde(default)]
808 attempts: PersistentDecisionLog<DecisionAttemptRecord>,
809 #[serde(default, with = "decision_archive_receipts_serde")]
810 archive_receipt_buckets: OrdMap<
811 DecisionArchivePageKey,
812 OrdMap<DecisionHistoryKey, CompactDecisionArchiveReceipt>,
813 >,
814 #[serde(default, with = "decision_archive_page_directory_serde")]
815 archive_bucket_page_ids: OrdMap<DecisionArchivePageKey, String>,
816 #[serde(default)]
817 archive_receipt_count: u64,
818 }
819
820 let persisted = PersistedDecisionState::deserialize(deserializer)?;
821 let attempts_by_request = persisted
822 .attempts
823 .ordinals()
824 .zip(persisted.attempts.iter())
825 .map(|(ordinal, attempt)| (attempt.request_id, ordinal))
826 .collect();
827 let mut deadline_index = OrdMap::<SimTime, OrdSet<DecisionTicketId>>::new();
828 for ticket in persisted.tickets.values().filter(|ticket| ticket.is_open()) {
829 if let Some(deadline) = ticket.deadline {
830 deadline_index
831 .entry(deadline)
832 .or_default()
833 .insert(ticket.id);
834 }
835 }
836 Self {
837 controllers: persisted.controllers,
838 tickets: persisted.tickets,
839 traces: persisted.traces,
840 attempts: persisted.attempts,
841 attempts_by_request,
842 deadline_index,
843 archive_receipt_buckets: persisted.archive_receipt_buckets,
844 archive_bucket_page_ids: persisted.archive_bucket_page_ids,
845 archive_receipt_count: persisted.archive_receipt_count,
846 hot_history_accumulator: DecisionHotHistoryAccumulator::default(),
847 }
848 .with_rebuilt_hot_history_accumulator()
849 .map_err(serde::de::Error::custom)?
850 .with_rebuilt_archive_bucket_pages()
851 .map_err(serde::de::Error::custom)
852 }
853}
854
855impl DecisionState {
856 fn with_rebuilt_hot_history_accumulator(mut self) -> Result<Self, DecisionError> {
857 self.rebuild_hot_history_accumulator()?;
858 Ok(self)
859 }
860
861 fn rebuild_hot_history_accumulator(&mut self) -> Result<(), DecisionError> {
862 let mut rebuilt = DecisionHotHistoryAccumulator::default();
863 for ticket in self.tickets.values() {
864 let key = DecisionHistoryKey::Ticket(ticket.id);
865 rebuilt.insert(
866 &key,
867 &DecisionArchiveRecord::Ticket {
868 ticket: ticket.as_ref().clone(),
869 },
870 )?;
871 }
872 for attempt in self.attempts.iter() {
873 let key = DecisionHistoryKey::Attempt(attempt.request_id);
874 rebuilt.insert(
875 &key,
876 &DecisionArchiveRecord::Attempt {
877 attempt: attempt.clone(),
878 },
879 )?;
880 }
881 for trace in self.traces.iter() {
882 let key = DecisionHistoryKey::Trace(trace.id);
883 rebuilt.insert(
884 &key,
885 &DecisionArchiveRecord::Trace {
886 trace: trace.clone(),
887 },
888 )?;
889 }
890 self.hot_history_accumulator = rebuilt;
891 Ok(())
892 }
893
894 fn insert_hot_history_record(
895 &mut self,
896 record: &DecisionArchiveRecord,
897 ) -> Result<(), DecisionError> {
898 self.hot_history_accumulator.insert(&record.key(), record)
899 }
900
901 fn replace_hot_history_record(
902 &mut self,
903 previous: &DecisionArchiveRecord,
904 next: &DecisionArchiveRecord,
905 ) -> Result<(), DecisionError> {
906 if previous.key() != next.key() {
907 return Err(archive_error(
908 "decision hot-history replacement changed record identity",
909 ));
910 }
911 self.hot_history_accumulator
912 .remove(&previous.key(), previous)?;
913 self.hot_history_accumulator.insert(&next.key(), next)
914 }
915
916 fn remove_hot_history_record(
917 &mut self,
918 record: &DecisionArchiveRecord,
919 ) -> Result<(), DecisionError> {
920 self.hot_history_accumulator.remove(&record.key(), record)
921 }
922
923 #[must_use]
924 pub fn is_empty(&self) -> bool {
925 self.controllers.is_empty()
926 && self.tickets.is_empty()
927 && self.traces.is_empty()
928 && self.attempts.is_empty()
929 && self.archive_receipt_count == 0
930 }
931
932 fn compact_archive_receipt(
933 &self,
934 key: &DecisionHistoryKey,
935 ) -> Option<&CompactDecisionArchiveReceipt> {
936 let bucket = decision_history_page_key(key).ok()?;
937 self.archive_receipt_buckets.get(&bucket)?.get(key)
938 }
939
940 fn contains_archived_key(&self, key: &DecisionHistoryKey) -> bool {
941 self.compact_archive_receipt(key).is_some()
942 }
943
944 fn insert_archive_receipt(
945 &mut self,
946 receipt: &DecisionArchiveReceipt,
947 ) -> Result<Option<CompactDecisionArchiveReceipt>, DecisionError> {
948 let bucket = decision_history_page_key(&receipt.key)?;
949 let previous = self
950 .archive_receipt_buckets
951 .entry(bucket)
952 .or_default()
953 .insert(
954 receipt.key.clone(),
955 CompactDecisionArchiveReceipt::from_receipt(receipt)?,
956 );
957 self.refresh_archive_bucket_page_id(bucket)?;
958 if previous.is_none() {
959 self.archive_receipt_count = self
960 .archive_receipt_count
961 .checked_add(1)
962 .ok_or_else(|| archive_error("decision archive receipt count overflowed"))?;
963 }
964 Ok(previous)
965 }
966
967 fn with_rebuilt_archive_bucket_pages(mut self) -> Result<Self, DecisionError> {
968 let resident_count =
969 self.archive_receipt_buckets
970 .values()
971 .try_fold(0_u64, |total, receipts| {
972 total
973 .checked_add(receipts.len() as u64)
974 .ok_or_else(|| archive_error("decision archive receipt count overflowed"))
975 })?;
976 if resident_count > 0 {
977 let buckets = self
978 .archive_receipt_buckets
979 .keys()
980 .copied()
981 .collect::<Vec<_>>();
982 if self.archive_bucket_page_ids.is_empty() {
983 self.archive_receipt_count = resident_count;
984 for bucket in buckets {
985 self.refresh_archive_bucket_page_id(bucket)?;
986 }
987 } else {
988 for bucket in buckets {
989 let resident_page_id = self
990 .decision_archive_bucket_page(bucket)?
991 .ok_or_else(|| archive_error("resident decision archive bucket vanished"))?
992 .state_page_id()?;
993 if self.archive_bucket_page_ids.get(&bucket) != Some(&resident_page_id) {
994 return Err(archive_error(
995 "decision archive bucket directory disagrees with resident receipts",
996 ));
997 }
998 }
999 }
1000 }
1001 self.validate_archive_directory_shape()?;
1002 Ok(self)
1003 }
1004
1005 fn refresh_archive_bucket_page_id(
1006 &mut self,
1007 bucket: DecisionArchivePageKey,
1008 ) -> Result<(), DecisionError> {
1009 let page = self
1010 .decision_archive_bucket_page(bucket)?
1011 .ok_or_else(|| archive_error("decision archive bucket unexpectedly disappeared"))?;
1012 self.archive_bucket_page_ids
1013 .insert(bucket, page.state_page_id()?);
1014 Ok(())
1015 }
1016
1017 pub fn decision_archive_bucket_page(
1018 &self,
1019 bucket: DecisionArchivePageKey,
1020 ) -> Result<Option<DecisionArchiveBucketPage>, DecisionError> {
1021 let Some(receipts) = self.archive_receipt_buckets.get(&bucket) else {
1022 return Ok(None);
1023 };
1024 let page = DecisionArchiveBucketPage {
1025 format_version: DECISION_ARCHIVE_BUCKET_PAGE_FORMAT_VERSION,
1026 bucket: bucket.bucket,
1027 segment: bucket.segment,
1028 receipts: receipts
1029 .iter()
1030 .map(|(key, receipt)| receipt.to_receipt(key))
1031 .collect(),
1032 };
1033 page.validate()?;
1034 Ok(Some(page))
1035 }
1036
1037 #[must_use]
1038 pub fn decision_archive_bucket_page_ids(&self) -> &OrdMap<DecisionArchivePageKey, String> {
1039 &self.archive_bucket_page_ids
1040 }
1041
1042 #[must_use]
1046 pub fn paged_checkpoint_hot_state(&self) -> Self {
1047 let mut hot = self.clone();
1048 hot.archive_receipt_buckets.clear();
1049 hot.archive_bucket_page_ids.clear();
1050 hot.archive_receipt_count = 0;
1051 hot
1052 }
1053
1054 #[must_use]
1055 pub fn archived_history_count(&self) -> usize {
1056 usize::try_from(self.archive_receipt_count).unwrap_or(usize::MAX)
1057 }
1058
1059 #[must_use]
1060 pub fn controller(&self, id: &str) -> Option<&DecisionControllerBinding> {
1061 self.controllers.get(id).map(Arc::as_ref)
1062 }
1063
1064 #[must_use]
1065 pub fn ticket(&self, id: DecisionTicketId) -> Option<&DecisionTicket> {
1066 self.tickets.get(&id).map(Arc::as_ref)
1067 }
1068
1069 #[must_use]
1070 pub fn trace(&self, id: DecisionTraceId) -> Option<&DecisionTrace> {
1071 self.traces.get_ordinal(id.get())
1072 }
1073
1074 #[must_use]
1075 pub fn attempt(&self, id: canwu_core::DecisionRequestId) -> Option<&DecisionAttemptRecord> {
1076 self.attempts_by_request
1077 .get(&id)
1078 .and_then(|ordinal| self.attempts.get_ordinal(*ordinal))
1079 }
1080
1081 #[must_use]
1082 pub const fn attempts(&self) -> &PersistentDecisionLog<DecisionAttemptRecord> {
1083 &self.attempts
1084 }
1085
1086 #[must_use]
1089 pub fn decision_locator(&self, key: &DecisionHistoryKey) -> DecisionHistoryLocation {
1090 let hot = match key {
1091 DecisionHistoryKey::Ticket(id) => self.tickets.contains_key(id),
1092 DecisionHistoryKey::Attempt(id) => self.attempts_by_request.contains_key(id),
1093 DecisionHistoryKey::Trace(id) => self.traces.get_ordinal(id.get()).is_some(),
1094 };
1095 if hot {
1096 DecisionHistoryLocation::Hot
1097 } else if let Some(receipt) = self.compact_archive_receipt(key) {
1098 DecisionHistoryLocation::Archived {
1099 locator: encode_archive_hash(receipt.content_id),
1100 }
1101 } else if let Ok(page) = decision_history_page_key(key)
1102 && self.archive_bucket_page_ids.contains_key(&page)
1103 {
1104 DecisionHistoryLocation::Unresolved {
1105 bucket: page.bucket,
1106 segment: page.segment,
1107 }
1108 } else {
1109 DecisionHistoryLocation::Absent
1110 }
1111 }
1112
1113 pub fn decision_locator_with_provider(
1117 &self,
1118 key: &DecisionHistoryKey,
1119 provider: &dyn DecisionArchiveProvider,
1120 ) -> Result<DecisionHistoryLocation, DecisionError> {
1121 let local = self.decision_locator(key);
1122 match local {
1123 DecisionHistoryLocation::Hot
1124 | DecisionHistoryLocation::Archived { .. }
1125 | DecisionHistoryLocation::Absent => return Ok(local),
1126 DecisionHistoryLocation::Unresolved { .. } => {}
1127 }
1128 let bucket = decision_history_page_key(key)?;
1129 let Some(page_id) = self.archive_bucket_page_ids.get(&bucket) else {
1130 return Ok(DecisionHistoryLocation::Absent);
1131 };
1132 let page = provider
1133 .load_decision_archive_bucket_page(page_id)?
1134 .ok_or_else(|| history_unavailable("decision locator bucket page is unavailable"))?;
1135 page.validate()?;
1136 if page.bucket != bucket.bucket
1137 || page.segment != bucket.segment
1138 || page.state_page_id()? != *page_id
1139 {
1140 return Err(history_unavailable(
1141 "decision locator provider returned a mismatched bucket page",
1142 ));
1143 }
1144 match page
1145 .receipts
1146 .binary_search_by(|receipt| receipt.key.cmp(key))
1147 {
1148 Ok(index) => Ok(DecisionHistoryLocation::Archived {
1149 locator: page.receipts[index].locator.clone(),
1150 }),
1151 Err(_) => Ok(DecisionHistoryLocation::Absent),
1152 }
1153 }
1154
1155 #[must_use]
1156 pub fn archive_receipt(&self, key: &DecisionHistoryKey) -> Option<DecisionArchiveReceipt> {
1157 self.compact_archive_receipt(key)
1158 .map(|receipt| receipt.to_receipt(key))
1159 }
1160
1161 fn archive_receipt_with_provider(
1162 &self,
1163 key: &DecisionHistoryKey,
1164 provider: &dyn DecisionArchiveProvider,
1165 ) -> Result<Option<DecisionArchiveReceipt>, DecisionError> {
1166 if let Some(receipt) = self.archive_receipt(key) {
1167 return Ok(Some(receipt));
1168 }
1169 let bucket = decision_history_page_key(key)?;
1170 let Some(page_id) = self.archive_bucket_page_ids.get(&bucket) else {
1171 return Ok(None);
1172 };
1173 let page = provider
1174 .load_decision_archive_bucket_page(page_id)?
1175 .ok_or_else(|| history_unavailable("decision locator bucket page is unavailable"))?;
1176 page.validate()?;
1177 if page.bucket != bucket.bucket
1178 || page.segment != bucket.segment
1179 || page.state_page_id()? != *page_id
1180 {
1181 return Err(history_unavailable(
1182 "decision locator provider returned a mismatched bucket page",
1183 ));
1184 }
1185 Ok(page
1186 .receipts
1187 .binary_search_by(|receipt| receipt.key.cmp(key))
1188 .ok()
1189 .map(|index| page.receipts[index].clone()))
1190 }
1191
1192 pub fn archived_history_keys(&self) -> impl Iterator<Item = &DecisionHistoryKey> {
1193 self.archive_receipt_buckets.values().flat_map(OrdMap::keys)
1194 }
1195
1196 pub fn required_archived_dependency_page_keys(
1201 &self,
1202 ) -> Result<OrdSet<DecisionArchivePageKey>, DecisionError> {
1203 let mut pages = OrdSet::new();
1204 for trace in self.traces.iter() {
1205 if !self.tickets.contains_key(&trace.ticket_id) {
1206 pages.insert(decision_history_page_key(&DecisionHistoryKey::Ticket(
1207 trace.ticket_id,
1208 ))?);
1209 }
1210 }
1211 for ticket in self.tickets.values() {
1212 if let DecisionTicketState::Resolved { trace_id, .. } = ticket.state
1213 && self.traces.get_ordinal(trace_id.get()).is_none()
1214 {
1215 pages.insert(decision_history_page_key(&DecisionHistoryKey::Trace(
1216 trace_id,
1217 ))?);
1218 }
1219 if let Some(parent_id) = ticket.parent_ticket
1220 && !self.tickets.contains_key(&parent_id)
1221 {
1222 pages.insert(decision_history_page_key(&DecisionHistoryKey::Ticket(
1223 parent_id,
1224 ))?);
1225 }
1226 }
1227 for attempt in self.attempts.iter() {
1228 if let DecisionAttemptOutcome::Accepted {
1229 trace_id: Some(trace_id),
1230 ..
1231 } = &attempt.outcome
1232 && self.traces.get_ordinal(trace_id.get()).is_none()
1233 {
1234 pages.insert(decision_history_page_key(&DecisionHistoryKey::Trace(
1235 *trace_id,
1236 ))?);
1237 }
1238 }
1239 Ok(pages)
1240 }
1241
1242 pub fn paged_checkpoint_parts(
1246 &self,
1247 ) -> Result<
1248 (
1249 Self,
1250 std::collections::BTreeMap<DecisionArchivePageKey, Vec<DecisionArchiveReceipt>>,
1251 ),
1252 DecisionError,
1253 > {
1254 let hot = self.paged_checkpoint_hot_state();
1255 let mut buckets =
1256 std::collections::BTreeMap::<DecisionArchivePageKey, Vec<DecisionArchiveReceipt>>::new(
1257 );
1258 for (bucket, receipts) in &self.archive_receipt_buckets {
1259 let decoded = receipts
1260 .iter()
1261 .map(|(key, receipt)| receipt.to_receipt(key))
1262 .collect();
1263 buckets.insert(*bucket, decoded);
1264 }
1265 Ok((hot, buckets))
1266 }
1267
1268 pub fn from_paged_checkpoint_parts(
1269 mut hot: Self,
1270 buckets: impl IntoIterator<Item = (DecisionArchivePageKey, Vec<DecisionArchiveReceipt>)>,
1271 ) -> Result<Self, DecisionError> {
1272 if !hot.archive_receipt_buckets.is_empty() {
1273 return Err(archive_error(
1274 "paged decision hot state must not duplicate archive receipts",
1275 ));
1276 }
1277 hot.archive_bucket_page_ids.clear();
1278 hot.archive_receipt_count = 0;
1279 for (bucket, receipts) in buckets {
1280 if bucket.bucket >= DECISION_HISTORY_BUCKET_COUNT
1281 || receipts.windows(2).any(|pair| pair[0].key >= pair[1].key)
1282 {
1283 return Err(archive_error(
1284 "paged decision archive bucket is out of range or not strictly ordered",
1285 ));
1286 }
1287 for receipt in receipts {
1288 if decision_history_page_key(&receipt.key)? != bucket
1289 || hot.insert_archive_receipt(&receipt)?.is_some()
1290 {
1291 return Err(archive_error(
1292 "paged decision archive bucket contains a misplaced or duplicate receipt",
1293 ));
1294 }
1295 }
1296 }
1297 hot.validate()?;
1298 Ok(hot)
1299 }
1300
1301 pub fn from_paged_checkpoint_root(
1304 hot: Self,
1305 archive_bucket_page_ids: OrdMap<DecisionArchivePageKey, String>,
1306 archive_receipt_count: u64,
1307 archive_receipt_root: &str,
1308 ) -> Result<Self, DecisionError> {
1309 Self::from_paged_checkpoint_root_with_resident_pages(
1310 hot,
1311 archive_bucket_page_ids,
1312 archive_receipt_count,
1313 archive_receipt_root,
1314 std::iter::empty(),
1315 )
1316 }
1317
1318 pub fn from_paged_checkpoint_root_with_resident_pages(
1323 mut hot: Self,
1324 archive_bucket_page_ids: OrdMap<DecisionArchivePageKey, String>,
1325 archive_receipt_count: u64,
1326 archive_receipt_root: &str,
1327 resident_pages: impl IntoIterator<Item = DecisionArchiveBucketPage>,
1328 ) -> Result<Self, DecisionError> {
1329 if !hot.archive_receipt_buckets.is_empty()
1330 || !hot.archive_bucket_page_ids.is_empty()
1331 || hot.archive_receipt_count != 0
1332 {
1333 return Err(archive_error(
1334 "paged decision hot state must not contain archive receipts or archive metadata",
1335 ));
1336 }
1337 hot.archive_bucket_page_ids = archive_bucket_page_ids;
1338 hot.archive_receipt_count = archive_receipt_count;
1339 for page in resident_pages {
1340 page.validate()?;
1341 let page_key = DecisionArchivePageKey {
1342 bucket: page.bucket,
1343 segment: page.segment,
1344 };
1345 let page_id = page.state_page_id()?;
1346 if hot.archive_bucket_page_ids.get(&page_key) != Some(&page_id)
1347 || hot.archive_receipt_buckets.contains_key(&page_key)
1348 {
1349 return Err(archive_error(
1350 "resident decision archive page is absent from or disagrees with the committed directory",
1351 ));
1352 }
1353 let mut receipts = OrdMap::new();
1354 for receipt in page.receipts {
1355 if receipts
1356 .insert(
1357 receipt.key.clone(),
1358 CompactDecisionArchiveReceipt::from_receipt(&receipt)?,
1359 )
1360 .is_some()
1361 {
1362 return Err(archive_error(
1363 "resident decision archive page contains duplicate receipts",
1364 ));
1365 }
1366 }
1367 hot.archive_receipt_buckets.insert(page_key, receipts);
1368 }
1369 hot.validate()?;
1370 if hot.archive_receipt_root()? != archive_receipt_root {
1371 return Err(archive_error(
1372 "paged decision archive directory root is inconsistent",
1373 ));
1374 }
1375 Ok(hot)
1376 }
1377
1378 pub fn archive_receipt_commitment(&self) -> Result<String, DecisionError> {
1379 self.archive_receipt_root()
1380 }
1381
1382 pub fn authoritative_commitment(&self) -> Result<String, DecisionError> {
1385 decision_archive_hash(
1386 "canwu.commitment.decisions.v2",
1387 &(
1388 &self.controllers,
1389 self.hot_history_root()?,
1390 self.archive_receipt_root()?,
1391 ),
1392 )
1393 }
1394
1395 pub fn archived_decision_history_page(
1396 &self,
1397 cursor: Option<&DecisionHistoryCursor>,
1398 budget: DecisionHistoryQueryBudget,
1399 provider: &dyn DecisionArchiveProvider,
1400 ) -> Result<DecisionHistoryPage, DecisionError> {
1401 if budget.max_results == 0
1402 || budget.max_results > MAX_DECISION_HISTORY_PAGE_SIZE
1403 || budget.max_provider_calls == 0
1404 || budget.max_decoded_bytes == 0
1405 || budget.max_decoded_bytes > MAX_DECISION_HISTORY_PAGE_BYTES
1406 {
1407 return Err(history_budget_error(
1408 "decision history page budget is zero, inconsistent, or exceeds the hard limit",
1409 ));
1410 }
1411 let archive_root = self.archive_receipt_root()?;
1412 if cursor.is_some_and(|cursor| cursor.archive_root != archive_root) {
1413 return Err(history_unavailable(
1414 "decision history cursor belongs to a different archive generation",
1415 ));
1416 }
1417 if cursor.is_some_and(|cursor| {
1418 cursor.bucket >= DECISION_HISTORY_BUCKET_COUNT
1419 || decision_history_page_key(&cursor.after).ok()
1420 != Some(DecisionArchivePageKey {
1421 bucket: cursor.bucket,
1422 segment: cursor.segment,
1423 })
1424 }) {
1425 return Err(history_unavailable(
1426 "decision history cursor has an invalid bucket position",
1427 ));
1428 }
1429 let start_bucket = cursor.map_or(
1430 DecisionArchivePageKey {
1431 bucket: 0,
1432 segment: 0,
1433 },
1434 |cursor| DecisionArchivePageKey {
1435 bucket: cursor.bucket,
1436 segment: cursor.segment,
1437 },
1438 );
1439 let mut selected = Vec::<(DecisionArchivePageKey, DecisionArchiveReceipt)>::new();
1440 let mut provider_calls = 0_usize;
1441 for (bucket, page_id) in self.archive_bucket_page_ids.range(start_bucket..) {
1442 let loaded;
1443 let receipts = if let Some(receipts) = self.archive_receipt_buckets.get(bucket) {
1444 receipts
1445 .iter()
1446 .map(|(key, receipt)| receipt.to_receipt(key))
1447 .collect::<Vec<_>>()
1448 } else {
1449 provider_calls = provider_calls.checked_add(1).ok_or_else(|| {
1450 history_budget_error("decision history provider-call count overflowed")
1451 })?;
1452 if provider_calls > budget.max_provider_calls {
1453 return Err(history_budget_error(
1454 "decision history page exceeded its provider-call budget",
1455 ));
1456 }
1457 loaded = provider
1458 .load_decision_archive_bucket_page(page_id)?
1459 .ok_or_else(|| {
1460 history_unavailable("decision locator bucket page is unavailable")
1461 })?;
1462 loaded.validate()?;
1463 if loaded.bucket != bucket.bucket
1464 || loaded.segment != bucket.segment
1465 || loaded.state_page_id()? != *page_id
1466 {
1467 return Err(history_unavailable(
1468 "decision locator provider returned a mismatched bucket page",
1469 ));
1470 }
1471 loaded.receipts
1472 };
1473 if let Some(cursor) = cursor
1474 .filter(|cursor| cursor.bucket == bucket.bucket && cursor.segment == bucket.segment)
1475 {
1476 let after = &cursor.after;
1477 for receipt in receipts.iter().filter(|receipt| receipt.key > *after) {
1478 selected.push((*bucket, receipt.clone()));
1479 if selected.len() > budget.max_results {
1480 break;
1481 }
1482 }
1483 } else {
1484 for receipt in receipts {
1485 selected.push((*bucket, receipt));
1486 if selected.len() > budget.max_results {
1487 break;
1488 }
1489 }
1490 }
1491 if selected.len() > budget.max_results {
1492 break;
1493 }
1494 }
1495 let has_more = selected.len() > budget.max_results;
1496 let selected = &selected[..selected.len().min(budget.max_results)];
1497 let mut records = Vec::with_capacity(selected.len());
1498 let mut decoded_bytes = 0_u64;
1499 for (_, receipt) in selected {
1500 provider_calls = provider_calls.checked_add(1).ok_or_else(|| {
1501 history_budget_error("decision history provider-call count overflowed")
1502 })?;
1503 if provider_calls > budget.max_provider_calls {
1504 return Err(history_budget_error(
1505 "decision history page exceeded its provider-call budget",
1506 ));
1507 }
1508 let blob = provider
1509 .load_decision_archive(&receipt.locator)?
1510 .ok_or_else(|| {
1511 history_unavailable("decision archive provider omitted a committed page member")
1512 })?;
1513 blob.validate()?;
1514 if blob.key != receipt.key || blob.content_id()? != receipt.content_id {
1515 return Err(history_unavailable(
1516 "decision archive provider returned a mismatched page member",
1517 ));
1518 }
1519 decoded_bytes = decoded_bytes
1520 .checked_add(receipt.encoded_bytes)
1521 .ok_or_else(|| history_budget_error("decision history byte count overflowed"))?;
1522 if decoded_bytes > budget.max_decoded_bytes {
1523 return Err(history_budget_error(
1524 "decision history page exceeded its decoded-byte budget",
1525 ));
1526 }
1527 records.push(blob.record);
1528 }
1529 let next_cursor = if has_more {
1530 selected
1531 .last()
1532 .map(|(bucket, receipt)| DecisionHistoryCursor {
1533 archive_root: archive_root.clone(),
1534 bucket: bucket.bucket,
1535 segment: bucket.segment,
1536 after: receipt.key.clone(),
1537 })
1538 } else {
1539 None
1540 };
1541 Ok(DecisionHistoryPage {
1542 archive_root,
1543 records,
1544 next_cursor,
1545 provider_calls: provider_calls as u64,
1546 decoded_bytes,
1547 })
1548 }
1549
1550 pub fn archive_reachability(
1554 &self,
1555 provider: &dyn DecisionArchiveProvider,
1556 ) -> Result<DecisionArchiveReachability, DecisionError> {
1557 let mut reachable = DecisionArchiveReachability::default();
1558 let mut receipt_count = 0_u64;
1559 for (bucket, page_id) in &self.archive_bucket_page_ids {
1560 let page = provider
1561 .load_decision_archive_bucket_page(page_id)?
1562 .ok_or_else(|| {
1563 history_unavailable("decision archive bucket page is unavailable")
1564 })?;
1565 page.validate()?;
1566 if page.bucket != bucket.bucket
1567 || page.segment != bucket.segment
1568 || page.state_page_id()? != *page_id
1569 {
1570 return Err(history_unavailable(
1571 "decision archive reachability provider returned a mismatched page",
1572 ));
1573 }
1574 reachable.bucket_page_ids.insert(page_id.clone());
1575 for receipt in page.receipts {
1576 receipt_count = receipt_count
1577 .checked_add(1)
1578 .ok_or_else(|| archive_error("decision archive receipt count overflowed"))?;
1579 reachable.blob_locators.insert(receipt.locator);
1580 }
1581 }
1582 if receipt_count != self.archive_receipt_count {
1583 return Err(archive_error(
1584 "decision archive reachability count disagrees with its directory",
1585 ));
1586 }
1587 Ok(reachable)
1588 }
1589
1590 pub fn load_decision_history(
1591 &self,
1592 key: &DecisionHistoryKey,
1593 provider: &dyn DecisionArchiveProvider,
1594 ) -> Result<Option<DecisionArchiveRecord>, DecisionError> {
1595 if let Some(record) = self.hot_archive_record(key) {
1596 return Ok(Some(record));
1597 }
1598 let Some(receipt) = self.archive_receipt_with_provider(key, provider)? else {
1599 return Ok(None);
1600 };
1601 let blob = provider
1602 .load_decision_archive(&receipt.locator)?
1603 .ok_or_else(|| archive_error("decision archive provider omitted a committed blob"))?;
1604 blob.validate()?;
1605 let content_id = blob.content_id()?;
1606 if blob.key != *key || content_id != receipt.content_id {
1607 return Err(archive_error(
1608 "decision archive provider returned content that disagrees with its receipt",
1609 ));
1610 }
1611 Ok(Some(blob.record))
1612 }
1613
1614 pub fn prepare_decision_archive(
1615 &self,
1616 keys: &[DecisionHistoryKey],
1617 ) -> Result<PreparedDecisionArchive, DecisionError> {
1618 if keys.is_empty() || keys.len() > MAX_DECISION_ARCHIVE_BATCH_ENTRIES {
1619 return Err(archive_error(
1620 "decision archive batch is empty or exceeds the bounded entry limit",
1621 ));
1622 }
1623 let mut canonical_keys = keys.to_vec();
1624 canonical_keys.sort();
1625 canonical_keys.dedup();
1626 if canonical_keys.len() != keys.len() {
1627 return Err(archive_error(
1628 "decision archive batch contains duplicate keys",
1629 ));
1630 }
1631 let source_root = self.hot_history_root()?;
1632 let mut blobs = Vec::with_capacity(canonical_keys.len());
1633 let mut receipts = Vec::with_capacity(canonical_keys.len());
1634 for key in canonical_keys {
1635 if self.contains_archived_key(&key) {
1636 return Err(archive_error("decision history is already archived"));
1637 }
1638 let record = self
1639 .hot_archive_record(&key)
1640 .ok_or_else(|| archive_error("decision archive key is not resident"))?;
1641 self.validate_archive_eligibility(&record)?;
1642 let blob = DecisionArchiveBlob {
1643 format_version: DECISION_ARCHIVE_FORMAT_VERSION,
1644 key: key.clone(),
1645 record,
1646 };
1647 let content_id = blob.content_id()?;
1648 let encoded_bytes = u64::try_from(
1649 serde_json::to_vec(&blob)
1650 .map_err(|error| {
1651 archive_error(format!("cannot encode decision archive: {error}"))
1652 })?
1653 .len(),
1654 )
1655 .map_err(|_| {
1656 archive_error("decision archive blob exceeds the persistent byte range")
1657 })?;
1658 receipts.push(DecisionArchiveReceipt {
1659 format_version: DECISION_ARCHIVE_FORMAT_VERSION,
1660 key,
1661 locator: content_id.clone(),
1662 content_id,
1663 encoded_bytes,
1664 });
1665 blobs.push(blob);
1666 }
1667 let token = decision_archive_hash(
1668 "canwu.decision.archive-token.v1",
1669 &(&source_root, &receipts),
1670 )?;
1671 Ok(PreparedDecisionArchive {
1672 source_root,
1673 token,
1674 blobs,
1675 receipts,
1676 })
1677 }
1678
1679 pub fn commit_decision_archive(
1680 &self,
1681 prepared: &PreparedDecisionArchive,
1682 provider: &dyn DecisionArchiveProvider,
1683 ) -> Result<Self, DecisionError> {
1684 let verified = self.verify_decision_archive(prepared, provider)?;
1685 self.commit_verified_decision_archive(&verified)
1686 }
1687
1688 pub fn verify_decision_archive(
1689 &self,
1690 prepared: &PreparedDecisionArchive,
1691 provider: &dyn DecisionArchiveProvider,
1692 ) -> Result<VerifiedDecisionArchiveCommit, DecisionError> {
1693 let keys = prepared
1694 .receipts
1695 .iter()
1696 .map(|receipt| receipt.key.clone())
1697 .collect::<Vec<_>>();
1698 let expected = self.prepare_decision_archive(&keys)?;
1699 if expected != *prepared {
1700 return Err(archive_error(
1701 "prepared decision archive is stale or altered",
1702 ));
1703 }
1704 for (blob, receipt) in prepared.blobs.iter().zip(&prepared.receipts) {
1705 let stored = provider
1706 .load_decision_archive(&receipt.locator)?
1707 .ok_or_else(|| {
1708 archive_error("prepared decision archive is not durably readable")
1709 })?;
1710 if stored != *blob || stored.content_id()? != receipt.content_id {
1711 return Err(archive_error(
1712 "stored decision archive blob failed verification",
1713 ));
1714 }
1715 }
1716 let mut additions = std::collections::BTreeMap::<
1717 DecisionArchivePageKey,
1718 Vec<&DecisionArchiveReceipt>,
1719 >::new();
1720 for receipt in &prepared.receipts {
1721 additions
1722 .entry(decision_history_page_key(&receipt.key)?)
1723 .or_default()
1724 .push(receipt);
1725 }
1726 let mut page_replacements = Vec::with_capacity(additions.len());
1727 for (page_key, receipts) in additions {
1728 let previous_page_id = self.archive_bucket_page_ids.get(&page_key).cloned();
1729 let mut merged = if let Some(resident) = self.archive_receipt_buckets.get(&page_key) {
1730 resident
1731 .iter()
1732 .map(|(key, receipt)| (key.clone(), receipt.to_receipt(key)))
1733 .collect::<std::collections::BTreeMap<_, _>>()
1734 } else if let Some(page_id) = previous_page_id.as_deref() {
1735 let page = provider
1736 .load_decision_archive_bucket_page(page_id)?
1737 .ok_or_else(|| {
1738 history_unavailable(
1739 "decision archive locator page is unavailable for authenticated merge",
1740 )
1741 })?;
1742 page.validate()?;
1743 if page.bucket != page_key.bucket
1744 || page.segment != page_key.segment
1745 || page.state_page_id()? != page_id
1746 {
1747 return Err(history_unavailable(
1748 "decision archive locator provider returned a mismatched merge page",
1749 ));
1750 }
1751 page.receipts
1752 .into_iter()
1753 .map(|receipt| (receipt.key.clone(), receipt))
1754 .collect::<std::collections::BTreeMap<_, _>>()
1755 } else {
1756 std::collections::BTreeMap::new()
1757 };
1758 for receipt in receipts {
1759 if merged
1760 .insert(receipt.key.clone(), receipt.clone())
1761 .is_some()
1762 {
1763 return Err(archive_error(
1764 "decision archive locator merge would replace an existing receipt",
1765 ));
1766 }
1767 }
1768 let page = DecisionArchiveBucketPage {
1769 format_version: DECISION_ARCHIVE_BUCKET_PAGE_FORMAT_VERSION,
1770 bucket: page_key.bucket,
1771 segment: page_key.segment,
1772 receipts: merged.into_values().collect(),
1773 };
1774 page.validate()?;
1775 page_replacements.push(VerifiedDecisionArchivePageReplacement {
1776 page_key,
1777 previous_page_id,
1778 page,
1779 });
1780 }
1781 Ok(VerifiedDecisionArchiveCommit {
1782 format_version: DECISION_ARCHIVE_FORMAT_VERSION,
1783 source_root: prepared.source_root.clone(),
1784 token: prepared.token.clone(),
1785 receipts: prepared.receipts.clone(),
1786 page_replacements,
1787 })
1788 }
1789
1790 pub fn commit_verified_decision_archive(
1791 &self,
1792 verified: &VerifiedDecisionArchiveCommit,
1793 ) -> Result<Self, DecisionError> {
1794 if verified.format_version != DECISION_ARCHIVE_FORMAT_VERSION {
1795 return Err(archive_error(
1796 "verified decision archive uses an unsupported format",
1797 ));
1798 }
1799 let mut replacement_ids = std::collections::BTreeMap::new();
1800 for replacement in &verified.page_replacements {
1801 replacement.page.validate()?;
1802 if replacement.page_key.bucket != replacement.page.bucket
1803 || replacement.page_key.segment != replacement.page.segment
1804 || replacement_ids
1805 .insert(replacement.page_key, replacement.page.state_page_id()?)
1806 .is_some()
1807 || replacement
1808 .previous_page_id
1809 .as_ref()
1810 .is_some_and(|page_id| !canonical_archive_hash(page_id))
1811 {
1812 return Err(archive_error(
1813 "verified decision archive page replacement is malformed",
1814 ));
1815 }
1816 }
1817 let mut receipts_by_page = std::collections::BTreeMap::<
1818 DecisionArchivePageKey,
1819 Vec<&DecisionArchiveReceipt>,
1820 >::new();
1821 for receipt in &verified.receipts {
1822 receipts_by_page
1823 .entry(decision_history_page_key(&receipt.key)?)
1824 .or_default()
1825 .push(receipt);
1826 }
1827 if receipts_by_page.len() != replacement_ids.len()
1828 || receipts_by_page
1829 .keys()
1830 .any(|page_key| !replacement_ids.contains_key(page_key))
1831 {
1832 return Err(archive_error(
1833 "verified decision archive does not replace every touched locator page",
1834 ));
1835 }
1836 let already_applied = verified.receipts.iter().all(|receipt| {
1837 !matches!(
1838 self.decision_locator(&receipt.key),
1839 DecisionHistoryLocation::Hot
1840 )
1841 }) && replacement_ids
1842 .iter()
1843 .all(|(page_key, page_id)| self.archive_bucket_page_ids.get(page_key) == Some(page_id));
1844 if already_applied {
1845 return Ok(self.clone());
1846 }
1847 for replacement in &verified.page_replacements {
1848 if self.archive_bucket_page_ids.get(&replacement.page_key)
1849 != replacement.previous_page_id.as_ref()
1850 {
1851 return Err(archive_error(
1852 "verified decision archive locator base changed before commit",
1853 ));
1854 }
1855 let receipts = receipts_by_page.get(&replacement.page_key).ok_or_else(|| {
1856 archive_error("verified decision archive locator page has no touched receipts")
1857 })?;
1858 for receipt in receipts {
1859 if replacement
1860 .page
1861 .receipts
1862 .binary_search_by(|candidate| candidate.key.cmp(&receipt.key))
1863 .ok()
1864 .is_none_or(|index| replacement.page.receipts[index] != **receipt)
1865 {
1866 return Err(archive_error(
1867 "verified decision archive locator page omits a touched receipt",
1868 ));
1869 }
1870 }
1871 }
1872 let keys = verified
1873 .receipts
1874 .iter()
1875 .map(|receipt| receipt.key.clone())
1876 .collect::<Vec<_>>();
1877 let expected = self.prepare_decision_archive(&keys)?;
1878 if expected.source_root != verified.source_root
1879 || expected.token != verified.token
1880 || expected.receipts != verified.receipts
1881 {
1882 return Err(archive_error(
1883 "verified decision archive is stale or altered",
1884 ));
1885 }
1886 let mut next = self.clone();
1887 for receipt in &verified.receipts {
1888 let hot_record = next
1889 .hot_archive_record(&receipt.key)
1890 .ok_or_else(|| archive_error("decision history disappeared during commit"))?;
1891 next.remove_hot_history_record(&hot_record)?;
1892 match receipt.key {
1893 DecisionHistoryKey::Ticket(id) => {
1894 next.tickets.remove(&id);
1895 }
1896 DecisionHistoryKey::Attempt(id) => {
1897 let ordinal = next.attempts_by_request.remove(&id).ok_or_else(|| {
1898 archive_error("decision attempt disappeared during commit")
1899 })?;
1900 next.attempts.remove_ordinal(ordinal);
1901 }
1902 DecisionHistoryKey::Trace(id) => {
1903 next.traces.remove_ordinal(id.get());
1904 }
1905 }
1906 }
1907 next.archive_receipt_count = next
1908 .archive_receipt_count
1909 .checked_add(verified.receipts.len() as u64)
1910 .ok_or_else(|| archive_error("decision archive receipt count overflowed"))?;
1911 for replacement in &verified.page_replacements {
1912 let receipts = replacement
1913 .page
1914 .receipts
1915 .iter()
1916 .map(|receipt| {
1917 Ok((
1918 receipt.key.clone(),
1919 CompactDecisionArchiveReceipt::from_receipt(receipt)?,
1920 ))
1921 })
1922 .collect::<Result<OrdMap<_, _>, DecisionError>>()?;
1923 next.archive_receipt_buckets
1924 .insert(replacement.page_key, receipts);
1925 next.archive_bucket_page_ids
1926 .insert(replacement.page_key, replacement.page.state_page_id()?);
1927 }
1928 next.validate_archive_directory_shape()?;
1929 for receipt in &verified.receipts {
1930 if matches!(
1931 next.decision_locator(&receipt.key),
1932 DecisionHistoryLocation::Hot
1933 ) {
1934 return Err(archive_error(
1935 "decision archive commit left a touched key in hot state",
1936 ));
1937 }
1938 }
1939 Ok(next)
1940 }
1941
1942 fn hot_archive_record(&self, key: &DecisionHistoryKey) -> Option<DecisionArchiveRecord> {
1943 match key {
1944 DecisionHistoryKey::Ticket(id) => self
1945 .tickets
1946 .get(id)
1947 .map(|ticket| ticket.as_ref().clone())
1948 .map(|ticket| DecisionArchiveRecord::Ticket { ticket }),
1949 DecisionHistoryKey::Attempt(id) => self
1950 .attempt(*id)
1951 .cloned()
1952 .map(|attempt| DecisionArchiveRecord::Attempt { attempt }),
1953 DecisionHistoryKey::Trace(id) => self
1954 .traces
1955 .get_ordinal(id.get())
1956 .cloned()
1957 .map(|trace| DecisionArchiveRecord::Trace { trace }),
1958 }
1959 }
1960
1961 fn validate_archive_eligibility(
1962 &self,
1963 record: &DecisionArchiveRecord,
1964 ) -> Result<(), DecisionError> {
1965 match record {
1966 DecisionArchiveRecord::Ticket { ticket } if ticket.is_open() => Err(archive_error(
1967 "open decision tickets cannot leave the hot mutation path",
1968 )),
1969 DecisionArchiveRecord::Trace { trace } => {
1970 let terminal = self
1971 .tickets
1972 .get(&trace.ticket_id)
1973 .is_some_and(|ticket| !ticket.is_open())
1974 || self.contains_archived_key(&DecisionHistoryKey::Ticket(trace.ticket_id));
1975 if terminal {
1976 Ok(())
1977 } else {
1978 Err(archive_error(
1979 "decision traces can archive only after their ticket is terminal",
1980 ))
1981 }
1982 }
1983 DecisionArchiveRecord::Ticket { .. } | DecisionArchiveRecord::Attempt { .. } => Ok(()),
1984 }
1985 }
1986
1987 fn hot_history_root(&self) -> Result<String, DecisionError> {
1988 let resident_count = (self.tickets.len() as u64)
1989 .checked_add(self.traces.len() as u64)
1990 .and_then(|count| count.checked_add(self.attempts.len() as u64))
1991 .ok_or_else(|| archive_error("decision hot-history count overflowed"))?;
1992 if resident_count != self.hot_history_accumulator.count {
1993 return Err(archive_error(
1994 "decision hot-history accumulator count is inconsistent",
1995 ));
1996 }
1997 decision_archive_hash(
1998 "canwu.decision.hot-history-root.v2",
1999 &(
2000 self.hot_history_accumulator.count,
2001 self.hot_history_accumulator.xor,
2002 self.hot_history_accumulator.sum,
2003 ),
2004 )
2005 }
2006
2007 pub fn hot_history_commitment(&self) -> Result<String, DecisionError> {
2008 self.hot_history_root()
2009 }
2010
2011 fn archive_receipt_root(&self) -> Result<String, DecisionError> {
2012 let mut hasher = blake3::Hasher::new();
2013 hasher.update(b"canwu.decision.archive-receipt-root.v3");
2014 hasher.update(&[0]);
2015 hasher.update(&self.archive_receipt_count.to_be_bytes());
2016 hasher.update(&(self.archive_bucket_page_ids.len() as u64).to_be_bytes());
2017 for (bucket, page_id) in &self.archive_bucket_page_ids {
2018 hasher.update(&bucket.bucket.to_be_bytes());
2019 hasher.update(&[bucket.segment]);
2020 hasher.update(&decode_archive_hash(page_id)?);
2021 }
2022 Ok(hasher.finalize().to_hex().to_string())
2023 }
2024
2025 #[must_use]
2026 pub fn decision_hot_state(&self) -> DecisionHotState {
2027 DecisionHotState {
2028 ticket_count: self.tickets.len() as u64,
2029 attempt_count: self.attempts.len() as u64,
2030 trace_count: self.traces.len() as u64,
2031 }
2032 }
2033
2034 pub fn append_attempt(&mut self, attempt: DecisionAttemptRecord) -> Result<(), DecisionError> {
2035 validate_attempt_shape(self, &attempt)?;
2036 if self.attempts_by_request.contains_key(&attempt.request_id) {
2037 return Err(DecisionError::new(
2038 DecisionErrorCode::InvalidDecision,
2039 "decision attempts must use unique request IDs",
2040 ));
2041 }
2042 let request_id = attempt.request_id;
2043 self.insert_hot_history_record(&DecisionArchiveRecord::Attempt {
2044 attempt: attempt.clone(),
2045 })?;
2046 let ordinal = self.attempts.push(attempt);
2047 self.attempts_by_request.insert(request_id, ordinal);
2048 Ok(())
2049 }
2050
2051 pub fn open_tickets(&self) -> impl Iterator<Item = &DecisionTicket> {
2052 self.tickets
2053 .values()
2054 .map(Arc::as_ref)
2055 .filter(|ticket| ticket.is_open())
2056 }
2057
2058 pub fn validate(&self) -> Result<(), DecisionError> {
2059 self.validate_archive_directory_shape()?;
2060 if self.attempts_by_request.len() != self.attempts.len() {
2061 return Err(DecisionError::new(
2062 DecisionErrorCode::InvalidDecision,
2063 "decision attempt request index is inconsistent",
2064 ));
2065 }
2066 let mut expected_deadlines = OrdMap::<SimTime, OrdSet<DecisionTicketId>>::new();
2067 for ticket in self.tickets.values().filter(|ticket| ticket.is_open()) {
2068 if let Some(deadline) = ticket.deadline {
2069 expected_deadlines
2070 .entry(deadline)
2071 .or_default()
2072 .insert(ticket.id);
2073 }
2074 }
2075 if expected_deadlines != self.deadline_index {
2076 return Err(DecisionError::new(
2077 DecisionErrorCode::InvalidDecision,
2078 "decision deadline index is inconsistent",
2079 ));
2080 }
2081 for (id, controller) in &self.controllers {
2082 if id != &controller.id {
2083 return Err(DecisionError::new(
2084 DecisionErrorCode::InvalidController,
2085 "controller map key does not match its persisted identity",
2086 ));
2087 }
2088 controller.validate()?;
2089 }
2090 for (bucket, receipts) in &self.archive_receipt_buckets {
2091 if bucket.bucket >= DECISION_HISTORY_BUCKET_COUNT || receipts.is_empty() {
2092 return Err(archive_error("decision archive receipt bucket is invalid"));
2093 }
2094 for (key, receipt) in receipts {
2095 if decision_history_page_key(key)? != *bucket
2096 || receipt.encoded_bytes == 0
2097 || matches!(self.decision_locator(key), DecisionHistoryLocation::Hot)
2098 {
2099 return Err(archive_error(
2100 "decision archive receipt is invalid or overlaps hot history",
2101 ));
2102 }
2103 }
2104 let expected_page_id = self
2105 .decision_archive_bucket_page(*bucket)?
2106 .ok_or_else(|| archive_error("nonempty archive bucket has no page"))?
2107 .state_page_id()?;
2108 if self.archive_bucket_page_ids.get(bucket) != Some(&expected_page_id) {
2109 return Err(archive_error(
2110 "decision archive bucket page commitment is inconsistent",
2111 ));
2112 }
2113 }
2114 if self.attempts.iter().any(|attempt| {
2115 attempt.request_commitment.len() != 64
2116 || !attempt
2117 .request_commitment
2118 .bytes()
2119 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
2120 }) {
2121 return Err(DecisionError::new(
2122 DecisionErrorCode::InvalidDecision,
2123 "decision attempts require a canonical request commitment",
2124 ));
2125 }
2126 for (id, ticket) in &self.tickets {
2127 if id != &ticket.id || !self.controllers.contains_key(&ticket.assigned_controller) {
2128 return Err(DecisionError::new(
2129 DecisionErrorCode::InvalidDecision,
2130 "ticket identity or assigned controller is invalid",
2131 ));
2132 }
2133 ticket.validate()?;
2134 if let Some(parent_id) = ticket.parent_ticket {
2135 let lineage_valid = match self.tickets.get(&parent_id) {
2136 Some(parent) => {
2137 !parent.is_open()
2138 && parent.decision_maker == ticket.decision_maker
2139 && parent.updated_at <= ticket.opened_at
2140 }
2141 None => self.contains_archived_key(&DecisionHistoryKey::Ticket(parent_id)),
2142 };
2143 if !lineage_valid {
2144 return Err(DecisionError::new(
2145 DecisionErrorCode::InvalidDecision,
2146 "decision ticket lineage references an open, foreign, later, or unavailable parent",
2147 ));
2148 }
2149 }
2150 }
2151 for (ordinal, trace) in self.traces.ordinals().zip(self.traces.iter()) {
2152 if trace.id.get() != ordinal || ordinal >= self.traces.next_ordinal() {
2153 return Err(DecisionError::new(
2154 DecisionErrorCode::InvalidDecision,
2155 "decision trace identity disagrees with its persistent log ordinal",
2156 ));
2157 }
2158 let ticket_version_valid = self.tickets.get(&trace.ticket_id).is_some_and(|ticket| {
2159 trace.ticket_version != 0
2160 && trace.ticket_version <= ticket.version
2161 && trace.parent_ticket == ticket.parent_ticket
2162 }) || self
2163 .contains_archived_key(&DecisionHistoryKey::Ticket(trace.ticket_id));
2164 if !ticket_version_valid || !self.controllers.contains_key(&trace.controller_id) {
2165 return Err(DecisionError::new(
2166 DecisionErrorCode::InvalidDecision,
2167 "decision trace ticket, version, or controller is invalid",
2168 ));
2169 }
2170 }
2171 for ticket in self.tickets.values() {
2172 if let DecisionTicketState::Resolved { trace_id, .. } = ticket.state {
2173 let hot_trace_matches = self
2174 .traces
2175 .get_ordinal(trace_id.get())
2176 .is_some_and(|trace| trace.id == trace_id && trace.ticket_id == ticket.id);
2177 let archived_trace_exists =
2178 self.contains_archived_key(&DecisionHistoryKey::Trace(trace_id));
2179 if !hot_trace_matches && !archived_trace_exists {
2180 return Err(DecisionError::new(
2181 DecisionErrorCode::InvalidDecision,
2182 "resolved ticket does not reference hot or archived trace evidence",
2183 ));
2184 }
2185 }
2186 }
2187 let mut request_ids = std::collections::BTreeSet::new();
2188 for (ordinal, attempt) in self.attempts.ordinals().zip(self.attempts.iter()) {
2189 if attempt.request_id.get() == 0 || !request_ids.insert(attempt.request_id) {
2190 return Err(DecisionError::new(
2191 DecisionErrorCode::InvalidDecision,
2192 "decision attempts must use unique nonzero request IDs",
2193 ));
2194 }
2195 if self.attempts_by_request.get(&attempt.request_id) != Some(&ordinal) {
2196 return Err(DecisionError::new(
2197 DecisionErrorCode::InvalidDecision,
2198 "decision attempt request index is inconsistent",
2199 ));
2200 }
2201 match &attempt.outcome {
2202 DecisionAttemptOutcome::Accepted {
2203 trace_id,
2204 command_request_id,
2205 } => {
2206 if command_request_id.is_some() && trace_id.is_none() {
2207 return Err(DecisionError::new(
2208 DecisionErrorCode::InvalidDecision,
2209 "accepted decision commands require a decision trace",
2210 ));
2211 }
2212 if trace_id.is_some_and(|trace_id| {
2213 self.traces.get_ordinal(trace_id.get()).is_none()
2214 && !self.contains_archived_key(&DecisionHistoryKey::Trace(trace_id))
2215 }) {
2216 return Err(DecisionError::new(
2217 DecisionErrorCode::InvalidDecision,
2218 "accepted decision attempt references unavailable trace evidence",
2219 ));
2220 }
2221 }
2222 DecisionAttemptOutcome::Rejected { message, .. } => {
2223 require_text(message, "decision rejection message")?;
2224 }
2225 }
2226 }
2227 Ok(())
2228 }
2229
2230 fn validate_archive_directory_shape(&self) -> Result<(), DecisionError> {
2231 if self
2232 .archive_bucket_page_ids
2233 .iter()
2234 .any(|(bucket, page_id)| {
2235 bucket.bucket >= DECISION_HISTORY_BUCKET_COUNT || !canonical_archive_hash(page_id)
2236 })
2237 {
2238 return Err(archive_error(
2239 "decision archive bucket directory is malformed",
2240 ));
2241 }
2242 if (self.archive_receipt_count == 0) != self.archive_bucket_page_ids.is_empty() {
2243 return Err(archive_error(
2244 "decision archive receipt count and page directory disagree",
2245 ));
2246 }
2247 if !self.archive_receipt_buckets.is_empty() {
2248 let resident_count =
2249 self.archive_receipt_buckets
2250 .values()
2251 .try_fold(0_u64, |total, receipts| {
2252 total.checked_add(receipts.len() as u64).ok_or_else(|| {
2253 archive_error("decision archive receipt count overflowed")
2254 })
2255 })?;
2256 if resident_count > self.archive_receipt_count {
2257 return Err(archive_error(
2258 "resident decision archive receipt count exceeds the committed total",
2259 ));
2260 }
2261 }
2262 Ok(())
2263 }
2264
2265 pub fn apply(
2266 &mut self,
2267 mutation: DecisionMutation,
2268 at: SimTime,
2269 trace_id: Option<DecisionTraceId>,
2270 ) -> Result<PreparedDecision, DecisionError> {
2271 let prepared = match mutation {
2272 DecisionMutation::RegisterController { controller } => {
2273 controller.validate()?;
2274 if self.controllers.contains_key(&controller.id) {
2275 return Err(DecisionError::new(
2276 DecisionErrorCode::DuplicateController,
2277 format!(
2278 "decision controller {} is already registered",
2279 controller.id
2280 ),
2281 ));
2282 }
2283 self.controllers
2284 .insert(controller.id.clone(), Arc::new(controller));
2285 PreparedDecision::default()
2286 }
2287 DecisionMutation::Open { mut ticket } => {
2288 ticket.validate()?;
2289 if self.tickets.contains_key(&ticket.id) {
2290 return Err(DecisionError::new(
2291 DecisionErrorCode::DuplicateTicket,
2292 format!("decision ticket {} is already present", ticket.id),
2293 ));
2294 }
2295 if !self.controllers.contains_key(&ticket.assigned_controller) {
2296 return Err(DecisionError::new(
2297 DecisionErrorCode::InvalidController,
2298 format!(
2299 "decision ticket {} names unknown controller {}",
2300 ticket.id, ticket.assigned_controller
2301 ),
2302 ));
2303 }
2304 if ticket.deadline.is_some_and(|deadline| deadline < at) {
2305 return Err(DecisionError::new(
2306 DecisionErrorCode::InvalidDecision,
2307 "decision deadline precedes its admission time",
2308 ));
2309 }
2310 if let Some(parent_id) = ticket.parent_ticket {
2311 self.validate_parent_admission(parent_id, &ticket.decision_maker)?;
2312 }
2313 let persisted = DecisionTicket {
2314 id: ticket.id,
2315 definition: ticket.definition,
2316 decision_maker: ticket.decision_maker,
2317 assigned_controller: ticket.assigned_controller,
2318 summary: ticket.summary,
2319 context: ticket.context,
2320 options: std::mem::take(&mut ticket.options),
2321 opened_at: at,
2322 updated_at: at,
2323 deadline: ticket.deadline,
2324 version: 1,
2325 state: DecisionTicketState::Open,
2326 parent_ticket: ticket.parent_ticket,
2327 };
2328 if let Some(deadline) = persisted.deadline {
2329 self.deadline_index
2330 .entry(deadline)
2331 .or_default()
2332 .insert(persisted.id);
2333 }
2334 self.insert_hot_history_record(&DecisionArchiveRecord::Ticket {
2335 ticket: persisted.clone(),
2336 })?;
2337 self.tickets.insert(persisted.id, Arc::new(persisted));
2338 PreparedDecision::default()
2339 }
2340 DecisionMutation::ReplaceOptions {
2341 ticket_id,
2342 expected_version,
2343 context,
2344 mut options,
2345 } => {
2346 context.validate()?;
2347 canonicalize_options(&mut options)?;
2348 let previous = self
2349 .tickets
2350 .get(&ticket_id)
2351 .map(|ticket| ticket.as_ref().clone())
2352 .ok_or_else(|| {
2353 DecisionError::new(
2354 DecisionErrorCode::TicketNotFound,
2355 format!("decision ticket {ticket_id} was not found"),
2356 )
2357 })?;
2358 let ticket = self.open_ticket_mut(ticket_id, expected_version, at)?;
2359 ticket.context = context;
2360 ticket.options = options;
2361 ticket.updated_at = at;
2362 ticket.version = ticket.version.checked_add(1).ok_or_else(|| {
2363 DecisionError::new(
2364 DecisionErrorCode::InvalidDecision,
2365 "decision ticket version is exhausted",
2366 )
2367 })?;
2368 ticket.validate()?;
2369 let updated = ticket.clone();
2370 let _ = ticket;
2371 self.replace_hot_history_record(
2372 &DecisionArchiveRecord::Ticket { ticket: previous },
2373 &DecisionArchiveRecord::Ticket { ticket: updated },
2374 )?;
2375 PreparedDecision::default()
2376 }
2377 DecisionMutation::Resolve {
2378 ticket_id,
2379 expected_version,
2380 controller_id,
2381 policy,
2382 decision,
2383 command_request_id,
2384 } => self.resolve(
2385 ticket_id,
2386 expected_version,
2387 &controller_id,
2388 policy,
2389 decision,
2390 command_request_id,
2391 at,
2392 trace_id.ok_or_else(|| {
2393 DecisionError::new(
2394 DecisionErrorCode::InvalidDecision,
2395 "decision resolution requires a claimed trace ID",
2396 )
2397 })?,
2398 )?,
2399 DecisionMutation::Cancel {
2400 ticket_id,
2401 expected_version,
2402 reason,
2403 } => {
2404 require_text(&reason, "decision cancellation reason")?;
2405 let previous = self
2406 .tickets
2407 .get(&ticket_id)
2408 .map(|ticket| ticket.as_ref().clone())
2409 .ok_or_else(|| {
2410 DecisionError::new(
2411 DecisionErrorCode::TicketNotFound,
2412 format!("decision ticket {ticket_id} was not found"),
2413 )
2414 })?;
2415 let ticket = self.open_ticket_mut(ticket_id, expected_version, at)?;
2416 ticket.updated_at = at;
2417 ticket.version = ticket.version.checked_add(1).ok_or_else(|| {
2418 DecisionError::new(
2419 DecisionErrorCode::InvalidDecision,
2420 "decision ticket version is exhausted",
2421 )
2422 })?;
2423 ticket.state = DecisionTicketState::Cancelled { reason };
2424 ticket.validate()?;
2425 let deadline = ticket.deadline;
2426 let updated = ticket.clone();
2427 let _ = ticket;
2428 self.replace_hot_history_record(
2429 &DecisionArchiveRecord::Ticket { ticket: previous },
2430 &DecisionArchiveRecord::Ticket { ticket: updated },
2431 )?;
2432 self.remove_deadline(ticket_id, deadline);
2433 PreparedDecision::default()
2434 }
2435 };
2436 self.advance_time(at)?;
2437 Ok(prepared)
2438 }
2439
2440 fn validate_parent_admission(
2446 &self,
2447 parent_id: DecisionTicketId,
2448 decision_maker: &canwu_core::EntityRef,
2449 ) -> Result<(), DecisionError> {
2450 let Some(parent) = self.ticket(parent_id) else {
2451 return Err(DecisionError::new(
2452 DecisionErrorCode::TicketNotFound,
2453 format!("decision ticket parent {parent_id} is not in hot decision history"),
2454 ));
2455 };
2456 if parent.is_open() {
2457 return Err(DecisionError::new(
2458 DecisionErrorCode::InvalidDecision,
2459 format!("decision ticket parent {parent_id} is still open"),
2460 ));
2461 }
2462 if &parent.decision_maker != decision_maker {
2463 return Err(DecisionError::new(
2464 DecisionErrorCode::InvalidDecision,
2465 format!("decision ticket parent {parent_id} belongs to a different decision maker"),
2466 ));
2467 }
2468 Ok(())
2469 }
2470
2471 fn open_ticket_mut(
2472 &mut self,
2473 ticket_id: DecisionTicketId,
2474 expected_version: u64,
2475 at: SimTime,
2476 ) -> Result<&mut DecisionTicket, DecisionError> {
2477 let ticket = self.tickets.get_mut(&ticket_id).ok_or_else(|| {
2478 DecisionError::new(
2479 DecisionErrorCode::TicketNotFound,
2480 format!("decision ticket {ticket_id} was not found"),
2481 )
2482 })?;
2483 let ticket = Arc::make_mut(ticket);
2484 if !ticket.is_open() || ticket.deadline.is_some_and(|deadline| deadline < at) {
2485 return Err(DecisionError::new(
2486 DecisionErrorCode::ClosedTicket,
2487 format!("decision ticket {ticket_id} is not open"),
2488 ));
2489 }
2490 if ticket.version != expected_version {
2491 return Err(DecisionError::new(
2492 DecisionErrorCode::VersionConflict,
2493 format!(
2494 "decision ticket {ticket_id} is at version {}, expected {expected_version}",
2495 ticket.version
2496 ),
2497 ));
2498 }
2499 Ok(ticket)
2500 }
2501
2502 #[allow(clippy::too_many_arguments)]
2503 fn resolve(
2504 &mut self,
2505 ticket_id: DecisionTicketId,
2506 expected_version: u64,
2507 controller_id: &str,
2508 policy: DecisionPolicyIdentity,
2509 decision: PolicyDecision,
2510 command_request_id: Option<CommandRequestId>,
2511 at: SimTime,
2512 trace_id: DecisionTraceId,
2513 ) -> Result<PreparedDecision, DecisionError> {
2514 let controller = self.controllers.get(controller_id).ok_or_else(|| {
2515 DecisionError::new(
2516 DecisionErrorCode::InvalidController,
2517 format!("decision controller {controller_id} was not found"),
2518 )
2519 })?;
2520 if controller.policy != policy {
2521 return Err(DecisionError::new(
2522 DecisionErrorCode::PolicyMismatch,
2523 "decision resolution policy does not match the persisted controller binding",
2524 ));
2525 }
2526 let random_evidence_matches_policy = match controller.policy.kind {
2531 crate::DecisionPolicyKind::Random => {
2532 decision.random.is_some() && decision.stage.is_none()
2533 }
2534 crate::DecisionPolicyKind::Utility => {
2535 decision.random.is_some() == decision.is_random_tie_break()
2536 && (controller.random_tie_break || decision.random.is_none())
2537 }
2538 crate::DecisionPolicyKind::Rule
2539 | crate::DecisionPolicyKind::Human
2540 | crate::DecisionPolicyKind::External
2541 | crate::DecisionPolicyKind::Llm => {
2542 decision.random.is_none() && !decision.is_random_tie_break()
2543 }
2544 };
2545 if !random_evidence_matches_policy {
2546 return Err(DecisionError::new(
2547 DecisionErrorCode::PolicyMismatch,
2548 "random decision controllers require random draw evidence, utility controllers accept it only for a random tie-break, and other controllers reject it",
2549 ));
2550 }
2551 let previous_ticket = self
2552 .tickets
2553 .get(&ticket_id)
2554 .map(|ticket| ticket.as_ref().clone())
2555 .ok_or_else(|| {
2556 DecisionError::new(
2557 DecisionErrorCode::TicketNotFound,
2558 format!("decision ticket {ticket_id} was not found"),
2559 )
2560 })?;
2561 let ticket = self.open_ticket_mut(ticket_id, expected_version, at)?;
2562 if ticket.assigned_controller != controller_id {
2563 return Err(DecisionError::new(
2564 DecisionErrorCode::InvalidController,
2565 "decision resolution came from a controller not assigned to the ticket",
2566 ));
2567 }
2568 decision.validate(ticket)?;
2569 let action = match &decision.outcome {
2570 DecisionOutcome::Selected { option_id } => {
2571 ticket.option(option_id).map(|option| option.action.clone())
2572 }
2573 DecisionOutcome::Deferred { .. } => None,
2574 DecisionOutcome::Pending { .. } | DecisionOutcome::PendingRandom { .. } => {
2575 return Err(DecisionError::new(
2576 DecisionErrorCode::InvalidDecision,
2577 "pending policy outcomes are not authoritative decision mutations",
2578 ));
2579 }
2580 };
2581 if matches!(action, Some(DecisionAction::Command { .. })) != command_request_id.is_some() {
2582 return Err(DecisionError::new(
2583 DecisionErrorCode::InvalidDecision,
2584 "command actions require exactly one command request ID",
2585 ));
2586 }
2587 let trace = DecisionTrace {
2588 id: trace_id,
2589 ticket_id,
2590 ticket_version: ticket.version,
2591 controller_id: controller_id.to_owned(),
2592 policy,
2593 decided_at: at,
2594 outcome: decision.outcome.clone(),
2595 summary: decision.summary,
2596 evaluations: decision.evaluations,
2597 external: decision.external,
2598 random: decision.random,
2599 command_request_id,
2600 stage: decision.stage,
2601 fired_guards: decision.fired_guards,
2602 parent_ticket: ticket.parent_ticket,
2603 };
2604 ticket.updated_at = at;
2605 ticket.version = ticket.version.checked_add(1).ok_or_else(|| {
2606 DecisionError::new(
2607 DecisionErrorCode::InvalidDecision,
2608 "decision ticket version is exhausted",
2609 )
2610 })?;
2611 if let DecisionOutcome::Selected { option_id } = &trace.outcome {
2612 ticket.state = DecisionTicketState::Resolved {
2613 option_id: option_id.clone(),
2614 trace_id,
2615 };
2616 }
2617 ticket.validate()?;
2618 let deadline = ticket.deadline;
2619 let updated_ticket = ticket.clone();
2620 let _ = ticket;
2621 self.replace_hot_history_record(
2622 &DecisionArchiveRecord::Ticket {
2623 ticket: previous_ticket,
2624 },
2625 &DecisionArchiveRecord::Ticket {
2626 ticket: updated_ticket,
2627 },
2628 )?;
2629 self.remove_deadline(ticket_id, deadline);
2630 self.insert_hot_history_record(&DecisionArchiveRecord::Trace {
2631 trace: trace.clone(),
2632 })?;
2633 let ordinal = self.traces.push(trace.clone());
2634 debug_assert_eq!(ordinal, trace.id.get());
2635 Ok(PreparedDecision {
2636 trace: Some(trace),
2637 action,
2638 })
2639 }
2640
2641 pub fn advance_time(&mut self, at: SimTime) -> Result<(), DecisionError> {
2642 let due = self
2643 .deadline_index
2644 .range(..at)
2645 .map(|(deadline, tickets)| (*deadline, tickets.iter().copied().collect::<Vec<_>>()))
2646 .collect::<Vec<_>>();
2647 for (deadline, ticket_ids) in due {
2648 self.deadline_index.remove(&deadline);
2649 for ticket_id in ticket_ids {
2650 let previous = self
2651 .tickets
2652 .get(&ticket_id)
2653 .map(|ticket| ticket.as_ref().clone())
2654 .ok_or_else(|| {
2655 DecisionError::new(
2656 DecisionErrorCode::InvalidDecision,
2657 "decision ticket index changed during time advancement",
2658 )
2659 })?;
2660 let Some(ticket) = self.tickets.get_mut(&ticket_id) else {
2661 return Err(DecisionError::new(
2662 DecisionErrorCode::InvalidDecision,
2663 "decision ticket index changed during time advancement",
2664 ));
2665 };
2666 let ticket = Arc::make_mut(ticket);
2667 let updated =
2668 if ticket.is_open() && ticket.deadline.is_some_and(|deadline| deadline < at) {
2669 ticket.updated_at = at;
2670 ticket.version = ticket.version.checked_add(1).ok_or_else(|| {
2671 DecisionError::new(
2672 DecisionErrorCode::InvalidDecision,
2673 "decision ticket version is exhausted",
2674 )
2675 })?;
2676 ticket.state = DecisionTicketState::Expired;
2677 ticket.validate()?;
2678 Some(ticket.clone())
2679 } else {
2680 None
2681 };
2682 let _ = ticket;
2683 if let Some(updated) = updated {
2684 self.replace_hot_history_record(
2685 &DecisionArchiveRecord::Ticket { ticket: previous },
2686 &DecisionArchiveRecord::Ticket { ticket: updated },
2687 )?;
2688 }
2689 }
2690 }
2691 Ok(())
2692 }
2693
2694 fn remove_deadline(&mut self, ticket_id: DecisionTicketId, deadline: Option<SimTime>) {
2695 let Some(deadline) = deadline else {
2696 return;
2697 };
2698 if let Some(mut tickets) = self.deadline_index.get(&deadline).cloned() {
2699 tickets.remove(&ticket_id);
2700 if tickets.is_empty() {
2701 self.deadline_index.remove(&deadline);
2702 } else {
2703 self.deadline_index.insert(deadline, tickets);
2704 }
2705 }
2706 }
2707}
2708
2709fn validate_attempt_shape(
2710 state: &DecisionState,
2711 attempt: &DecisionAttemptRecord,
2712) -> Result<(), DecisionError> {
2713 if attempt.request_id.get() == 0 || !canonical_archive_hash(&attempt.request_commitment) {
2714 return Err(DecisionError::new(
2715 DecisionErrorCode::InvalidDecision,
2716 "decision attempt identity or request commitment is invalid",
2717 ));
2718 }
2719 match &attempt.outcome {
2720 DecisionAttemptOutcome::Accepted {
2721 trace_id,
2722 command_request_id,
2723 } => {
2724 if command_request_id.is_some() && trace_id.is_none() {
2725 return Err(DecisionError::new(
2726 DecisionErrorCode::InvalidDecision,
2727 "accepted decision commands require a decision trace",
2728 ));
2729 }
2730 if trace_id.is_some_and(|trace_id| {
2731 state.traces.get_ordinal(trace_id.get()).is_none()
2732 && !state.contains_archived_key(&DecisionHistoryKey::Trace(trace_id))
2733 }) {
2734 return Err(DecisionError::new(
2735 DecisionErrorCode::InvalidDecision,
2736 "accepted decision attempt references unavailable trace evidence",
2737 ));
2738 }
2739 }
2740 DecisionAttemptOutcome::Rejected { message, .. } => {
2741 require_text(message, "decision rejection message")?;
2742 }
2743 }
2744 Ok(())
2745}
2746
2747fn decision_hot_leaf_hash(
2748 key: &DecisionHistoryKey,
2749 record: &DecisionArchiveRecord,
2750) -> Result<[u8; 32], DecisionError> {
2751 if record.key() != *key {
2752 return Err(archive_error(
2753 "decision hot-history leaf identity is inconsistent",
2754 ));
2755 }
2756 let bytes = serde_json::to_vec(&(key, record)).map_err(|error| {
2757 archive_error(format!("cannot encode decision hot-history leaf: {error}"))
2758 })?;
2759 let mut hasher = blake3::Hasher::new();
2760 hasher.update(b"canwu.decision.hot-history-leaf.v2");
2761 hasher.update(&[0]);
2762 hasher.update(&bytes);
2763 Ok(*hasher.finalize().as_bytes())
2764}
2765
2766fn add_digest_mod_256(target: &mut [u8; 32], digest: [u8; 32]) {
2767 let mut carry = 0_u16;
2768 for index in (0..target.len()).rev() {
2769 let value = u16::from(target[index]) + u16::from(digest[index]) + carry;
2770 target[index] = u8::try_from(value % 256).expect("modulo 256 always fits in u8");
2771 carry = value >> 8;
2772 }
2773}
2774
2775fn subtract_digest_mod_256(target: &mut [u8; 32], digest: [u8; 32]) {
2776 let mut borrow = 0_i16;
2777 for index in (0..target.len()).rev() {
2778 let value = i16::from(target[index]) - i16::from(digest[index]) - borrow;
2779 if value < 0 {
2780 target[index] =
2781 u8::try_from(value + 256).expect("normalized byte subtraction fits in u8");
2782 borrow = 1;
2783 } else {
2784 target[index] = u8::try_from(value).expect("non-negative byte subtraction fits in u8");
2785 borrow = 0;
2786 }
2787 }
2788}
2789
2790pub fn decision_history_bucket(key: &DecisionHistoryKey) -> Result<u16, DecisionError> {
2791 Ok(decision_history_page_key(key)?.bucket)
2792}
2793
2794pub fn decision_history_page_key(
2795 key: &DecisionHistoryKey,
2796) -> Result<DecisionArchivePageKey, DecisionError> {
2797 let encoded = serde_json::to_vec(key)
2798 .map_err(|error| archive_error(format!("cannot encode decision history key: {error}")))?;
2799 let digest = blake3::hash(&encoded);
2800 let bytes = digest.as_bytes();
2801 Ok(DecisionArchivePageKey {
2802 bucket: (u16::from(bytes[0]) << 4) | u16::from(bytes[1] >> 4),
2803 segment: bytes[1] & 0x0f,
2804 })
2805}
2806
2807#[derive(Default)]
2808struct DecisionScaleArchive {
2809 blobs: RefCell<StdBTreeMap<String, DecisionArchiveBlob>>,
2810 pages: RefCell<StdBTreeMap<String, DecisionArchiveBucketPage>>,
2811}
2812
2813impl DecisionArchiveProvider for DecisionScaleArchive {
2814 fn load_decision_archive(
2815 &self,
2816 locator: &str,
2817 ) -> Result<Option<DecisionArchiveBlob>, DecisionError> {
2818 Ok(self.blobs.borrow().get(locator).cloned())
2819 }
2820
2821 fn load_decision_archive_bucket_page(
2822 &self,
2823 page_id: &str,
2824 ) -> Result<Option<DecisionArchiveBucketPage>, DecisionError> {
2825 Ok(self.pages.borrow().get(page_id).cloned())
2826 }
2827}
2828
2829impl DecisionArchiveStore for DecisionScaleArchive {
2830 fn store_decision_archive(
2831 &self,
2832 blob: &DecisionArchiveBlob,
2833 ) -> Result<DecisionArchiveStoreOutcome, DecisionError> {
2834 let locator = blob.content_id()?;
2835 let mut blobs = self.blobs.borrow_mut();
2836 if let Some(existing) = blobs.get(&locator) {
2837 if existing != blob {
2838 return Err(archive_error(
2839 "decision scale archive locator contains different content",
2840 ));
2841 }
2842 return Ok(DecisionArchiveStoreOutcome::AlreadyStored);
2843 }
2844 blobs.insert(locator, blob.clone());
2845 Ok(DecisionArchiveStoreOutcome::Stored)
2846 }
2847}
2848
2849#[doc(hidden)]
2853pub fn format8_decision_locator_scale_fixture(
2854 key_count: usize,
2855) -> Result<DecisionLocatorScaleFixture, DecisionError> {
2856 let mut state = DecisionState::default();
2857 let archive = DecisionScaleArchive::default();
2858 let mut archive_batches = 0_u64;
2859 let mut next_ordinal = 1_usize;
2860 while next_ordinal <= key_count {
2861 let batch_end = next_ordinal
2862 .saturating_add(MAX_DECISION_ARCHIVE_BATCH_ENTRIES - 1)
2863 .min(key_count);
2864 let mut keys = Vec::with_capacity(batch_end - next_ordinal + 1);
2865 for ordinal in next_ordinal..=batch_end {
2866 let ordinal = u64::try_from(ordinal)
2867 .map_err(|_| archive_error("decision locator scale key exceeds u64"))?;
2868 let request_id = canwu_core::DecisionRequestId::new(ordinal);
2869 let request_commitment = blake3::hash(&ordinal.to_be_bytes()).to_hex().to_string();
2870 state.append_attempt(DecisionAttemptRecord {
2871 request_id,
2872 request_commitment,
2873 at: SimTime::from_minutes(
2874 i64::try_from(ordinal)
2875 .map_err(|_| archive_error("decision scale time exceeds i64"))?,
2876 ),
2877 revision_before: ordinal - 1,
2878 expected_revision: ordinal - 1,
2879 outcome: DecisionAttemptOutcome::Rejected {
2880 code: crate::DecisionAttemptErrorCode::InvalidDecision,
2881 message: "format8 production-path scale attempt".to_owned(),
2882 },
2883 })?;
2884 keys.push(DecisionHistoryKey::Attempt(request_id));
2885 }
2886 let prepared = state.prepare_decision_archive(&keys)?;
2887 for blob in &prepared.blobs {
2888 let _ = archive.store_decision_archive(blob)?;
2889 }
2890 let verified = state.verify_decision_archive(&prepared, &archive)?;
2891 state = state.commit_verified_decision_archive(&verified)?;
2892 let touched_pages = keys
2893 .iter()
2894 .map(decision_history_page_key)
2895 .collect::<Result<OrdSet<_>, _>>()?;
2896 for page_key in touched_pages {
2897 let page = state
2898 .decision_archive_bucket_page(page_key)?
2899 .ok_or_else(|| archive_error("committed decision locator page is missing"))?;
2900 let page_id = page.state_page_id()?;
2901 archive.pages.borrow_mut().insert(page_id, page);
2902 }
2903 archive_batches = archive_batches
2904 .checked_add(1)
2905 .ok_or_else(|| archive_error("decision archive batch count overflowed"))?;
2906 next_ordinal = batch_end.saturating_add(1);
2907 }
2908
2909 let archive_root = state.archive_receipt_root()?;
2910 let hot_bytes = serde_json::to_vec(&state.paged_checkpoint_hot_state()).map_err(|error| {
2911 archive_error(format!(
2912 "cannot persist root-only decision hot state: {error}"
2913 ))
2914 })?;
2915 let hot = serde_json::from_slice::<DecisionState>(&hot_bytes).map_err(|error| {
2916 archive_error(format!(
2917 "cannot restart root-only decision hot state: {error}"
2918 ))
2919 })?;
2920 let restarted = DecisionState::from_paged_checkpoint_root(
2921 hot,
2922 state.archive_bucket_page_ids.clone(),
2923 state.archive_receipt_count,
2924 &archive_root,
2925 )?;
2926 let samples: OrdSet<usize> = [1_usize, key_count.saturating_div(2).max(1), key_count]
2927 .into_iter()
2928 .filter(|ordinal| *ordinal <= key_count)
2929 .collect::<OrdSet<_>>();
2930 for ordinal in &samples {
2931 let key = DecisionHistoryKey::Attempt(canwu_core::DecisionRequestId::new(
2932 u64::try_from(*ordinal)
2933 .map_err(|_| archive_error("decision scale sample exceeds u64"))?,
2934 ));
2935 if restarted.load_decision_history(&key, &archive)?.is_none() {
2936 return Err(archive_error(
2937 "root-only decision restart lost exact archived history",
2938 ));
2939 }
2940 }
2941 let reachable = restarted.archive_reachability(&archive)?;
2942 let mut max_page_entries = 0_u64;
2943 let mut max_page_encoded_bytes = 0_u64;
2944 for page_id in restarted.archive_bucket_page_ids.values() {
2945 let page = archive
2946 .load_decision_archive_bucket_page(page_id)?
2947 .ok_or_else(|| archive_error("decision scale locator page is unavailable"))?;
2948 max_page_entries = max_page_entries.max(page.receipts.len() as u64);
2949 max_page_encoded_bytes = max_page_encoded_bytes.max(
2950 serde_json::to_vec(&page)
2951 .map_err(|error| archive_error(format!("cannot size locator page: {error}")))?
2952 .len() as u64,
2953 );
2954 }
2955 let entry_bytes = size_of::<DecisionHistoryKey>()
2956 .saturating_add(size_of::<CompactDecisionArchiveReceipt>())
2957 .saturating_add(48);
2958 let metrics = DecisionLocatorScaleMetrics {
2959 entries: restarted.archived_history_count() as u64,
2960 locator_pages: restarted.archive_bucket_page_ids.len() as u64,
2961 max_page_entries,
2962 max_page_encoded_bytes,
2963 archive_batches,
2964 exact_restart_queries: samples.len() as u64,
2965 reachable_blob_locators: reachable.blob_locators.len() as u64,
2966 estimated_resident_structural_bytes: (restarted.archived_history_count() as u64)
2967 .saturating_mul(entry_bytes as u64),
2968 root_hash: archive_root,
2969 };
2970 Ok(DecisionLocatorScaleFixture {
2971 state,
2972 archive_blobs: archive.blobs.into_inner().into_values().collect(),
2973 metrics,
2974 })
2975}
2976
2977pub fn format8_decision_locator_scale_probe(
2980 key_count: usize,
2981) -> Result<DecisionLocatorScaleMetrics, DecisionError> {
2982 Ok(format8_decision_locator_scale_fixture(key_count)?.metrics)
2983}
2984
2985pub fn format8_trace_locator_scale_probe(
2989 trace_count: usize,
2990) -> Result<TraceLocatorScaleMetrics, DecisionError> {
2991 if trace_count == 0 {
2992 return Err(archive_error(
2993 "trace locator scale probe requires at least one trace",
2994 ));
2995 }
2996 let mut state = DecisionState::default();
2997 let controller_id = "format8-trace-controller".to_owned();
2998 let policy =
2999 DecisionPolicyIdentity::new(crate::DecisionPolicyKind::Rule, "format8-trace-policy", "1");
3000 state.controllers.insert(
3001 controller_id.clone(),
3002 Arc::new(DecisionControllerBinding::new(
3003 controller_id.clone(),
3004 policy.clone(),
3005 crate::DecisionAuthority::NoResponsibleActor {
3006 reason: "Format-8 trace scale fixture".to_owned(),
3007 },
3008 )),
3009 );
3010 let ticket_id = DecisionTicketId::new(1);
3011 let ticket = DecisionTicket {
3012 id: ticket_id,
3013 definition: "format8-trace-scale".to_owned(),
3014 decision_maker: canwu_core::EntityRef::Person(canwu_core::PersonId::new(1)),
3015 assigned_controller: controller_id.clone(),
3016 summary: "Format-8 trace scale ticket".to_owned(),
3017 context: crate::DecisionContext::new("format8-trace-scale", serde_json::json!({})),
3018 options: vec![crate::DecisionOption::new("defer", "Defer")],
3019 opened_at: SimTime::EPOCH,
3020 updated_at: SimTime::EPOCH,
3021 deadline: None,
3022 version: 1,
3023 state: DecisionTicketState::Cancelled {
3024 reason: "Scale fixture is terminal".to_owned(),
3025 },
3026 parent_ticket: None,
3027 };
3028 ticket.validate()?;
3029 state.insert_hot_history_record(&DecisionArchiveRecord::Ticket {
3030 ticket: ticket.clone(),
3031 })?;
3032 state.tickets.insert(ticket_id, Arc::new(ticket));
3033 for ordinal in 1..=trace_count {
3034 let ordinal = u64::try_from(ordinal)
3035 .map_err(|_| archive_error("trace locator scale ordinal exceeds u64"))?;
3036 let trace = DecisionTrace {
3037 id: DecisionTraceId::new(ordinal),
3038 ticket_id,
3039 ticket_version: 1,
3040 controller_id: controller_id.clone(),
3041 policy: policy.clone(),
3042 decided_at: SimTime::from_minutes(
3043 i64::try_from(ordinal)
3044 .map_err(|_| archive_error("trace locator scale time exceeds i64"))?,
3045 ),
3046 outcome: DecisionOutcome::Deferred {
3047 reason: "trace-scale".to_owned(),
3048 },
3049 summary: "trace-scale".to_owned(),
3050 evaluations: Vec::new(),
3051 external: None,
3052 random: None,
3053 command_request_id: None,
3054 stage: None,
3055 fired_guards: Vec::new(),
3056 parent_ticket: None,
3057 };
3058 state.insert_hot_history_record(&DecisionArchiveRecord::Trace {
3059 trace: trace.clone(),
3060 })?;
3061 let inserted = state.traces.push(trace);
3062 if inserted != ordinal {
3063 return Err(archive_error(
3064 "trace locator scale ordinal insertion is inconsistent",
3065 ));
3066 }
3067 }
3068 state.validate()?;
3069 let samples: OrdSet<usize> = [1_usize, trace_count.saturating_div(2).max(1), trace_count]
3070 .into_iter()
3071 .filter(|ordinal| *ordinal <= trace_count)
3072 .collect::<OrdSet<_>>();
3073 for ordinal in &samples {
3074 let key = DecisionHistoryKey::Trace(DecisionTraceId::new(
3075 u64::try_from(*ordinal)
3076 .map_err(|_| archive_error("trace locator scale sample exceeds u64"))?,
3077 ));
3078 if state.decision_locator(&key) != DecisionHistoryLocation::Hot {
3079 return Err(archive_error(
3080 "trace locator ordinal index lost a retained hot trace",
3081 ));
3082 }
3083 }
3084 let target = DecisionHistoryKey::Trace(DecisionTraceId::new(
3085 u64::try_from(trace_count)
3086 .map_err(|_| archive_error("trace locator scale target exceeds u64"))?,
3087 ));
3088 let archive = DecisionScaleArchive::default();
3089 let prepared = state.prepare_decision_archive(std::slice::from_ref(&target))?;
3090 for blob in &prepared.blobs {
3091 let _ = archive.store_decision_archive(blob)?;
3092 }
3093 let verified = state.verify_decision_archive(&prepared, &archive)?;
3094 let committed = state.commit_verified_decision_archive(&verified)?;
3095 let target_archived = matches!(
3096 committed.decision_locator(&target),
3097 DecisionHistoryLocation::Archived { .. }
3098 );
3099 if !target_archived {
3100 return Err(archive_error(
3101 "trace locator scale archive commit did not release its target",
3102 ));
3103 }
3104 Ok(TraceLocatorScaleMetrics {
3105 hot_trace_entries: trace_count as u64,
3106 indexed_lookup_samples: samples.len() as u64,
3107 archive_commit_entries: verified.receipts.len() as u64,
3108 target_archived,
3109 })
3110}
3111
3112#[cfg(test)]
3113mod archive_restart_tests {
3114 use super::*;
3115
3116 fn linked_decision_state(resolved_ticket: bool, accepted_attempt: bool) -> DecisionState {
3117 let mut state = DecisionState::default();
3118 let controller_id = "restart-controller".to_owned();
3119 let policy =
3120 DecisionPolicyIdentity::new(crate::DecisionPolicyKind::Rule, "restart-policy", "1");
3121 state.controllers.insert(
3122 controller_id.clone(),
3123 Arc::new(DecisionControllerBinding::new(
3124 controller_id.clone(),
3125 policy.clone(),
3126 crate::DecisionAuthority::NoResponsibleActor {
3127 reason: "restart dependency fixture".to_owned(),
3128 },
3129 )),
3130 );
3131 let ticket_id = DecisionTicketId::new(91);
3132 let trace_id = DecisionTraceId::new(1);
3133 let ticket = DecisionTicket {
3134 id: ticket_id,
3135 definition: "restart-dependency".to_owned(),
3136 decision_maker: canwu_core::EntityRef::Person(canwu_core::PersonId::new(9)),
3137 assigned_controller: controller_id.clone(),
3138 summary: "Restart dependency fixture".to_owned(),
3139 context: crate::DecisionContext::new("restart-dependency", serde_json::json!({})),
3140 options: vec![crate::DecisionOption::new("accept", "Accept")],
3141 opened_at: SimTime::EPOCH,
3142 updated_at: SimTime::from_minutes(1),
3143 deadline: None,
3144 version: 2,
3145 state: if resolved_ticket {
3146 DecisionTicketState::Resolved {
3147 option_id: "accept".to_owned(),
3148 trace_id,
3149 }
3150 } else {
3151 DecisionTicketState::Cancelled {
3152 reason: "terminal fixture".to_owned(),
3153 }
3154 },
3155 parent_ticket: None,
3156 };
3157 state
3158 .insert_hot_history_record(&DecisionArchiveRecord::Ticket {
3159 ticket: ticket.clone(),
3160 })
3161 .expect("insert fixture ticket history");
3162 state.tickets.insert(ticket_id, Arc::new(ticket));
3163 let trace = DecisionTrace {
3164 id: trace_id,
3165 ticket_id,
3166 ticket_version: 1,
3167 controller_id,
3168 policy,
3169 decided_at: SimTime::from_minutes(1),
3170 outcome: if resolved_ticket {
3171 DecisionOutcome::Selected {
3172 option_id: "accept".to_owned(),
3173 }
3174 } else {
3175 DecisionOutcome::Deferred {
3176 reason: "terminal fixture".to_owned(),
3177 }
3178 },
3179 summary: "Restart dependency trace".to_owned(),
3180 evaluations: Vec::new(),
3181 external: None,
3182 random: None,
3183 command_request_id: None,
3184 stage: None,
3185 fired_guards: Vec::new(),
3186 parent_ticket: None,
3187 };
3188 state
3189 .insert_hot_history_record(&DecisionArchiveRecord::Trace {
3190 trace: trace.clone(),
3191 })
3192 .expect("insert fixture trace history");
3193 assert_eq!(state.traces.push(trace), trace_id.get());
3194 if accepted_attempt {
3195 state
3196 .append_attempt(DecisionAttemptRecord {
3197 request_id: canwu_core::DecisionRequestId::new(92),
3198 request_commitment: "a".repeat(64),
3199 at: SimTime::from_minutes(1),
3200 revision_before: 1,
3201 expected_revision: 1,
3202 outcome: DecisionAttemptOutcome::Accepted {
3203 trace_id: Some(trace_id),
3204 command_request_id: None,
3205 },
3206 })
3207 .expect("append fixture attempt");
3208 }
3209 state.validate().expect("fixture state validates");
3210 state
3211 }
3212
3213 fn restart_with_exact_dependency_pages(state: &DecisionState) -> DecisionState {
3214 let hot_bytes = serde_json::to_vec(&state.paged_checkpoint_hot_state())
3215 .expect("encode paged hot state");
3216 let hot: DecisionState =
3217 serde_json::from_slice(&hot_bytes).expect("decode paged hot state");
3218 let required_pages = hot
3219 .required_archived_dependency_page_keys()
3220 .expect("derive exact dependency pages");
3221 assert!(!required_pages.is_empty());
3222 let resident_pages = required_pages
3223 .iter()
3224 .map(|page_key| {
3225 state
3226 .decision_archive_bucket_page(*page_key)
3227 .expect("encode dependency page")
3228 .expect("dependency page is resident")
3229 })
3230 .collect::<Vec<_>>();
3231 let restarted = DecisionState::from_paged_checkpoint_root_with_resident_pages(
3232 hot,
3233 state.archive_bucket_page_ids.clone(),
3234 state.archive_receipt_count,
3235 &state.archive_receipt_root().expect("archive root"),
3236 resident_pages,
3237 )
3238 .expect("restart with exact dependency pages");
3239 let encoded = serde_json::to_vec(&restarted).expect("encode sparse restarted state");
3240 serde_json::from_slice(&encoded).expect("sparse restarted state remains restartable")
3241 }
3242
3243 #[test]
3244 fn paged_checkpoint_root_rejects_hot_archive_metadata() {
3245 let hot = DecisionState {
3246 archive_receipt_count: 1,
3247 ..DecisionState::default()
3248 };
3249 let error = DecisionState::from_paged_checkpoint_root(
3250 hot,
3251 OrdMap::new(),
3252 0,
3253 &decision_archive_hash(
3254 "canwu.decision.archive-receipts.v3",
3255 &(0_u64, OrdMap::<DecisionArchivePageKey, String>::new()),
3256 )
3257 .expect("empty archive root"),
3258 )
3259 .expect_err("hot archive metadata must be rejected");
3260 assert_eq!(error.code, DecisionErrorCode::InvalidDecision);
3261 }
3262
3263 fn append_terminal_attempt(state: &mut DecisionState, ordinal: u64) {
3264 state
3265 .append_attempt(DecisionAttemptRecord {
3266 request_id: canwu_core::DecisionRequestId::new(ordinal),
3267 request_commitment: blake3::hash(&ordinal.to_be_bytes()).to_hex().to_string(),
3268 at: SimTime::from_minutes(i64::try_from(ordinal).expect("test time fits i64")),
3269 revision_before: ordinal.saturating_sub(1),
3270 expected_revision: ordinal.saturating_sub(1),
3271 outcome: DecisionAttemptOutcome::Rejected {
3272 code: crate::DecisionAttemptErrorCode::InvalidDecision,
3273 message: "terminal test attempt".to_owned(),
3274 },
3275 })
3276 .expect("append terminal attempt");
3277 }
3278
3279 fn archive_keys(
3280 state: &DecisionState,
3281 archive: &DecisionScaleArchive,
3282 keys: &[DecisionHistoryKey],
3283 ) -> (DecisionState, VerifiedDecisionArchiveCommit) {
3284 let prepared = state
3285 .prepare_decision_archive(keys)
3286 .expect("prepare archive");
3287 for blob in &prepared.blobs {
3288 archive.store_decision_archive(blob).expect("store blob");
3289 }
3290 let verified = state
3291 .verify_decision_archive(&prepared, archive)
3292 .expect("verify archive");
3293 let next = state
3294 .commit_verified_decision_archive(&verified)
3295 .expect("commit archive");
3296 for page_key in keys
3297 .iter()
3298 .map(decision_history_page_key)
3299 .collect::<Result<OrdSet<_>, _>>()
3300 .expect("derive touched pages")
3301 {
3302 let page = next
3303 .decision_archive_bucket_page(page_key)
3304 .expect("build page")
3305 .expect("page is resident");
3306 archive
3307 .pages
3308 .borrow_mut()
3309 .insert(page.state_page_id().expect("page id"), page);
3310 }
3311 (next, verified)
3312 }
3313
3314 #[test]
3315 fn root_only_restart_can_archive_again_into_the_same_segment() {
3316 let first = 1_u64;
3317 let first_key = DecisionHistoryKey::Attempt(canwu_core::DecisionRequestId::new(first));
3318 let page_key = decision_history_page_key(&first_key).expect("first page key");
3319 let second = (2_u64..1_000_000)
3320 .find(|ordinal| {
3321 decision_history_page_key(&DecisionHistoryKey::Attempt(
3322 canwu_core::DecisionRequestId::new(*ordinal),
3323 ))
3324 .ok()
3325 == Some(page_key)
3326 })
3327 .expect("find same-segment request id");
3328 let second_key = DecisionHistoryKey::Attempt(canwu_core::DecisionRequestId::new(second));
3329 let archive = DecisionScaleArchive::default();
3330 let mut state = DecisionState::default();
3331 append_terminal_attempt(&mut state, first);
3332 let (state, _) = archive_keys(&state, &archive, std::slice::from_ref(&first_key));
3333 let archive_root = state.archive_receipt_root().expect("archive root");
3334 let hot = state.paged_checkpoint_hot_state();
3335 assert!(hot.archive_receipt_buckets.is_empty());
3336 assert!(hot.archive_bucket_page_ids.is_empty());
3337 assert_eq!(hot.archive_receipt_count, 0);
3338 let mut restarted = DecisionState::from_paged_checkpoint_root(
3339 hot,
3340 state.archive_bucket_page_ids.clone(),
3341 state.archive_receipt_count,
3342 &archive_root,
3343 )
3344 .expect("root-only restart");
3345 append_terminal_attempt(&mut restarted, second);
3346 let (restarted, verified) =
3347 archive_keys(&restarted, &archive, std::slice::from_ref(&second_key));
3348
3349 assert_eq!(restarted.archived_history_count(), 2);
3350 assert!(
3351 restarted
3352 .load_decision_history(&first_key, &archive)
3353 .expect("load first")
3354 .is_some()
3355 );
3356 assert!(
3357 restarted
3358 .load_decision_history(&second_key, &archive)
3359 .expect("load second")
3360 .is_some()
3361 );
3362 assert_eq!(
3363 restarted
3364 .archive_reachability(&archive)
3365 .expect("enumerate reachability")
3366 .blob_locators
3367 .len(),
3368 2
3369 );
3370 assert_eq!(
3371 restarted
3372 .commit_verified_decision_archive(&verified)
3373 .expect("replay commit"),
3374 restarted
3375 );
3376 }
3377
3378 #[test]
3379 fn paged_restart_loads_archived_ticket_referenced_by_hot_trace() {
3380 let archive = DecisionScaleArchive::default();
3381 let state = linked_decision_state(true, false);
3382 let ticket_id = DecisionTicketId::new(91);
3383 let (archived, _) =
3384 archive_keys(&state, &archive, &[DecisionHistoryKey::Ticket(ticket_id)]);
3385 let restarted = restart_with_exact_dependency_pages(&archived);
3386
3387 assert!(restarted.trace(DecisionTraceId::new(1)).is_some());
3388 assert!(matches!(
3389 restarted.decision_locator(&DecisionHistoryKey::Ticket(ticket_id)),
3390 DecisionHistoryLocation::Archived { .. }
3391 ));
3392 restarted
3393 .validate()
3394 .expect("hot trace dependency validates");
3395 }
3396
3397 #[test]
3398 fn paged_restart_loads_archived_trace_referenced_by_resolved_hot_ticket() {
3399 let archive = DecisionScaleArchive::default();
3400 let state = linked_decision_state(true, false);
3401 let trace_id = DecisionTraceId::new(1);
3402 let (archived, _) = archive_keys(&state, &archive, &[DecisionHistoryKey::Trace(trace_id)]);
3403 let restarted = restart_with_exact_dependency_pages(&archived);
3404
3405 assert!(restarted.ticket(DecisionTicketId::new(91)).is_some());
3406 assert!(matches!(
3407 restarted.decision_locator(&DecisionHistoryKey::Trace(trace_id)),
3408 DecisionHistoryLocation::Archived { .. }
3409 ));
3410 restarted
3411 .validate()
3412 .expect("resolved hot ticket dependency validates");
3413 }
3414
3415 #[test]
3416 fn paged_restart_loads_archived_trace_referenced_by_accepted_hot_attempt() {
3417 let archive = DecisionScaleArchive::default();
3418 let state = linked_decision_state(false, true);
3419 let trace_id = DecisionTraceId::new(1);
3420 let (archived, _) = archive_keys(&state, &archive, &[DecisionHistoryKey::Trace(trace_id)]);
3421 let restarted = restart_with_exact_dependency_pages(&archived);
3422
3423 assert!(
3424 restarted
3425 .attempt(canwu_core::DecisionRequestId::new(92))
3426 .is_some()
3427 );
3428 assert!(matches!(
3429 restarted.decision_locator(&DecisionHistoryKey::Trace(trace_id)),
3430 DecisionHistoryLocation::Archived { .. }
3431 ));
3432 restarted
3433 .validate()
3434 .expect("accepted hot attempt dependency validates");
3435 }
3436}
3437
3438#[derive(Clone, Debug, Default, Eq, PartialEq)]
3439pub struct PreparedDecision {
3440 pub trace: Option<DecisionTrace>,
3441 pub action: Option<DecisionAction>,
3442}
3443
3444#[derive(Clone, Debug, Eq, PartialEq)]
3445pub enum ControllerDecision {
3446 Authoritative {
3447 decision: PolicyDecision,
3448 action: Option<DecisionAction>,
3449 },
3450 Pending(PolicyDecision),
3451}
3452
3453pub struct DecisionController;
3454
3455impl DecisionController {
3456 pub fn evaluate(
3457 ticket: &DecisionTicket,
3458 controller: &DecisionControllerBinding,
3459 policy: &dyn DecisionPolicy,
3460 ) -> Result<ControllerDecision, DecisionError> {
3461 if !ticket.is_open() {
3462 return Err(DecisionError::new(
3463 DecisionErrorCode::ClosedTicket,
3464 "only open tickets can be evaluated",
3465 ));
3466 }
3467 if ticket.assigned_controller != controller.id || policy.identity() != &controller.policy {
3468 return Err(DecisionError::new(
3469 DecisionErrorCode::PolicyMismatch,
3470 "runtime policy identity does not match the ticket controller binding",
3471 ));
3472 }
3473 let decision = policy.decide(ticket)?;
3474 decision.validate(ticket)?;
3475 if decision.random.is_some() {
3476 return Err(DecisionError::new(
3477 DecisionErrorCode::PolicyMismatch,
3478 "policies cannot supply random draw evidence; draws come from a boundary resolution",
3479 ));
3480 }
3481 if matches!(decision.outcome, DecisionOutcome::PendingRandom { .. })
3482 && !controller.random_tie_break
3483 {
3484 return Err(DecisionError::new(
3485 DecisionErrorCode::PolicyMismatch,
3486 "the controller binding does not permit random tie-breaks",
3487 ));
3488 }
3489 if decision.outcome.is_pending() {
3490 return Ok(ControllerDecision::Pending(decision));
3491 }
3492 let action = match &decision.outcome {
3493 DecisionOutcome::Selected { option_id } => {
3494 ticket.option(option_id).map(|option| option.action.clone())
3495 }
3496 DecisionOutcome::Deferred { .. }
3497 | DecisionOutcome::Pending { .. }
3498 | DecisionOutcome::PendingRandom { .. } => None,
3499 };
3500 Ok(ControllerDecision::Authoritative { action, decision })
3501 }
3502}