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 }
1220 for attempt in self.attempts.iter() {
1221 if let DecisionAttemptOutcome::Accepted {
1222 trace_id: Some(trace_id),
1223 ..
1224 } = &attempt.outcome
1225 && self.traces.get_ordinal(trace_id.get()).is_none()
1226 {
1227 pages.insert(decision_history_page_key(&DecisionHistoryKey::Trace(
1228 *trace_id,
1229 ))?);
1230 }
1231 }
1232 Ok(pages)
1233 }
1234
1235 pub fn paged_checkpoint_parts(
1239 &self,
1240 ) -> Result<
1241 (
1242 Self,
1243 std::collections::BTreeMap<DecisionArchivePageKey, Vec<DecisionArchiveReceipt>>,
1244 ),
1245 DecisionError,
1246 > {
1247 let hot = self.paged_checkpoint_hot_state();
1248 let mut buckets =
1249 std::collections::BTreeMap::<DecisionArchivePageKey, Vec<DecisionArchiveReceipt>>::new(
1250 );
1251 for (bucket, receipts) in &self.archive_receipt_buckets {
1252 let decoded = receipts
1253 .iter()
1254 .map(|(key, receipt)| receipt.to_receipt(key))
1255 .collect();
1256 buckets.insert(*bucket, decoded);
1257 }
1258 Ok((hot, buckets))
1259 }
1260
1261 pub fn from_paged_checkpoint_parts(
1262 mut hot: Self,
1263 buckets: impl IntoIterator<Item = (DecisionArchivePageKey, Vec<DecisionArchiveReceipt>)>,
1264 ) -> Result<Self, DecisionError> {
1265 if !hot.archive_receipt_buckets.is_empty() {
1266 return Err(archive_error(
1267 "paged decision hot state must not duplicate archive receipts",
1268 ));
1269 }
1270 hot.archive_bucket_page_ids.clear();
1271 hot.archive_receipt_count = 0;
1272 for (bucket, receipts) in buckets {
1273 if bucket.bucket >= DECISION_HISTORY_BUCKET_COUNT
1274 || receipts.windows(2).any(|pair| pair[0].key >= pair[1].key)
1275 {
1276 return Err(archive_error(
1277 "paged decision archive bucket is out of range or not strictly ordered",
1278 ));
1279 }
1280 for receipt in receipts {
1281 if decision_history_page_key(&receipt.key)? != bucket
1282 || hot.insert_archive_receipt(&receipt)?.is_some()
1283 {
1284 return Err(archive_error(
1285 "paged decision archive bucket contains a misplaced or duplicate receipt",
1286 ));
1287 }
1288 }
1289 }
1290 hot.validate()?;
1291 Ok(hot)
1292 }
1293
1294 pub fn from_paged_checkpoint_root(
1297 hot: Self,
1298 archive_bucket_page_ids: OrdMap<DecisionArchivePageKey, String>,
1299 archive_receipt_count: u64,
1300 archive_receipt_root: &str,
1301 ) -> Result<Self, DecisionError> {
1302 Self::from_paged_checkpoint_root_with_resident_pages(
1303 hot,
1304 archive_bucket_page_ids,
1305 archive_receipt_count,
1306 archive_receipt_root,
1307 std::iter::empty(),
1308 )
1309 }
1310
1311 pub fn from_paged_checkpoint_root_with_resident_pages(
1316 mut hot: Self,
1317 archive_bucket_page_ids: OrdMap<DecisionArchivePageKey, String>,
1318 archive_receipt_count: u64,
1319 archive_receipt_root: &str,
1320 resident_pages: impl IntoIterator<Item = DecisionArchiveBucketPage>,
1321 ) -> Result<Self, DecisionError> {
1322 if !hot.archive_receipt_buckets.is_empty()
1323 || !hot.archive_bucket_page_ids.is_empty()
1324 || hot.archive_receipt_count != 0
1325 {
1326 return Err(archive_error(
1327 "paged decision hot state must not contain archive receipts or archive metadata",
1328 ));
1329 }
1330 hot.archive_bucket_page_ids = archive_bucket_page_ids;
1331 hot.archive_receipt_count = archive_receipt_count;
1332 for page in resident_pages {
1333 page.validate()?;
1334 let page_key = DecisionArchivePageKey {
1335 bucket: page.bucket,
1336 segment: page.segment,
1337 };
1338 let page_id = page.state_page_id()?;
1339 if hot.archive_bucket_page_ids.get(&page_key) != Some(&page_id)
1340 || hot.archive_receipt_buckets.contains_key(&page_key)
1341 {
1342 return Err(archive_error(
1343 "resident decision archive page is absent from or disagrees with the committed directory",
1344 ));
1345 }
1346 let mut receipts = OrdMap::new();
1347 for receipt in page.receipts {
1348 if receipts
1349 .insert(
1350 receipt.key.clone(),
1351 CompactDecisionArchiveReceipt::from_receipt(&receipt)?,
1352 )
1353 .is_some()
1354 {
1355 return Err(archive_error(
1356 "resident decision archive page contains duplicate receipts",
1357 ));
1358 }
1359 }
1360 hot.archive_receipt_buckets.insert(page_key, receipts);
1361 }
1362 hot.validate()?;
1363 if hot.archive_receipt_root()? != archive_receipt_root {
1364 return Err(archive_error(
1365 "paged decision archive directory root is inconsistent",
1366 ));
1367 }
1368 Ok(hot)
1369 }
1370
1371 pub fn archive_receipt_commitment(&self) -> Result<String, DecisionError> {
1372 self.archive_receipt_root()
1373 }
1374
1375 pub fn authoritative_commitment(&self) -> Result<String, DecisionError> {
1378 decision_archive_hash(
1379 "canwu.commitment.decisions.v2",
1380 &(
1381 &self.controllers,
1382 self.hot_history_root()?,
1383 self.archive_receipt_root()?,
1384 ),
1385 )
1386 }
1387
1388 pub fn archived_decision_history_page(
1389 &self,
1390 cursor: Option<&DecisionHistoryCursor>,
1391 budget: DecisionHistoryQueryBudget,
1392 provider: &dyn DecisionArchiveProvider,
1393 ) -> Result<DecisionHistoryPage, DecisionError> {
1394 if budget.max_results == 0
1395 || budget.max_results > MAX_DECISION_HISTORY_PAGE_SIZE
1396 || budget.max_provider_calls == 0
1397 || budget.max_decoded_bytes == 0
1398 || budget.max_decoded_bytes > MAX_DECISION_HISTORY_PAGE_BYTES
1399 {
1400 return Err(history_budget_error(
1401 "decision history page budget is zero, inconsistent, or exceeds the hard limit",
1402 ));
1403 }
1404 let archive_root = self.archive_receipt_root()?;
1405 if cursor.is_some_and(|cursor| cursor.archive_root != archive_root) {
1406 return Err(history_unavailable(
1407 "decision history cursor belongs to a different archive generation",
1408 ));
1409 }
1410 if cursor.is_some_and(|cursor| {
1411 cursor.bucket >= DECISION_HISTORY_BUCKET_COUNT
1412 || decision_history_page_key(&cursor.after).ok()
1413 != Some(DecisionArchivePageKey {
1414 bucket: cursor.bucket,
1415 segment: cursor.segment,
1416 })
1417 }) {
1418 return Err(history_unavailable(
1419 "decision history cursor has an invalid bucket position",
1420 ));
1421 }
1422 let start_bucket = cursor.map_or(
1423 DecisionArchivePageKey {
1424 bucket: 0,
1425 segment: 0,
1426 },
1427 |cursor| DecisionArchivePageKey {
1428 bucket: cursor.bucket,
1429 segment: cursor.segment,
1430 },
1431 );
1432 let mut selected = Vec::<(DecisionArchivePageKey, DecisionArchiveReceipt)>::new();
1433 let mut provider_calls = 0_usize;
1434 for (bucket, page_id) in self.archive_bucket_page_ids.range(start_bucket..) {
1435 let loaded;
1436 let receipts = if let Some(receipts) = self.archive_receipt_buckets.get(bucket) {
1437 receipts
1438 .iter()
1439 .map(|(key, receipt)| receipt.to_receipt(key))
1440 .collect::<Vec<_>>()
1441 } else {
1442 provider_calls = provider_calls.checked_add(1).ok_or_else(|| {
1443 history_budget_error("decision history provider-call count overflowed")
1444 })?;
1445 if provider_calls > budget.max_provider_calls {
1446 return Err(history_budget_error(
1447 "decision history page exceeded its provider-call budget",
1448 ));
1449 }
1450 loaded = provider
1451 .load_decision_archive_bucket_page(page_id)?
1452 .ok_or_else(|| {
1453 history_unavailable("decision locator bucket page is unavailable")
1454 })?;
1455 loaded.validate()?;
1456 if loaded.bucket != bucket.bucket
1457 || loaded.segment != bucket.segment
1458 || loaded.state_page_id()? != *page_id
1459 {
1460 return Err(history_unavailable(
1461 "decision locator provider returned a mismatched bucket page",
1462 ));
1463 }
1464 loaded.receipts
1465 };
1466 if let Some(cursor) = cursor
1467 .filter(|cursor| cursor.bucket == bucket.bucket && cursor.segment == bucket.segment)
1468 {
1469 let after = &cursor.after;
1470 for receipt in receipts.iter().filter(|receipt| receipt.key > *after) {
1471 selected.push((*bucket, receipt.clone()));
1472 if selected.len() > budget.max_results {
1473 break;
1474 }
1475 }
1476 } else {
1477 for receipt in receipts {
1478 selected.push((*bucket, receipt));
1479 if selected.len() > budget.max_results {
1480 break;
1481 }
1482 }
1483 }
1484 if selected.len() > budget.max_results {
1485 break;
1486 }
1487 }
1488 let has_more = selected.len() > budget.max_results;
1489 let selected = &selected[..selected.len().min(budget.max_results)];
1490 let mut records = Vec::with_capacity(selected.len());
1491 let mut decoded_bytes = 0_u64;
1492 for (_, receipt) in selected {
1493 provider_calls = provider_calls.checked_add(1).ok_or_else(|| {
1494 history_budget_error("decision history provider-call count overflowed")
1495 })?;
1496 if provider_calls > budget.max_provider_calls {
1497 return Err(history_budget_error(
1498 "decision history page exceeded its provider-call budget",
1499 ));
1500 }
1501 let blob = provider
1502 .load_decision_archive(&receipt.locator)?
1503 .ok_or_else(|| {
1504 history_unavailable("decision archive provider omitted a committed page member")
1505 })?;
1506 blob.validate()?;
1507 if blob.key != receipt.key || blob.content_id()? != receipt.content_id {
1508 return Err(history_unavailable(
1509 "decision archive provider returned a mismatched page member",
1510 ));
1511 }
1512 decoded_bytes = decoded_bytes
1513 .checked_add(receipt.encoded_bytes)
1514 .ok_or_else(|| history_budget_error("decision history byte count overflowed"))?;
1515 if decoded_bytes > budget.max_decoded_bytes {
1516 return Err(history_budget_error(
1517 "decision history page exceeded its decoded-byte budget",
1518 ));
1519 }
1520 records.push(blob.record);
1521 }
1522 let next_cursor = if has_more {
1523 selected
1524 .last()
1525 .map(|(bucket, receipt)| DecisionHistoryCursor {
1526 archive_root: archive_root.clone(),
1527 bucket: bucket.bucket,
1528 segment: bucket.segment,
1529 after: receipt.key.clone(),
1530 })
1531 } else {
1532 None
1533 };
1534 Ok(DecisionHistoryPage {
1535 archive_root,
1536 records,
1537 next_cursor,
1538 provider_calls: provider_calls as u64,
1539 decoded_bytes,
1540 })
1541 }
1542
1543 pub fn archive_reachability(
1547 &self,
1548 provider: &dyn DecisionArchiveProvider,
1549 ) -> Result<DecisionArchiveReachability, DecisionError> {
1550 let mut reachable = DecisionArchiveReachability::default();
1551 let mut receipt_count = 0_u64;
1552 for (bucket, page_id) in &self.archive_bucket_page_ids {
1553 let page = provider
1554 .load_decision_archive_bucket_page(page_id)?
1555 .ok_or_else(|| {
1556 history_unavailable("decision archive bucket page is unavailable")
1557 })?;
1558 page.validate()?;
1559 if page.bucket != bucket.bucket
1560 || page.segment != bucket.segment
1561 || page.state_page_id()? != *page_id
1562 {
1563 return Err(history_unavailable(
1564 "decision archive reachability provider returned a mismatched page",
1565 ));
1566 }
1567 reachable.bucket_page_ids.insert(page_id.clone());
1568 for receipt in page.receipts {
1569 receipt_count = receipt_count
1570 .checked_add(1)
1571 .ok_or_else(|| archive_error("decision archive receipt count overflowed"))?;
1572 reachable.blob_locators.insert(receipt.locator);
1573 }
1574 }
1575 if receipt_count != self.archive_receipt_count {
1576 return Err(archive_error(
1577 "decision archive reachability count disagrees with its directory",
1578 ));
1579 }
1580 Ok(reachable)
1581 }
1582
1583 pub fn load_decision_history(
1584 &self,
1585 key: &DecisionHistoryKey,
1586 provider: &dyn DecisionArchiveProvider,
1587 ) -> Result<Option<DecisionArchiveRecord>, DecisionError> {
1588 if let Some(record) = self.hot_archive_record(key) {
1589 return Ok(Some(record));
1590 }
1591 let Some(receipt) = self.archive_receipt_with_provider(key, provider)? else {
1592 return Ok(None);
1593 };
1594 let blob = provider
1595 .load_decision_archive(&receipt.locator)?
1596 .ok_or_else(|| archive_error("decision archive provider omitted a committed blob"))?;
1597 blob.validate()?;
1598 let content_id = blob.content_id()?;
1599 if blob.key != *key || content_id != receipt.content_id {
1600 return Err(archive_error(
1601 "decision archive provider returned content that disagrees with its receipt",
1602 ));
1603 }
1604 Ok(Some(blob.record))
1605 }
1606
1607 pub fn prepare_decision_archive(
1608 &self,
1609 keys: &[DecisionHistoryKey],
1610 ) -> Result<PreparedDecisionArchive, DecisionError> {
1611 if keys.is_empty() || keys.len() > MAX_DECISION_ARCHIVE_BATCH_ENTRIES {
1612 return Err(archive_error(
1613 "decision archive batch is empty or exceeds the bounded entry limit",
1614 ));
1615 }
1616 let mut canonical_keys = keys.to_vec();
1617 canonical_keys.sort();
1618 canonical_keys.dedup();
1619 if canonical_keys.len() != keys.len() {
1620 return Err(archive_error(
1621 "decision archive batch contains duplicate keys",
1622 ));
1623 }
1624 let source_root = self.hot_history_root()?;
1625 let mut blobs = Vec::with_capacity(canonical_keys.len());
1626 let mut receipts = Vec::with_capacity(canonical_keys.len());
1627 for key in canonical_keys {
1628 if self.contains_archived_key(&key) {
1629 return Err(archive_error("decision history is already archived"));
1630 }
1631 let record = self
1632 .hot_archive_record(&key)
1633 .ok_or_else(|| archive_error("decision archive key is not resident"))?;
1634 self.validate_archive_eligibility(&record)?;
1635 let blob = DecisionArchiveBlob {
1636 format_version: DECISION_ARCHIVE_FORMAT_VERSION,
1637 key: key.clone(),
1638 record,
1639 };
1640 let content_id = blob.content_id()?;
1641 let encoded_bytes = u64::try_from(
1642 serde_json::to_vec(&blob)
1643 .map_err(|error| {
1644 archive_error(format!("cannot encode decision archive: {error}"))
1645 })?
1646 .len(),
1647 )
1648 .map_err(|_| {
1649 archive_error("decision archive blob exceeds the persistent byte range")
1650 })?;
1651 receipts.push(DecisionArchiveReceipt {
1652 format_version: DECISION_ARCHIVE_FORMAT_VERSION,
1653 key,
1654 locator: content_id.clone(),
1655 content_id,
1656 encoded_bytes,
1657 });
1658 blobs.push(blob);
1659 }
1660 let token = decision_archive_hash(
1661 "canwu.decision.archive-token.v1",
1662 &(&source_root, &receipts),
1663 )?;
1664 Ok(PreparedDecisionArchive {
1665 source_root,
1666 token,
1667 blobs,
1668 receipts,
1669 })
1670 }
1671
1672 pub fn commit_decision_archive(
1673 &self,
1674 prepared: &PreparedDecisionArchive,
1675 provider: &dyn DecisionArchiveProvider,
1676 ) -> Result<Self, DecisionError> {
1677 let verified = self.verify_decision_archive(prepared, provider)?;
1678 self.commit_verified_decision_archive(&verified)
1679 }
1680
1681 pub fn verify_decision_archive(
1682 &self,
1683 prepared: &PreparedDecisionArchive,
1684 provider: &dyn DecisionArchiveProvider,
1685 ) -> Result<VerifiedDecisionArchiveCommit, DecisionError> {
1686 let keys = prepared
1687 .receipts
1688 .iter()
1689 .map(|receipt| receipt.key.clone())
1690 .collect::<Vec<_>>();
1691 let expected = self.prepare_decision_archive(&keys)?;
1692 if expected != *prepared {
1693 return Err(archive_error(
1694 "prepared decision archive is stale or altered",
1695 ));
1696 }
1697 for (blob, receipt) in prepared.blobs.iter().zip(&prepared.receipts) {
1698 let stored = provider
1699 .load_decision_archive(&receipt.locator)?
1700 .ok_or_else(|| {
1701 archive_error("prepared decision archive is not durably readable")
1702 })?;
1703 if stored != *blob || stored.content_id()? != receipt.content_id {
1704 return Err(archive_error(
1705 "stored decision archive blob failed verification",
1706 ));
1707 }
1708 }
1709 let mut additions = std::collections::BTreeMap::<
1710 DecisionArchivePageKey,
1711 Vec<&DecisionArchiveReceipt>,
1712 >::new();
1713 for receipt in &prepared.receipts {
1714 additions
1715 .entry(decision_history_page_key(&receipt.key)?)
1716 .or_default()
1717 .push(receipt);
1718 }
1719 let mut page_replacements = Vec::with_capacity(additions.len());
1720 for (page_key, receipts) in additions {
1721 let previous_page_id = self.archive_bucket_page_ids.get(&page_key).cloned();
1722 let mut merged = if let Some(resident) = self.archive_receipt_buckets.get(&page_key) {
1723 resident
1724 .iter()
1725 .map(|(key, receipt)| (key.clone(), receipt.to_receipt(key)))
1726 .collect::<std::collections::BTreeMap<_, _>>()
1727 } else if let Some(page_id) = previous_page_id.as_deref() {
1728 let page = provider
1729 .load_decision_archive_bucket_page(page_id)?
1730 .ok_or_else(|| {
1731 history_unavailable(
1732 "decision archive locator page is unavailable for authenticated merge",
1733 )
1734 })?;
1735 page.validate()?;
1736 if page.bucket != page_key.bucket
1737 || page.segment != page_key.segment
1738 || page.state_page_id()? != page_id
1739 {
1740 return Err(history_unavailable(
1741 "decision archive locator provider returned a mismatched merge page",
1742 ));
1743 }
1744 page.receipts
1745 .into_iter()
1746 .map(|receipt| (receipt.key.clone(), receipt))
1747 .collect::<std::collections::BTreeMap<_, _>>()
1748 } else {
1749 std::collections::BTreeMap::new()
1750 };
1751 for receipt in receipts {
1752 if merged
1753 .insert(receipt.key.clone(), receipt.clone())
1754 .is_some()
1755 {
1756 return Err(archive_error(
1757 "decision archive locator merge would replace an existing receipt",
1758 ));
1759 }
1760 }
1761 let page = DecisionArchiveBucketPage {
1762 format_version: DECISION_ARCHIVE_BUCKET_PAGE_FORMAT_VERSION,
1763 bucket: page_key.bucket,
1764 segment: page_key.segment,
1765 receipts: merged.into_values().collect(),
1766 };
1767 page.validate()?;
1768 page_replacements.push(VerifiedDecisionArchivePageReplacement {
1769 page_key,
1770 previous_page_id,
1771 page,
1772 });
1773 }
1774 Ok(VerifiedDecisionArchiveCommit {
1775 format_version: DECISION_ARCHIVE_FORMAT_VERSION,
1776 source_root: prepared.source_root.clone(),
1777 token: prepared.token.clone(),
1778 receipts: prepared.receipts.clone(),
1779 page_replacements,
1780 })
1781 }
1782
1783 pub fn commit_verified_decision_archive(
1784 &self,
1785 verified: &VerifiedDecisionArchiveCommit,
1786 ) -> Result<Self, DecisionError> {
1787 if verified.format_version != DECISION_ARCHIVE_FORMAT_VERSION {
1788 return Err(archive_error(
1789 "verified decision archive uses an unsupported format",
1790 ));
1791 }
1792 let mut replacement_ids = std::collections::BTreeMap::new();
1793 for replacement in &verified.page_replacements {
1794 replacement.page.validate()?;
1795 if replacement.page_key.bucket != replacement.page.bucket
1796 || replacement.page_key.segment != replacement.page.segment
1797 || replacement_ids
1798 .insert(replacement.page_key, replacement.page.state_page_id()?)
1799 .is_some()
1800 || replacement
1801 .previous_page_id
1802 .as_ref()
1803 .is_some_and(|page_id| !canonical_archive_hash(page_id))
1804 {
1805 return Err(archive_error(
1806 "verified decision archive page replacement is malformed",
1807 ));
1808 }
1809 }
1810 let mut receipts_by_page = std::collections::BTreeMap::<
1811 DecisionArchivePageKey,
1812 Vec<&DecisionArchiveReceipt>,
1813 >::new();
1814 for receipt in &verified.receipts {
1815 receipts_by_page
1816 .entry(decision_history_page_key(&receipt.key)?)
1817 .or_default()
1818 .push(receipt);
1819 }
1820 if receipts_by_page.len() != replacement_ids.len()
1821 || receipts_by_page
1822 .keys()
1823 .any(|page_key| !replacement_ids.contains_key(page_key))
1824 {
1825 return Err(archive_error(
1826 "verified decision archive does not replace every touched locator page",
1827 ));
1828 }
1829 let already_applied = verified.receipts.iter().all(|receipt| {
1830 !matches!(
1831 self.decision_locator(&receipt.key),
1832 DecisionHistoryLocation::Hot
1833 )
1834 }) && replacement_ids
1835 .iter()
1836 .all(|(page_key, page_id)| self.archive_bucket_page_ids.get(page_key) == Some(page_id));
1837 if already_applied {
1838 return Ok(self.clone());
1839 }
1840 for replacement in &verified.page_replacements {
1841 if self.archive_bucket_page_ids.get(&replacement.page_key)
1842 != replacement.previous_page_id.as_ref()
1843 {
1844 return Err(archive_error(
1845 "verified decision archive locator base changed before commit",
1846 ));
1847 }
1848 let receipts = receipts_by_page.get(&replacement.page_key).ok_or_else(|| {
1849 archive_error("verified decision archive locator page has no touched receipts")
1850 })?;
1851 for receipt in receipts {
1852 if replacement
1853 .page
1854 .receipts
1855 .binary_search_by(|candidate| candidate.key.cmp(&receipt.key))
1856 .ok()
1857 .is_none_or(|index| replacement.page.receipts[index] != **receipt)
1858 {
1859 return Err(archive_error(
1860 "verified decision archive locator page omits a touched receipt",
1861 ));
1862 }
1863 }
1864 }
1865 let keys = verified
1866 .receipts
1867 .iter()
1868 .map(|receipt| receipt.key.clone())
1869 .collect::<Vec<_>>();
1870 let expected = self.prepare_decision_archive(&keys)?;
1871 if expected.source_root != verified.source_root
1872 || expected.token != verified.token
1873 || expected.receipts != verified.receipts
1874 {
1875 return Err(archive_error(
1876 "verified decision archive is stale or altered",
1877 ));
1878 }
1879 let mut next = self.clone();
1880 for receipt in &verified.receipts {
1881 let hot_record = next
1882 .hot_archive_record(&receipt.key)
1883 .ok_or_else(|| archive_error("decision history disappeared during commit"))?;
1884 next.remove_hot_history_record(&hot_record)?;
1885 match receipt.key {
1886 DecisionHistoryKey::Ticket(id) => {
1887 next.tickets.remove(&id);
1888 }
1889 DecisionHistoryKey::Attempt(id) => {
1890 let ordinal = next.attempts_by_request.remove(&id).ok_or_else(|| {
1891 archive_error("decision attempt disappeared during commit")
1892 })?;
1893 next.attempts.remove_ordinal(ordinal);
1894 }
1895 DecisionHistoryKey::Trace(id) => {
1896 next.traces.remove_ordinal(id.get());
1897 }
1898 }
1899 }
1900 next.archive_receipt_count = next
1901 .archive_receipt_count
1902 .checked_add(verified.receipts.len() as u64)
1903 .ok_or_else(|| archive_error("decision archive receipt count overflowed"))?;
1904 for replacement in &verified.page_replacements {
1905 let receipts = replacement
1906 .page
1907 .receipts
1908 .iter()
1909 .map(|receipt| {
1910 Ok((
1911 receipt.key.clone(),
1912 CompactDecisionArchiveReceipt::from_receipt(receipt)?,
1913 ))
1914 })
1915 .collect::<Result<OrdMap<_, _>, DecisionError>>()?;
1916 next.archive_receipt_buckets
1917 .insert(replacement.page_key, receipts);
1918 next.archive_bucket_page_ids
1919 .insert(replacement.page_key, replacement.page.state_page_id()?);
1920 }
1921 next.validate_archive_directory_shape()?;
1922 for receipt in &verified.receipts {
1923 if matches!(
1924 next.decision_locator(&receipt.key),
1925 DecisionHistoryLocation::Hot
1926 ) {
1927 return Err(archive_error(
1928 "decision archive commit left a touched key in hot state",
1929 ));
1930 }
1931 }
1932 Ok(next)
1933 }
1934
1935 fn hot_archive_record(&self, key: &DecisionHistoryKey) -> Option<DecisionArchiveRecord> {
1936 match key {
1937 DecisionHistoryKey::Ticket(id) => self
1938 .tickets
1939 .get(id)
1940 .map(|ticket| ticket.as_ref().clone())
1941 .map(|ticket| DecisionArchiveRecord::Ticket { ticket }),
1942 DecisionHistoryKey::Attempt(id) => self
1943 .attempt(*id)
1944 .cloned()
1945 .map(|attempt| DecisionArchiveRecord::Attempt { attempt }),
1946 DecisionHistoryKey::Trace(id) => self
1947 .traces
1948 .get_ordinal(id.get())
1949 .cloned()
1950 .map(|trace| DecisionArchiveRecord::Trace { trace }),
1951 }
1952 }
1953
1954 fn validate_archive_eligibility(
1955 &self,
1956 record: &DecisionArchiveRecord,
1957 ) -> Result<(), DecisionError> {
1958 match record {
1959 DecisionArchiveRecord::Ticket { ticket } if ticket.is_open() => Err(archive_error(
1960 "open decision tickets cannot leave the hot mutation path",
1961 )),
1962 DecisionArchiveRecord::Trace { trace } => {
1963 let terminal = self
1964 .tickets
1965 .get(&trace.ticket_id)
1966 .is_some_and(|ticket| !ticket.is_open())
1967 || self.contains_archived_key(&DecisionHistoryKey::Ticket(trace.ticket_id));
1968 if terminal {
1969 Ok(())
1970 } else {
1971 Err(archive_error(
1972 "decision traces can archive only after their ticket is terminal",
1973 ))
1974 }
1975 }
1976 DecisionArchiveRecord::Ticket { .. } | DecisionArchiveRecord::Attempt { .. } => Ok(()),
1977 }
1978 }
1979
1980 fn hot_history_root(&self) -> Result<String, DecisionError> {
1981 let resident_count = (self.tickets.len() as u64)
1982 .checked_add(self.traces.len() as u64)
1983 .and_then(|count| count.checked_add(self.attempts.len() as u64))
1984 .ok_or_else(|| archive_error("decision hot-history count overflowed"))?;
1985 if resident_count != self.hot_history_accumulator.count {
1986 return Err(archive_error(
1987 "decision hot-history accumulator count is inconsistent",
1988 ));
1989 }
1990 decision_archive_hash(
1991 "canwu.decision.hot-history-root.v2",
1992 &(
1993 self.hot_history_accumulator.count,
1994 self.hot_history_accumulator.xor,
1995 self.hot_history_accumulator.sum,
1996 ),
1997 )
1998 }
1999
2000 pub fn hot_history_commitment(&self) -> Result<String, DecisionError> {
2001 self.hot_history_root()
2002 }
2003
2004 fn archive_receipt_root(&self) -> Result<String, DecisionError> {
2005 let mut hasher = blake3::Hasher::new();
2006 hasher.update(b"canwu.decision.archive-receipt-root.v3");
2007 hasher.update(&[0]);
2008 hasher.update(&self.archive_receipt_count.to_be_bytes());
2009 hasher.update(&(self.archive_bucket_page_ids.len() as u64).to_be_bytes());
2010 for (bucket, page_id) in &self.archive_bucket_page_ids {
2011 hasher.update(&bucket.bucket.to_be_bytes());
2012 hasher.update(&[bucket.segment]);
2013 hasher.update(&decode_archive_hash(page_id)?);
2014 }
2015 Ok(hasher.finalize().to_hex().to_string())
2016 }
2017
2018 #[must_use]
2019 pub fn decision_hot_state(&self) -> DecisionHotState {
2020 DecisionHotState {
2021 ticket_count: self.tickets.len() as u64,
2022 attempt_count: self.attempts.len() as u64,
2023 trace_count: self.traces.len() as u64,
2024 }
2025 }
2026
2027 pub fn append_attempt(&mut self, attempt: DecisionAttemptRecord) -> Result<(), DecisionError> {
2028 validate_attempt_shape(self, &attempt)?;
2029 if self.attempts_by_request.contains_key(&attempt.request_id) {
2030 return Err(DecisionError::new(
2031 DecisionErrorCode::InvalidDecision,
2032 "decision attempts must use unique request IDs",
2033 ));
2034 }
2035 let request_id = attempt.request_id;
2036 self.insert_hot_history_record(&DecisionArchiveRecord::Attempt {
2037 attempt: attempt.clone(),
2038 })?;
2039 let ordinal = self.attempts.push(attempt);
2040 self.attempts_by_request.insert(request_id, ordinal);
2041 Ok(())
2042 }
2043
2044 pub fn open_tickets(&self) -> impl Iterator<Item = &DecisionTicket> {
2045 self.tickets
2046 .values()
2047 .map(Arc::as_ref)
2048 .filter(|ticket| ticket.is_open())
2049 }
2050
2051 pub fn validate(&self) -> Result<(), DecisionError> {
2052 self.validate_archive_directory_shape()?;
2053 if self.attempts_by_request.len() != self.attempts.len() {
2054 return Err(DecisionError::new(
2055 DecisionErrorCode::InvalidDecision,
2056 "decision attempt request index is inconsistent",
2057 ));
2058 }
2059 let mut expected_deadlines = OrdMap::<SimTime, OrdSet<DecisionTicketId>>::new();
2060 for ticket in self.tickets.values().filter(|ticket| ticket.is_open()) {
2061 if let Some(deadline) = ticket.deadline {
2062 expected_deadlines
2063 .entry(deadline)
2064 .or_default()
2065 .insert(ticket.id);
2066 }
2067 }
2068 if expected_deadlines != self.deadline_index {
2069 return Err(DecisionError::new(
2070 DecisionErrorCode::InvalidDecision,
2071 "decision deadline index is inconsistent",
2072 ));
2073 }
2074 for (id, controller) in &self.controllers {
2075 if id != &controller.id {
2076 return Err(DecisionError::new(
2077 DecisionErrorCode::InvalidController,
2078 "controller map key does not match its persisted identity",
2079 ));
2080 }
2081 controller.validate()?;
2082 }
2083 for (bucket, receipts) in &self.archive_receipt_buckets {
2084 if bucket.bucket >= DECISION_HISTORY_BUCKET_COUNT || receipts.is_empty() {
2085 return Err(archive_error("decision archive receipt bucket is invalid"));
2086 }
2087 for (key, receipt) in receipts {
2088 if decision_history_page_key(key)? != *bucket
2089 || receipt.encoded_bytes == 0
2090 || matches!(self.decision_locator(key), DecisionHistoryLocation::Hot)
2091 {
2092 return Err(archive_error(
2093 "decision archive receipt is invalid or overlaps hot history",
2094 ));
2095 }
2096 }
2097 let expected_page_id = self
2098 .decision_archive_bucket_page(*bucket)?
2099 .ok_or_else(|| archive_error("nonempty archive bucket has no page"))?
2100 .state_page_id()?;
2101 if self.archive_bucket_page_ids.get(bucket) != Some(&expected_page_id) {
2102 return Err(archive_error(
2103 "decision archive bucket page commitment is inconsistent",
2104 ));
2105 }
2106 }
2107 if self.attempts.iter().any(|attempt| {
2108 attempt.request_commitment.len() != 64
2109 || !attempt
2110 .request_commitment
2111 .bytes()
2112 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
2113 }) {
2114 return Err(DecisionError::new(
2115 DecisionErrorCode::InvalidDecision,
2116 "decision attempts require a canonical request commitment",
2117 ));
2118 }
2119 for (id, ticket) in &self.tickets {
2120 if id != &ticket.id || !self.controllers.contains_key(&ticket.assigned_controller) {
2121 return Err(DecisionError::new(
2122 DecisionErrorCode::InvalidDecision,
2123 "ticket identity or assigned controller is invalid",
2124 ));
2125 }
2126 ticket.validate()?;
2127 }
2128 for (ordinal, trace) in self.traces.ordinals().zip(self.traces.iter()) {
2129 if trace.id.get() != ordinal || ordinal >= self.traces.next_ordinal() {
2130 return Err(DecisionError::new(
2131 DecisionErrorCode::InvalidDecision,
2132 "decision trace identity disagrees with its persistent log ordinal",
2133 ));
2134 }
2135 let ticket_version_valid = self.tickets.get(&trace.ticket_id).is_some_and(|ticket| {
2136 trace.ticket_version != 0 && trace.ticket_version <= ticket.version
2137 }) || self
2138 .contains_archived_key(&DecisionHistoryKey::Ticket(trace.ticket_id));
2139 if !ticket_version_valid || !self.controllers.contains_key(&trace.controller_id) {
2140 return Err(DecisionError::new(
2141 DecisionErrorCode::InvalidDecision,
2142 "decision trace ticket, version, or controller is invalid",
2143 ));
2144 }
2145 }
2146 for ticket in self.tickets.values() {
2147 if let DecisionTicketState::Resolved { trace_id, .. } = ticket.state {
2148 let hot_trace_matches = self
2149 .traces
2150 .get_ordinal(trace_id.get())
2151 .is_some_and(|trace| trace.id == trace_id && trace.ticket_id == ticket.id);
2152 let archived_trace_exists =
2153 self.contains_archived_key(&DecisionHistoryKey::Trace(trace_id));
2154 if !hot_trace_matches && !archived_trace_exists {
2155 return Err(DecisionError::new(
2156 DecisionErrorCode::InvalidDecision,
2157 "resolved ticket does not reference hot or archived trace evidence",
2158 ));
2159 }
2160 }
2161 }
2162 let mut request_ids = std::collections::BTreeSet::new();
2163 for (ordinal, attempt) in self.attempts.ordinals().zip(self.attempts.iter()) {
2164 if attempt.request_id.get() == 0 || !request_ids.insert(attempt.request_id) {
2165 return Err(DecisionError::new(
2166 DecisionErrorCode::InvalidDecision,
2167 "decision attempts must use unique nonzero request IDs",
2168 ));
2169 }
2170 if self.attempts_by_request.get(&attempt.request_id) != Some(&ordinal) {
2171 return Err(DecisionError::new(
2172 DecisionErrorCode::InvalidDecision,
2173 "decision attempt request index is inconsistent",
2174 ));
2175 }
2176 match &attempt.outcome {
2177 DecisionAttemptOutcome::Accepted {
2178 trace_id,
2179 command_request_id,
2180 } => {
2181 if command_request_id.is_some() && trace_id.is_none() {
2182 return Err(DecisionError::new(
2183 DecisionErrorCode::InvalidDecision,
2184 "accepted decision commands require a decision trace",
2185 ));
2186 }
2187 if trace_id.is_some_and(|trace_id| {
2188 self.traces.get_ordinal(trace_id.get()).is_none()
2189 && !self.contains_archived_key(&DecisionHistoryKey::Trace(trace_id))
2190 }) {
2191 return Err(DecisionError::new(
2192 DecisionErrorCode::InvalidDecision,
2193 "accepted decision attempt references unavailable trace evidence",
2194 ));
2195 }
2196 }
2197 DecisionAttemptOutcome::Rejected { message, .. } => {
2198 require_text(message, "decision rejection message")?;
2199 }
2200 }
2201 }
2202 Ok(())
2203 }
2204
2205 fn validate_archive_directory_shape(&self) -> Result<(), DecisionError> {
2206 if self
2207 .archive_bucket_page_ids
2208 .iter()
2209 .any(|(bucket, page_id)| {
2210 bucket.bucket >= DECISION_HISTORY_BUCKET_COUNT || !canonical_archive_hash(page_id)
2211 })
2212 {
2213 return Err(archive_error(
2214 "decision archive bucket directory is malformed",
2215 ));
2216 }
2217 if (self.archive_receipt_count == 0) != self.archive_bucket_page_ids.is_empty() {
2218 return Err(archive_error(
2219 "decision archive receipt count and page directory disagree",
2220 ));
2221 }
2222 if !self.archive_receipt_buckets.is_empty() {
2223 let resident_count =
2224 self.archive_receipt_buckets
2225 .values()
2226 .try_fold(0_u64, |total, receipts| {
2227 total.checked_add(receipts.len() as u64).ok_or_else(|| {
2228 archive_error("decision archive receipt count overflowed")
2229 })
2230 })?;
2231 if resident_count > self.archive_receipt_count {
2232 return Err(archive_error(
2233 "resident decision archive receipt count exceeds the committed total",
2234 ));
2235 }
2236 }
2237 Ok(())
2238 }
2239
2240 pub fn apply(
2241 &mut self,
2242 mutation: DecisionMutation,
2243 at: SimTime,
2244 trace_id: Option<DecisionTraceId>,
2245 ) -> Result<PreparedDecision, DecisionError> {
2246 let prepared = match mutation {
2247 DecisionMutation::RegisterController { controller } => {
2248 controller.validate()?;
2249 if self.controllers.contains_key(&controller.id) {
2250 return Err(DecisionError::new(
2251 DecisionErrorCode::DuplicateController,
2252 format!(
2253 "decision controller {} is already registered",
2254 controller.id
2255 ),
2256 ));
2257 }
2258 self.controllers
2259 .insert(controller.id.clone(), Arc::new(controller));
2260 PreparedDecision::default()
2261 }
2262 DecisionMutation::Open { mut ticket } => {
2263 ticket.validate()?;
2264 if self.tickets.contains_key(&ticket.id) {
2265 return Err(DecisionError::new(
2266 DecisionErrorCode::DuplicateTicket,
2267 format!("decision ticket {} is already present", ticket.id),
2268 ));
2269 }
2270 if !self.controllers.contains_key(&ticket.assigned_controller) {
2271 return Err(DecisionError::new(
2272 DecisionErrorCode::InvalidController,
2273 format!(
2274 "decision ticket {} names unknown controller {}",
2275 ticket.id, ticket.assigned_controller
2276 ),
2277 ));
2278 }
2279 if ticket.deadline.is_some_and(|deadline| deadline < at) {
2280 return Err(DecisionError::new(
2281 DecisionErrorCode::InvalidDecision,
2282 "decision deadline precedes its admission time",
2283 ));
2284 }
2285 let persisted = DecisionTicket {
2286 id: ticket.id,
2287 definition: ticket.definition,
2288 decision_maker: ticket.decision_maker,
2289 assigned_controller: ticket.assigned_controller,
2290 summary: ticket.summary,
2291 context: ticket.context,
2292 options: std::mem::take(&mut ticket.options),
2293 opened_at: at,
2294 updated_at: at,
2295 deadline: ticket.deadline,
2296 version: 1,
2297 state: DecisionTicketState::Open,
2298 };
2299 if let Some(deadline) = persisted.deadline {
2300 self.deadline_index
2301 .entry(deadline)
2302 .or_default()
2303 .insert(persisted.id);
2304 }
2305 self.insert_hot_history_record(&DecisionArchiveRecord::Ticket {
2306 ticket: persisted.clone(),
2307 })?;
2308 self.tickets.insert(persisted.id, Arc::new(persisted));
2309 PreparedDecision::default()
2310 }
2311 DecisionMutation::ReplaceOptions {
2312 ticket_id,
2313 expected_version,
2314 context,
2315 mut options,
2316 } => {
2317 context.validate()?;
2318 canonicalize_options(&mut options)?;
2319 let previous = self
2320 .tickets
2321 .get(&ticket_id)
2322 .map(|ticket| ticket.as_ref().clone())
2323 .ok_or_else(|| {
2324 DecisionError::new(
2325 DecisionErrorCode::TicketNotFound,
2326 format!("decision ticket {ticket_id} was not found"),
2327 )
2328 })?;
2329 let ticket = self.open_ticket_mut(ticket_id, expected_version, at)?;
2330 ticket.context = context;
2331 ticket.options = options;
2332 ticket.updated_at = at;
2333 ticket.version = ticket.version.checked_add(1).ok_or_else(|| {
2334 DecisionError::new(
2335 DecisionErrorCode::InvalidDecision,
2336 "decision ticket version is exhausted",
2337 )
2338 })?;
2339 ticket.validate()?;
2340 let updated = ticket.clone();
2341 let _ = ticket;
2342 self.replace_hot_history_record(
2343 &DecisionArchiveRecord::Ticket { ticket: previous },
2344 &DecisionArchiveRecord::Ticket { ticket: updated },
2345 )?;
2346 PreparedDecision::default()
2347 }
2348 DecisionMutation::Resolve {
2349 ticket_id,
2350 expected_version,
2351 controller_id,
2352 policy,
2353 decision,
2354 command_request_id,
2355 } => self.resolve(
2356 ticket_id,
2357 expected_version,
2358 &controller_id,
2359 policy,
2360 decision,
2361 command_request_id,
2362 at,
2363 trace_id.ok_or_else(|| {
2364 DecisionError::new(
2365 DecisionErrorCode::InvalidDecision,
2366 "decision resolution requires a claimed trace ID",
2367 )
2368 })?,
2369 )?,
2370 DecisionMutation::Cancel {
2371 ticket_id,
2372 expected_version,
2373 reason,
2374 } => {
2375 require_text(&reason, "decision cancellation reason")?;
2376 let previous = self
2377 .tickets
2378 .get(&ticket_id)
2379 .map(|ticket| ticket.as_ref().clone())
2380 .ok_or_else(|| {
2381 DecisionError::new(
2382 DecisionErrorCode::TicketNotFound,
2383 format!("decision ticket {ticket_id} was not found"),
2384 )
2385 })?;
2386 let ticket = self.open_ticket_mut(ticket_id, expected_version, at)?;
2387 ticket.updated_at = at;
2388 ticket.version = ticket.version.checked_add(1).ok_or_else(|| {
2389 DecisionError::new(
2390 DecisionErrorCode::InvalidDecision,
2391 "decision ticket version is exhausted",
2392 )
2393 })?;
2394 ticket.state = DecisionTicketState::Cancelled { reason };
2395 ticket.validate()?;
2396 let deadline = ticket.deadline;
2397 let updated = ticket.clone();
2398 let _ = ticket;
2399 self.replace_hot_history_record(
2400 &DecisionArchiveRecord::Ticket { ticket: previous },
2401 &DecisionArchiveRecord::Ticket { ticket: updated },
2402 )?;
2403 self.remove_deadline(ticket_id, deadline);
2404 PreparedDecision::default()
2405 }
2406 };
2407 self.advance_time(at)?;
2408 Ok(prepared)
2409 }
2410
2411 fn open_ticket_mut(
2412 &mut self,
2413 ticket_id: DecisionTicketId,
2414 expected_version: u64,
2415 at: SimTime,
2416 ) -> Result<&mut DecisionTicket, DecisionError> {
2417 let ticket = self.tickets.get_mut(&ticket_id).ok_or_else(|| {
2418 DecisionError::new(
2419 DecisionErrorCode::TicketNotFound,
2420 format!("decision ticket {ticket_id} was not found"),
2421 )
2422 })?;
2423 let ticket = Arc::make_mut(ticket);
2424 if !ticket.is_open() || ticket.deadline.is_some_and(|deadline| deadline < at) {
2425 return Err(DecisionError::new(
2426 DecisionErrorCode::ClosedTicket,
2427 format!("decision ticket {ticket_id} is not open"),
2428 ));
2429 }
2430 if ticket.version != expected_version {
2431 return Err(DecisionError::new(
2432 DecisionErrorCode::VersionConflict,
2433 format!(
2434 "decision ticket {ticket_id} is at version {}, expected {expected_version}",
2435 ticket.version
2436 ),
2437 ));
2438 }
2439 Ok(ticket)
2440 }
2441
2442 #[allow(clippy::too_many_arguments)]
2443 fn resolve(
2444 &mut self,
2445 ticket_id: DecisionTicketId,
2446 expected_version: u64,
2447 controller_id: &str,
2448 policy: DecisionPolicyIdentity,
2449 decision: PolicyDecision,
2450 command_request_id: Option<CommandRequestId>,
2451 at: SimTime,
2452 trace_id: DecisionTraceId,
2453 ) -> Result<PreparedDecision, DecisionError> {
2454 let controller = self.controllers.get(controller_id).ok_or_else(|| {
2455 DecisionError::new(
2456 DecisionErrorCode::InvalidController,
2457 format!("decision controller {controller_id} was not found"),
2458 )
2459 })?;
2460 if controller.policy != policy {
2461 return Err(DecisionError::new(
2462 DecisionErrorCode::PolicyMismatch,
2463 "decision resolution policy does not match the persisted controller binding",
2464 ));
2465 }
2466 if (controller.policy.kind == crate::DecisionPolicyKind::Random)
2467 != decision.random.is_some()
2468 {
2469 return Err(DecisionError::new(
2470 DecisionErrorCode::PolicyMismatch,
2471 "random decision controllers require random draw evidence, and other controllers reject it",
2472 ));
2473 }
2474 let previous_ticket = self
2475 .tickets
2476 .get(&ticket_id)
2477 .map(|ticket| ticket.as_ref().clone())
2478 .ok_or_else(|| {
2479 DecisionError::new(
2480 DecisionErrorCode::TicketNotFound,
2481 format!("decision ticket {ticket_id} was not found"),
2482 )
2483 })?;
2484 let ticket = self.open_ticket_mut(ticket_id, expected_version, at)?;
2485 if ticket.assigned_controller != controller_id {
2486 return Err(DecisionError::new(
2487 DecisionErrorCode::InvalidController,
2488 "decision resolution came from a controller not assigned to the ticket",
2489 ));
2490 }
2491 decision.validate(ticket)?;
2492 let action = match &decision.outcome {
2493 DecisionOutcome::Selected { option_id } => {
2494 ticket.option(option_id).map(|option| option.action.clone())
2495 }
2496 DecisionOutcome::Deferred { .. } => None,
2497 DecisionOutcome::Pending { .. } => {
2498 return Err(DecisionError::new(
2499 DecisionErrorCode::InvalidDecision,
2500 "pending policy outcomes are not authoritative decision mutations",
2501 ));
2502 }
2503 };
2504 if matches!(action, Some(DecisionAction::Command { .. })) != command_request_id.is_some() {
2505 return Err(DecisionError::new(
2506 DecisionErrorCode::InvalidDecision,
2507 "command actions require exactly one command request ID",
2508 ));
2509 }
2510 let trace = DecisionTrace {
2511 id: trace_id,
2512 ticket_id,
2513 ticket_version: ticket.version,
2514 controller_id: controller_id.to_owned(),
2515 policy,
2516 decided_at: at,
2517 outcome: decision.outcome.clone(),
2518 summary: decision.summary,
2519 evaluations: decision.evaluations,
2520 external: decision.external,
2521 random: decision.random,
2522 command_request_id,
2523 };
2524 ticket.updated_at = at;
2525 ticket.version = ticket.version.checked_add(1).ok_or_else(|| {
2526 DecisionError::new(
2527 DecisionErrorCode::InvalidDecision,
2528 "decision ticket version is exhausted",
2529 )
2530 })?;
2531 if let DecisionOutcome::Selected { option_id } = &trace.outcome {
2532 ticket.state = DecisionTicketState::Resolved {
2533 option_id: option_id.clone(),
2534 trace_id,
2535 };
2536 }
2537 ticket.validate()?;
2538 let deadline = ticket.deadline;
2539 let updated_ticket = ticket.clone();
2540 let _ = ticket;
2541 self.replace_hot_history_record(
2542 &DecisionArchiveRecord::Ticket {
2543 ticket: previous_ticket,
2544 },
2545 &DecisionArchiveRecord::Ticket {
2546 ticket: updated_ticket,
2547 },
2548 )?;
2549 self.remove_deadline(ticket_id, deadline);
2550 self.insert_hot_history_record(&DecisionArchiveRecord::Trace {
2551 trace: trace.clone(),
2552 })?;
2553 let ordinal = self.traces.push(trace.clone());
2554 debug_assert_eq!(ordinal, trace.id.get());
2555 Ok(PreparedDecision {
2556 trace: Some(trace),
2557 action,
2558 })
2559 }
2560
2561 pub fn advance_time(&mut self, at: SimTime) -> Result<(), DecisionError> {
2562 let due = self
2563 .deadline_index
2564 .range(..at)
2565 .map(|(deadline, tickets)| (*deadline, tickets.iter().copied().collect::<Vec<_>>()))
2566 .collect::<Vec<_>>();
2567 for (deadline, ticket_ids) in due {
2568 self.deadline_index.remove(&deadline);
2569 for ticket_id in ticket_ids {
2570 let previous = self
2571 .tickets
2572 .get(&ticket_id)
2573 .map(|ticket| ticket.as_ref().clone())
2574 .ok_or_else(|| {
2575 DecisionError::new(
2576 DecisionErrorCode::InvalidDecision,
2577 "decision ticket index changed during time advancement",
2578 )
2579 })?;
2580 let Some(ticket) = self.tickets.get_mut(&ticket_id) else {
2581 return Err(DecisionError::new(
2582 DecisionErrorCode::InvalidDecision,
2583 "decision ticket index changed during time advancement",
2584 ));
2585 };
2586 let ticket = Arc::make_mut(ticket);
2587 let updated =
2588 if ticket.is_open() && ticket.deadline.is_some_and(|deadline| deadline < at) {
2589 ticket.updated_at = at;
2590 ticket.version = ticket.version.checked_add(1).ok_or_else(|| {
2591 DecisionError::new(
2592 DecisionErrorCode::InvalidDecision,
2593 "decision ticket version is exhausted",
2594 )
2595 })?;
2596 ticket.state = DecisionTicketState::Expired;
2597 ticket.validate()?;
2598 Some(ticket.clone())
2599 } else {
2600 None
2601 };
2602 let _ = ticket;
2603 if let Some(updated) = updated {
2604 self.replace_hot_history_record(
2605 &DecisionArchiveRecord::Ticket { ticket: previous },
2606 &DecisionArchiveRecord::Ticket { ticket: updated },
2607 )?;
2608 }
2609 }
2610 }
2611 Ok(())
2612 }
2613
2614 fn remove_deadline(&mut self, ticket_id: DecisionTicketId, deadline: Option<SimTime>) {
2615 let Some(deadline) = deadline else {
2616 return;
2617 };
2618 if let Some(mut tickets) = self.deadline_index.get(&deadline).cloned() {
2619 tickets.remove(&ticket_id);
2620 if tickets.is_empty() {
2621 self.deadline_index.remove(&deadline);
2622 } else {
2623 self.deadline_index.insert(deadline, tickets);
2624 }
2625 }
2626 }
2627}
2628
2629fn validate_attempt_shape(
2630 state: &DecisionState,
2631 attempt: &DecisionAttemptRecord,
2632) -> Result<(), DecisionError> {
2633 if attempt.request_id.get() == 0 || !canonical_archive_hash(&attempt.request_commitment) {
2634 return Err(DecisionError::new(
2635 DecisionErrorCode::InvalidDecision,
2636 "decision attempt identity or request commitment is invalid",
2637 ));
2638 }
2639 match &attempt.outcome {
2640 DecisionAttemptOutcome::Accepted {
2641 trace_id,
2642 command_request_id,
2643 } => {
2644 if command_request_id.is_some() && trace_id.is_none() {
2645 return Err(DecisionError::new(
2646 DecisionErrorCode::InvalidDecision,
2647 "accepted decision commands require a decision trace",
2648 ));
2649 }
2650 if trace_id.is_some_and(|trace_id| {
2651 state.traces.get_ordinal(trace_id.get()).is_none()
2652 && !state.contains_archived_key(&DecisionHistoryKey::Trace(trace_id))
2653 }) {
2654 return Err(DecisionError::new(
2655 DecisionErrorCode::InvalidDecision,
2656 "accepted decision attempt references unavailable trace evidence",
2657 ));
2658 }
2659 }
2660 DecisionAttemptOutcome::Rejected { message, .. } => {
2661 require_text(message, "decision rejection message")?;
2662 }
2663 }
2664 Ok(())
2665}
2666
2667fn decision_hot_leaf_hash(
2668 key: &DecisionHistoryKey,
2669 record: &DecisionArchiveRecord,
2670) -> Result<[u8; 32], DecisionError> {
2671 if record.key() != *key {
2672 return Err(archive_error(
2673 "decision hot-history leaf identity is inconsistent",
2674 ));
2675 }
2676 let bytes = serde_json::to_vec(&(key, record)).map_err(|error| {
2677 archive_error(format!("cannot encode decision hot-history leaf: {error}"))
2678 })?;
2679 let mut hasher = blake3::Hasher::new();
2680 hasher.update(b"canwu.decision.hot-history-leaf.v2");
2681 hasher.update(&[0]);
2682 hasher.update(&bytes);
2683 Ok(*hasher.finalize().as_bytes())
2684}
2685
2686fn add_digest_mod_256(target: &mut [u8; 32], digest: [u8; 32]) {
2687 let mut carry = 0_u16;
2688 for index in (0..target.len()).rev() {
2689 let value = u16::from(target[index]) + u16::from(digest[index]) + carry;
2690 target[index] = u8::try_from(value % 256).expect("modulo 256 always fits in u8");
2691 carry = value >> 8;
2692 }
2693}
2694
2695fn subtract_digest_mod_256(target: &mut [u8; 32], digest: [u8; 32]) {
2696 let mut borrow = 0_i16;
2697 for index in (0..target.len()).rev() {
2698 let value = i16::from(target[index]) - i16::from(digest[index]) - borrow;
2699 if value < 0 {
2700 target[index] =
2701 u8::try_from(value + 256).expect("normalized byte subtraction fits in u8");
2702 borrow = 1;
2703 } else {
2704 target[index] = u8::try_from(value).expect("non-negative byte subtraction fits in u8");
2705 borrow = 0;
2706 }
2707 }
2708}
2709
2710pub fn decision_history_bucket(key: &DecisionHistoryKey) -> Result<u16, DecisionError> {
2711 Ok(decision_history_page_key(key)?.bucket)
2712}
2713
2714pub fn decision_history_page_key(
2715 key: &DecisionHistoryKey,
2716) -> Result<DecisionArchivePageKey, DecisionError> {
2717 let encoded = serde_json::to_vec(key)
2718 .map_err(|error| archive_error(format!("cannot encode decision history key: {error}")))?;
2719 let digest = blake3::hash(&encoded);
2720 let bytes = digest.as_bytes();
2721 Ok(DecisionArchivePageKey {
2722 bucket: (u16::from(bytes[0]) << 4) | u16::from(bytes[1] >> 4),
2723 segment: bytes[1] & 0x0f,
2724 })
2725}
2726
2727#[derive(Default)]
2728struct DecisionScaleArchive {
2729 blobs: RefCell<StdBTreeMap<String, DecisionArchiveBlob>>,
2730 pages: RefCell<StdBTreeMap<String, DecisionArchiveBucketPage>>,
2731}
2732
2733impl DecisionArchiveProvider for DecisionScaleArchive {
2734 fn load_decision_archive(
2735 &self,
2736 locator: &str,
2737 ) -> Result<Option<DecisionArchiveBlob>, DecisionError> {
2738 Ok(self.blobs.borrow().get(locator).cloned())
2739 }
2740
2741 fn load_decision_archive_bucket_page(
2742 &self,
2743 page_id: &str,
2744 ) -> Result<Option<DecisionArchiveBucketPage>, DecisionError> {
2745 Ok(self.pages.borrow().get(page_id).cloned())
2746 }
2747}
2748
2749impl DecisionArchiveStore for DecisionScaleArchive {
2750 fn store_decision_archive(
2751 &self,
2752 blob: &DecisionArchiveBlob,
2753 ) -> Result<DecisionArchiveStoreOutcome, DecisionError> {
2754 let locator = blob.content_id()?;
2755 let mut blobs = self.blobs.borrow_mut();
2756 if let Some(existing) = blobs.get(&locator) {
2757 if existing != blob {
2758 return Err(archive_error(
2759 "decision scale archive locator contains different content",
2760 ));
2761 }
2762 return Ok(DecisionArchiveStoreOutcome::AlreadyStored);
2763 }
2764 blobs.insert(locator, blob.clone());
2765 Ok(DecisionArchiveStoreOutcome::Stored)
2766 }
2767}
2768
2769#[doc(hidden)]
2773pub fn format8_decision_locator_scale_fixture(
2774 key_count: usize,
2775) -> Result<DecisionLocatorScaleFixture, DecisionError> {
2776 let mut state = DecisionState::default();
2777 let archive = DecisionScaleArchive::default();
2778 let mut archive_batches = 0_u64;
2779 let mut next_ordinal = 1_usize;
2780 while next_ordinal <= key_count {
2781 let batch_end = next_ordinal
2782 .saturating_add(MAX_DECISION_ARCHIVE_BATCH_ENTRIES - 1)
2783 .min(key_count);
2784 let mut keys = Vec::with_capacity(batch_end - next_ordinal + 1);
2785 for ordinal in next_ordinal..=batch_end {
2786 let ordinal = u64::try_from(ordinal)
2787 .map_err(|_| archive_error("decision locator scale key exceeds u64"))?;
2788 let request_id = canwu_core::DecisionRequestId::new(ordinal);
2789 let request_commitment = blake3::hash(&ordinal.to_be_bytes()).to_hex().to_string();
2790 state.append_attempt(DecisionAttemptRecord {
2791 request_id,
2792 request_commitment,
2793 at: SimTime::from_minutes(
2794 i64::try_from(ordinal)
2795 .map_err(|_| archive_error("decision scale time exceeds i64"))?,
2796 ),
2797 revision_before: ordinal - 1,
2798 expected_revision: ordinal - 1,
2799 outcome: DecisionAttemptOutcome::Rejected {
2800 code: crate::DecisionAttemptErrorCode::InvalidDecision,
2801 message: "format8 production-path scale attempt".to_owned(),
2802 },
2803 })?;
2804 keys.push(DecisionHistoryKey::Attempt(request_id));
2805 }
2806 let prepared = state.prepare_decision_archive(&keys)?;
2807 for blob in &prepared.blobs {
2808 let _ = archive.store_decision_archive(blob)?;
2809 }
2810 let verified = state.verify_decision_archive(&prepared, &archive)?;
2811 state = state.commit_verified_decision_archive(&verified)?;
2812 let touched_pages = keys
2813 .iter()
2814 .map(decision_history_page_key)
2815 .collect::<Result<OrdSet<_>, _>>()?;
2816 for page_key in touched_pages {
2817 let page = state
2818 .decision_archive_bucket_page(page_key)?
2819 .ok_or_else(|| archive_error("committed decision locator page is missing"))?;
2820 let page_id = page.state_page_id()?;
2821 archive.pages.borrow_mut().insert(page_id, page);
2822 }
2823 archive_batches = archive_batches
2824 .checked_add(1)
2825 .ok_or_else(|| archive_error("decision archive batch count overflowed"))?;
2826 next_ordinal = batch_end.saturating_add(1);
2827 }
2828
2829 let archive_root = state.archive_receipt_root()?;
2830 let hot_bytes = serde_json::to_vec(&state.paged_checkpoint_hot_state()).map_err(|error| {
2831 archive_error(format!(
2832 "cannot persist root-only decision hot state: {error}"
2833 ))
2834 })?;
2835 let hot = serde_json::from_slice::<DecisionState>(&hot_bytes).map_err(|error| {
2836 archive_error(format!(
2837 "cannot restart root-only decision hot state: {error}"
2838 ))
2839 })?;
2840 let restarted = DecisionState::from_paged_checkpoint_root(
2841 hot,
2842 state.archive_bucket_page_ids.clone(),
2843 state.archive_receipt_count,
2844 &archive_root,
2845 )?;
2846 let samples: OrdSet<usize> = [1_usize, key_count.saturating_div(2).max(1), key_count]
2847 .into_iter()
2848 .filter(|ordinal| *ordinal <= key_count)
2849 .collect::<OrdSet<_>>();
2850 for ordinal in &samples {
2851 let key = DecisionHistoryKey::Attempt(canwu_core::DecisionRequestId::new(
2852 u64::try_from(*ordinal)
2853 .map_err(|_| archive_error("decision scale sample exceeds u64"))?,
2854 ));
2855 if restarted.load_decision_history(&key, &archive)?.is_none() {
2856 return Err(archive_error(
2857 "root-only decision restart lost exact archived history",
2858 ));
2859 }
2860 }
2861 let reachable = restarted.archive_reachability(&archive)?;
2862 let mut max_page_entries = 0_u64;
2863 let mut max_page_encoded_bytes = 0_u64;
2864 for page_id in restarted.archive_bucket_page_ids.values() {
2865 let page = archive
2866 .load_decision_archive_bucket_page(page_id)?
2867 .ok_or_else(|| archive_error("decision scale locator page is unavailable"))?;
2868 max_page_entries = max_page_entries.max(page.receipts.len() as u64);
2869 max_page_encoded_bytes = max_page_encoded_bytes.max(
2870 serde_json::to_vec(&page)
2871 .map_err(|error| archive_error(format!("cannot size locator page: {error}")))?
2872 .len() as u64,
2873 );
2874 }
2875 let entry_bytes = size_of::<DecisionHistoryKey>()
2876 .saturating_add(size_of::<CompactDecisionArchiveReceipt>())
2877 .saturating_add(48);
2878 let metrics = DecisionLocatorScaleMetrics {
2879 entries: restarted.archived_history_count() as u64,
2880 locator_pages: restarted.archive_bucket_page_ids.len() as u64,
2881 max_page_entries,
2882 max_page_encoded_bytes,
2883 archive_batches,
2884 exact_restart_queries: samples.len() as u64,
2885 reachable_blob_locators: reachable.blob_locators.len() as u64,
2886 estimated_resident_structural_bytes: (restarted.archived_history_count() as u64)
2887 .saturating_mul(entry_bytes as u64),
2888 root_hash: archive_root,
2889 };
2890 Ok(DecisionLocatorScaleFixture {
2891 state,
2892 archive_blobs: archive.blobs.into_inner().into_values().collect(),
2893 metrics,
2894 })
2895}
2896
2897pub fn format8_decision_locator_scale_probe(
2900 key_count: usize,
2901) -> Result<DecisionLocatorScaleMetrics, DecisionError> {
2902 Ok(format8_decision_locator_scale_fixture(key_count)?.metrics)
2903}
2904
2905pub fn format8_trace_locator_scale_probe(
2909 trace_count: usize,
2910) -> Result<TraceLocatorScaleMetrics, DecisionError> {
2911 if trace_count == 0 {
2912 return Err(archive_error(
2913 "trace locator scale probe requires at least one trace",
2914 ));
2915 }
2916 let mut state = DecisionState::default();
2917 let controller_id = "format8-trace-controller".to_owned();
2918 let policy =
2919 DecisionPolicyIdentity::new(crate::DecisionPolicyKind::Rule, "format8-trace-policy", "1");
2920 state.controllers.insert(
2921 controller_id.clone(),
2922 Arc::new(DecisionControllerBinding::new(
2923 controller_id.clone(),
2924 policy.clone(),
2925 crate::DecisionAuthority::NoResponsibleActor {
2926 reason: "Format-8 trace scale fixture".to_owned(),
2927 },
2928 )),
2929 );
2930 let ticket_id = DecisionTicketId::new(1);
2931 let ticket = DecisionTicket {
2932 id: ticket_id,
2933 definition: "format8-trace-scale".to_owned(),
2934 decision_maker: canwu_core::EntityRef::Person(canwu_core::PersonId::new(1)),
2935 assigned_controller: controller_id.clone(),
2936 summary: "Format-8 trace scale ticket".to_owned(),
2937 context: crate::DecisionContext::new("format8-trace-scale", serde_json::json!({})),
2938 options: vec![crate::DecisionOption::new("defer", "Defer")],
2939 opened_at: SimTime::EPOCH,
2940 updated_at: SimTime::EPOCH,
2941 deadline: None,
2942 version: 1,
2943 state: DecisionTicketState::Cancelled {
2944 reason: "Scale fixture is terminal".to_owned(),
2945 },
2946 };
2947 ticket.validate()?;
2948 state.insert_hot_history_record(&DecisionArchiveRecord::Ticket {
2949 ticket: ticket.clone(),
2950 })?;
2951 state.tickets.insert(ticket_id, Arc::new(ticket));
2952 for ordinal in 1..=trace_count {
2953 let ordinal = u64::try_from(ordinal)
2954 .map_err(|_| archive_error("trace locator scale ordinal exceeds u64"))?;
2955 let trace = DecisionTrace {
2956 id: DecisionTraceId::new(ordinal),
2957 ticket_id,
2958 ticket_version: 1,
2959 controller_id: controller_id.clone(),
2960 policy: policy.clone(),
2961 decided_at: SimTime::from_minutes(
2962 i64::try_from(ordinal)
2963 .map_err(|_| archive_error("trace locator scale time exceeds i64"))?,
2964 ),
2965 outcome: DecisionOutcome::Deferred {
2966 reason: "trace-scale".to_owned(),
2967 },
2968 summary: "trace-scale".to_owned(),
2969 evaluations: Vec::new(),
2970 external: None,
2971 random: None,
2972 command_request_id: None,
2973 };
2974 state.insert_hot_history_record(&DecisionArchiveRecord::Trace {
2975 trace: trace.clone(),
2976 })?;
2977 let inserted = state.traces.push(trace);
2978 if inserted != ordinal {
2979 return Err(archive_error(
2980 "trace locator scale ordinal insertion is inconsistent",
2981 ));
2982 }
2983 }
2984 state.validate()?;
2985 let samples: OrdSet<usize> = [1_usize, trace_count.saturating_div(2).max(1), trace_count]
2986 .into_iter()
2987 .filter(|ordinal| *ordinal <= trace_count)
2988 .collect::<OrdSet<_>>();
2989 for ordinal in &samples {
2990 let key = DecisionHistoryKey::Trace(DecisionTraceId::new(
2991 u64::try_from(*ordinal)
2992 .map_err(|_| archive_error("trace locator scale sample exceeds u64"))?,
2993 ));
2994 if state.decision_locator(&key) != DecisionHistoryLocation::Hot {
2995 return Err(archive_error(
2996 "trace locator ordinal index lost a retained hot trace",
2997 ));
2998 }
2999 }
3000 let target = DecisionHistoryKey::Trace(DecisionTraceId::new(
3001 u64::try_from(trace_count)
3002 .map_err(|_| archive_error("trace locator scale target exceeds u64"))?,
3003 ));
3004 let archive = DecisionScaleArchive::default();
3005 let prepared = state.prepare_decision_archive(std::slice::from_ref(&target))?;
3006 for blob in &prepared.blobs {
3007 let _ = archive.store_decision_archive(blob)?;
3008 }
3009 let verified = state.verify_decision_archive(&prepared, &archive)?;
3010 let committed = state.commit_verified_decision_archive(&verified)?;
3011 let target_archived = matches!(
3012 committed.decision_locator(&target),
3013 DecisionHistoryLocation::Archived { .. }
3014 );
3015 if !target_archived {
3016 return Err(archive_error(
3017 "trace locator scale archive commit did not release its target",
3018 ));
3019 }
3020 Ok(TraceLocatorScaleMetrics {
3021 hot_trace_entries: trace_count as u64,
3022 indexed_lookup_samples: samples.len() as u64,
3023 archive_commit_entries: verified.receipts.len() as u64,
3024 target_archived,
3025 })
3026}
3027
3028#[cfg(test)]
3029mod archive_restart_tests {
3030 use super::*;
3031
3032 fn linked_decision_state(resolved_ticket: bool, accepted_attempt: bool) -> DecisionState {
3033 let mut state = DecisionState::default();
3034 let controller_id = "restart-controller".to_owned();
3035 let policy =
3036 DecisionPolicyIdentity::new(crate::DecisionPolicyKind::Rule, "restart-policy", "1");
3037 state.controllers.insert(
3038 controller_id.clone(),
3039 Arc::new(DecisionControllerBinding::new(
3040 controller_id.clone(),
3041 policy.clone(),
3042 crate::DecisionAuthority::NoResponsibleActor {
3043 reason: "restart dependency fixture".to_owned(),
3044 },
3045 )),
3046 );
3047 let ticket_id = DecisionTicketId::new(91);
3048 let trace_id = DecisionTraceId::new(1);
3049 let ticket = DecisionTicket {
3050 id: ticket_id,
3051 definition: "restart-dependency".to_owned(),
3052 decision_maker: canwu_core::EntityRef::Person(canwu_core::PersonId::new(9)),
3053 assigned_controller: controller_id.clone(),
3054 summary: "Restart dependency fixture".to_owned(),
3055 context: crate::DecisionContext::new("restart-dependency", serde_json::json!({})),
3056 options: vec![crate::DecisionOption::new("accept", "Accept")],
3057 opened_at: SimTime::EPOCH,
3058 updated_at: SimTime::from_minutes(1),
3059 deadline: None,
3060 version: 2,
3061 state: if resolved_ticket {
3062 DecisionTicketState::Resolved {
3063 option_id: "accept".to_owned(),
3064 trace_id,
3065 }
3066 } else {
3067 DecisionTicketState::Cancelled {
3068 reason: "terminal fixture".to_owned(),
3069 }
3070 },
3071 };
3072 state
3073 .insert_hot_history_record(&DecisionArchiveRecord::Ticket {
3074 ticket: ticket.clone(),
3075 })
3076 .expect("insert fixture ticket history");
3077 state.tickets.insert(ticket_id, Arc::new(ticket));
3078 let trace = DecisionTrace {
3079 id: trace_id,
3080 ticket_id,
3081 ticket_version: 1,
3082 controller_id,
3083 policy,
3084 decided_at: SimTime::from_minutes(1),
3085 outcome: if resolved_ticket {
3086 DecisionOutcome::Selected {
3087 option_id: "accept".to_owned(),
3088 }
3089 } else {
3090 DecisionOutcome::Deferred {
3091 reason: "terminal fixture".to_owned(),
3092 }
3093 },
3094 summary: "Restart dependency trace".to_owned(),
3095 evaluations: Vec::new(),
3096 external: None,
3097 random: None,
3098 command_request_id: None,
3099 };
3100 state
3101 .insert_hot_history_record(&DecisionArchiveRecord::Trace {
3102 trace: trace.clone(),
3103 })
3104 .expect("insert fixture trace history");
3105 assert_eq!(state.traces.push(trace), trace_id.get());
3106 if accepted_attempt {
3107 state
3108 .append_attempt(DecisionAttemptRecord {
3109 request_id: canwu_core::DecisionRequestId::new(92),
3110 request_commitment: "a".repeat(64),
3111 at: SimTime::from_minutes(1),
3112 revision_before: 1,
3113 expected_revision: 1,
3114 outcome: DecisionAttemptOutcome::Accepted {
3115 trace_id: Some(trace_id),
3116 command_request_id: None,
3117 },
3118 })
3119 .expect("append fixture attempt");
3120 }
3121 state.validate().expect("fixture state validates");
3122 state
3123 }
3124
3125 fn restart_with_exact_dependency_pages(state: &DecisionState) -> DecisionState {
3126 let hot_bytes = serde_json::to_vec(&state.paged_checkpoint_hot_state())
3127 .expect("encode paged hot state");
3128 let hot: DecisionState =
3129 serde_json::from_slice(&hot_bytes).expect("decode paged hot state");
3130 let required_pages = hot
3131 .required_archived_dependency_page_keys()
3132 .expect("derive exact dependency pages");
3133 assert!(!required_pages.is_empty());
3134 let resident_pages = required_pages
3135 .iter()
3136 .map(|page_key| {
3137 state
3138 .decision_archive_bucket_page(*page_key)
3139 .expect("encode dependency page")
3140 .expect("dependency page is resident")
3141 })
3142 .collect::<Vec<_>>();
3143 let restarted = DecisionState::from_paged_checkpoint_root_with_resident_pages(
3144 hot,
3145 state.archive_bucket_page_ids.clone(),
3146 state.archive_receipt_count,
3147 &state.archive_receipt_root().expect("archive root"),
3148 resident_pages,
3149 )
3150 .expect("restart with exact dependency pages");
3151 let encoded = serde_json::to_vec(&restarted).expect("encode sparse restarted state");
3152 serde_json::from_slice(&encoded).expect("sparse restarted state remains restartable")
3153 }
3154
3155 #[test]
3156 fn paged_checkpoint_root_rejects_hot_archive_metadata() {
3157 let hot = DecisionState {
3158 archive_receipt_count: 1,
3159 ..DecisionState::default()
3160 };
3161 let error = DecisionState::from_paged_checkpoint_root(
3162 hot,
3163 OrdMap::new(),
3164 0,
3165 &decision_archive_hash(
3166 "canwu.decision.archive-receipts.v3",
3167 &(0_u64, OrdMap::<DecisionArchivePageKey, String>::new()),
3168 )
3169 .expect("empty archive root"),
3170 )
3171 .expect_err("hot archive metadata must be rejected");
3172 assert_eq!(error.code, DecisionErrorCode::InvalidDecision);
3173 }
3174
3175 fn append_terminal_attempt(state: &mut DecisionState, ordinal: u64) {
3176 state
3177 .append_attempt(DecisionAttemptRecord {
3178 request_id: canwu_core::DecisionRequestId::new(ordinal),
3179 request_commitment: blake3::hash(&ordinal.to_be_bytes()).to_hex().to_string(),
3180 at: SimTime::from_minutes(i64::try_from(ordinal).expect("test time fits i64")),
3181 revision_before: ordinal.saturating_sub(1),
3182 expected_revision: ordinal.saturating_sub(1),
3183 outcome: DecisionAttemptOutcome::Rejected {
3184 code: crate::DecisionAttemptErrorCode::InvalidDecision,
3185 message: "terminal test attempt".to_owned(),
3186 },
3187 })
3188 .expect("append terminal attempt");
3189 }
3190
3191 fn archive_keys(
3192 state: &DecisionState,
3193 archive: &DecisionScaleArchive,
3194 keys: &[DecisionHistoryKey],
3195 ) -> (DecisionState, VerifiedDecisionArchiveCommit) {
3196 let prepared = state
3197 .prepare_decision_archive(keys)
3198 .expect("prepare archive");
3199 for blob in &prepared.blobs {
3200 archive.store_decision_archive(blob).expect("store blob");
3201 }
3202 let verified = state
3203 .verify_decision_archive(&prepared, archive)
3204 .expect("verify archive");
3205 let next = state
3206 .commit_verified_decision_archive(&verified)
3207 .expect("commit archive");
3208 for page_key in keys
3209 .iter()
3210 .map(decision_history_page_key)
3211 .collect::<Result<OrdSet<_>, _>>()
3212 .expect("derive touched pages")
3213 {
3214 let page = next
3215 .decision_archive_bucket_page(page_key)
3216 .expect("build page")
3217 .expect("page is resident");
3218 archive
3219 .pages
3220 .borrow_mut()
3221 .insert(page.state_page_id().expect("page id"), page);
3222 }
3223 (next, verified)
3224 }
3225
3226 #[test]
3227 fn root_only_restart_can_archive_again_into_the_same_segment() {
3228 let first = 1_u64;
3229 let first_key = DecisionHistoryKey::Attempt(canwu_core::DecisionRequestId::new(first));
3230 let page_key = decision_history_page_key(&first_key).expect("first page key");
3231 let second = (2_u64..1_000_000)
3232 .find(|ordinal| {
3233 decision_history_page_key(&DecisionHistoryKey::Attempt(
3234 canwu_core::DecisionRequestId::new(*ordinal),
3235 ))
3236 .ok()
3237 == Some(page_key)
3238 })
3239 .expect("find same-segment request id");
3240 let second_key = DecisionHistoryKey::Attempt(canwu_core::DecisionRequestId::new(second));
3241 let archive = DecisionScaleArchive::default();
3242 let mut state = DecisionState::default();
3243 append_terminal_attempt(&mut state, first);
3244 let (state, _) = archive_keys(&state, &archive, std::slice::from_ref(&first_key));
3245 let archive_root = state.archive_receipt_root().expect("archive root");
3246 let hot = state.paged_checkpoint_hot_state();
3247 assert!(hot.archive_receipt_buckets.is_empty());
3248 assert!(hot.archive_bucket_page_ids.is_empty());
3249 assert_eq!(hot.archive_receipt_count, 0);
3250 let mut restarted = DecisionState::from_paged_checkpoint_root(
3251 hot,
3252 state.archive_bucket_page_ids.clone(),
3253 state.archive_receipt_count,
3254 &archive_root,
3255 )
3256 .expect("root-only restart");
3257 append_terminal_attempt(&mut restarted, second);
3258 let (restarted, verified) =
3259 archive_keys(&restarted, &archive, std::slice::from_ref(&second_key));
3260
3261 assert_eq!(restarted.archived_history_count(), 2);
3262 assert!(
3263 restarted
3264 .load_decision_history(&first_key, &archive)
3265 .expect("load first")
3266 .is_some()
3267 );
3268 assert!(
3269 restarted
3270 .load_decision_history(&second_key, &archive)
3271 .expect("load second")
3272 .is_some()
3273 );
3274 assert_eq!(
3275 restarted
3276 .archive_reachability(&archive)
3277 .expect("enumerate reachability")
3278 .blob_locators
3279 .len(),
3280 2
3281 );
3282 assert_eq!(
3283 restarted
3284 .commit_verified_decision_archive(&verified)
3285 .expect("replay commit"),
3286 restarted
3287 );
3288 }
3289
3290 #[test]
3291 fn paged_restart_loads_archived_ticket_referenced_by_hot_trace() {
3292 let archive = DecisionScaleArchive::default();
3293 let state = linked_decision_state(true, false);
3294 let ticket_id = DecisionTicketId::new(91);
3295 let (archived, _) =
3296 archive_keys(&state, &archive, &[DecisionHistoryKey::Ticket(ticket_id)]);
3297 let restarted = restart_with_exact_dependency_pages(&archived);
3298
3299 assert!(restarted.trace(DecisionTraceId::new(1)).is_some());
3300 assert!(matches!(
3301 restarted.decision_locator(&DecisionHistoryKey::Ticket(ticket_id)),
3302 DecisionHistoryLocation::Archived { .. }
3303 ));
3304 restarted
3305 .validate()
3306 .expect("hot trace dependency validates");
3307 }
3308
3309 #[test]
3310 fn paged_restart_loads_archived_trace_referenced_by_resolved_hot_ticket() {
3311 let archive = DecisionScaleArchive::default();
3312 let state = linked_decision_state(true, false);
3313 let trace_id = DecisionTraceId::new(1);
3314 let (archived, _) = archive_keys(&state, &archive, &[DecisionHistoryKey::Trace(trace_id)]);
3315 let restarted = restart_with_exact_dependency_pages(&archived);
3316
3317 assert!(restarted.ticket(DecisionTicketId::new(91)).is_some());
3318 assert!(matches!(
3319 restarted.decision_locator(&DecisionHistoryKey::Trace(trace_id)),
3320 DecisionHistoryLocation::Archived { .. }
3321 ));
3322 restarted
3323 .validate()
3324 .expect("resolved hot ticket dependency validates");
3325 }
3326
3327 #[test]
3328 fn paged_restart_loads_archived_trace_referenced_by_accepted_hot_attempt() {
3329 let archive = DecisionScaleArchive::default();
3330 let state = linked_decision_state(false, true);
3331 let trace_id = DecisionTraceId::new(1);
3332 let (archived, _) = archive_keys(&state, &archive, &[DecisionHistoryKey::Trace(trace_id)]);
3333 let restarted = restart_with_exact_dependency_pages(&archived);
3334
3335 assert!(
3336 restarted
3337 .attempt(canwu_core::DecisionRequestId::new(92))
3338 .is_some()
3339 );
3340 assert!(matches!(
3341 restarted.decision_locator(&DecisionHistoryKey::Trace(trace_id)),
3342 DecisionHistoryLocation::Archived { .. }
3343 ));
3344 restarted
3345 .validate()
3346 .expect("accepted hot attempt dependency validates");
3347 }
3348}
3349
3350#[derive(Clone, Debug, Default, Eq, PartialEq)]
3351pub struct PreparedDecision {
3352 pub trace: Option<DecisionTrace>,
3353 pub action: Option<DecisionAction>,
3354}
3355
3356#[derive(Clone, Debug, Eq, PartialEq)]
3357pub enum ControllerDecision {
3358 Authoritative {
3359 decision: PolicyDecision,
3360 action: Option<DecisionAction>,
3361 },
3362 Pending(PolicyDecision),
3363}
3364
3365pub struct DecisionController;
3366
3367impl DecisionController {
3368 pub fn evaluate(
3369 ticket: &DecisionTicket,
3370 controller: &DecisionControllerBinding,
3371 policy: &dyn DecisionPolicy,
3372 ) -> Result<ControllerDecision, DecisionError> {
3373 if !ticket.is_open() {
3374 return Err(DecisionError::new(
3375 DecisionErrorCode::ClosedTicket,
3376 "only open tickets can be evaluated",
3377 ));
3378 }
3379 if ticket.assigned_controller != controller.id || policy.identity() != &controller.policy {
3380 return Err(DecisionError::new(
3381 DecisionErrorCode::PolicyMismatch,
3382 "runtime policy identity does not match the ticket controller binding",
3383 ));
3384 }
3385 let decision = policy.decide(ticket)?;
3386 decision.validate(ticket)?;
3387 if matches!(decision.outcome, DecisionOutcome::Pending { .. }) {
3388 return Ok(ControllerDecision::Pending(decision));
3389 }
3390 let action = match &decision.outcome {
3391 DecisionOutcome::Selected { option_id } => {
3392 ticket.option(option_id).map(|option| option.action.clone())
3393 }
3394 DecisionOutcome::Deferred { .. } | DecisionOutcome::Pending { .. } => None,
3395 };
3396 Ok(ControllerDecision::Authoritative { action, decision })
3397 }
3398}