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