1use super::{
2 ActorKnowledge, Army, ArmyId, BoundaryId, BoundaryKnowledgeChange, CanwuError, CauseRef,
3 CommandId, CommandRecord, CreatedPerson, DecisionAttemptRecord, DecisionControllerBinding,
4 DecisionRequestId, DecisionTicket, DecisionTicketId, DomainRecord, DomainRecordKind,
5 DomainRecordRef, DomainRecordType, DomainRecordVersionRef, DomainRecordVersionSource,
6 EntityRef, ErrorCode, EventId, EvidenceRef, Government, GovernmentId, HashSet, IngressId,
7 IngressPayload, IngressQueueKey, IngressRecord, KnowledgeHolderRef, KnowledgeQuery,
8 KnowledgeRecord, KnowledgeRecordId, Person, PersonAvailability, PersonId, PluginComponentKey,
9 PluginComponentRecord, RandomOperationTarget, RandomStreamKey, RefCell, ReservationAllocation,
10 ReservationRef, Route, RouteId, RuntimeCurrentState, RuntimeEvidence, RuntimeState, SimEvent,
11 SimTime, StateKey, Territory, TerritoryId, TypedDomainRecordRef, Value, component_key,
12 domain_record_candidates, live_current_domain_record_version, random, records,
13 retained_domain_record_version, retained_evidence_time, validate_domain_record_page_request,
14 validation,
15};
16use std::collections::{BTreeMap, BTreeSet};
17
18pub(super) enum SimulationViewState<'a> {
19 Runtime(&'a RuntimeState),
20 Boundary {
21 current: &'a RuntimeCurrentState,
22 now: SimTime,
23 runtime: &'a RuntimeState,
24 },
25}
26
27impl SimulationViewState<'_> {
28 const fn current(&self) -> &RuntimeCurrentState {
29 match self {
30 Self::Runtime(state) => &state.current,
31 Self::Boundary { current, .. } => current,
32 }
33 }
34
35 const fn now(&self) -> SimTime {
36 match self {
37 Self::Runtime(state) => state.scheduler.now,
38 Self::Boundary { now, .. } => *now,
39 }
40 }
41
42 const fn evidence(&self) -> &RuntimeEvidence {
43 match self {
44 Self::Runtime(state) => &state.evidence,
45 Self::Boundary { runtime, .. } => &runtime.evidence,
46 }
47 }
48
49 const fn runtime(&self) -> &RuntimeState {
50 match self {
51 Self::Runtime(state) | Self::Boundary { runtime: state, .. } => state,
52 }
53 }
54}
55
56pub struct SimulationView<'a> {
57 pub(super) state: SimulationViewState<'a>,
58 pub(super) state_owners: &'a BTreeMap<StateKey, String>,
59 pub(super) reader: Option<&'a str>,
60 pub(super) allowed_reads: Option<&'a [StateKey]>,
61 pub(super) allowed_ingress: Option<&'a HashSet<IngressId>>,
62 pub(super) ingress_plugin: Option<&'a str>,
63 pub(super) component_overlay: Option<&'a BTreeMap<PluginComponentKey, PluginComponentRecord>>,
64 pub(super) proposed_components: Option<&'a BTreeMap<PluginComponentKey, PluginComponentRecord>>,
65 pub(super) record_overlay: Option<&'a BTreeMap<DomainRecordRef, DomainRecord>>,
66 pub(super) proposed_records: Option<&'a BTreeMap<DomainRecordRef, DomainRecord>>,
67 pub(super) boundary_id: Option<BoundaryId>,
68 pub(super) proposal_evidence: Option<&'a BTreeSet<EvidenceRef>>,
69 pub(super) knowledge_overlay:
70 Option<&'a BTreeMap<KnowledgeHolderRef, BTreeMap<KnowledgeRecordId, KnowledgeRecord>>>,
71 pub(super) allocations: Option<&'a BTreeMap<ReservationRef, ReservationAllocation>>,
72 pub(super) allowed_reservations: Option<&'a [ReservationRef]>,
73 pub(super) random_session: Option<RefCell<random::RandomSession>>,
74 pub(super) plugin_archive_provider: &'a dyn super::PluginArchiveObjectProvider,
75 pub(super) transitions: Option<&'a super::transitions::BoundaryTransitionLedger>,
76}
77
78impl SimulationView<'_> {
79 pub fn plugin_archive_object(
83 &self,
84 namespace: &str,
85 object_id: &str,
86 ) -> Result<Option<Vec<u8>>, CanwuError> {
87 self.plugin_archive_provider
88 .load_plugin_archive_object(namespace, object_id)
89 }
90
91 #[must_use]
92 pub const fn time(&self) -> SimTime {
93 self.state.now()
94 }
95
96 pub fn army(&self, id: ArmyId) -> Result<Option<&Army>, CanwuError> {
97 self.require_read(&StateKey::core_armies())?;
98 Ok(self.state.current().armies.get(&id))
99 }
100
101 pub fn person(&self, id: PersonId) -> Result<Option<&Person>, CanwuError> {
102 self.require_read(&StateKey::core_people())?;
103 Ok(self.state.current().people.get(&id))
104 }
105
106 pub fn person_availability(
109 &self,
110 id: PersonId,
111 ) -> Result<Option<&PersonAvailability>, CanwuError> {
112 self.require_read(&StateKey::core_person_availability())?;
113 Ok(self.state.current().person_availability.get(&id))
114 }
115
116 pub fn persons_created_by_correlation(
121 &self,
122 plugin: &str,
123 correlation: &str,
124 ) -> Result<Vec<CreatedPerson>, CanwuError> {
125 self.require_read(&StateKey::core_people())?;
126 Ok(self
127 .state
128 .current()
129 .created_persons
130 .iter()
131 .filter(|created| created.plugin == plugin && created.correlation == correlation)
132 .cloned()
133 .collect())
134 }
135
136 pub fn government(&self, id: GovernmentId) -> Result<Option<&Government>, CanwuError> {
137 self.require_read(&StateKey::core_governments())?;
138 Ok(self.state.current().governments.get(&id))
139 }
140
141 pub fn territory(&self, id: TerritoryId) -> Result<Option<&Territory>, CanwuError> {
142 self.require_read(&StateKey::core_territories())?;
143 Ok(self.state.current().territories.get(&id))
144 }
145
146 pub fn route(&self, id: RouteId) -> Result<Option<&Route>, CanwuError> {
147 self.require_read(&StateKey::core_routes())?;
148 Ok(self.state.current().routes.get(&id))
149 }
150
151 pub fn actor_knowledge(&self, actor: PersonId) -> Result<Option<&ActorKnowledge>, CanwuError> {
152 self.require_read(&StateKey::core_knowledge())?;
153 Ok(self.state.current().knowledge.for_actor(actor))
154 }
155
156 pub fn knowledge_record_count_in_namespace(
158 &self,
159 namespace: &str,
160 ) -> Result<usize, CanwuError> {
161 self.require_read(&StateKey::core_knowledge())?;
162 let settled = self
163 .state
164 .current()
165 .knowledge
166 .record_count_in_namespace(namespace);
167 let proposed = self.knowledge_overlay.map_or(0, |overlay| {
168 overlay
169 .values()
170 .flat_map(BTreeMap::values)
171 .filter(|record| record.schema.kind.namespace == namespace)
172 .count()
173 });
174 settled.checked_add(proposed).ok_or_else(|| {
175 CanwuError::new(
176 ErrorCode::ValueOutOfRange,
177 "knowledge namespace record count overflowed",
178 )
179 })
180 }
181
182 pub fn knowledge_records(
187 &self,
188 holder: KnowledgeHolderRef,
189 query: &KnowledgeQuery,
190 ) -> Result<canwu_knowledge::KnowledgeQueryResult, CanwuError> {
191 self.require_read(&StateKey::core_knowledge())?;
192 let result = if let Some(overlay) = self.knowledge_overlay {
193 self.state.current().knowledge.query_with_overlay(
194 holder,
195 query,
196 self.boundary_id,
197 overlay,
198 )
199 } else {
200 self.state
201 .current()
202 .knowledge
203 .query_current(holder, query, self.boundary_id)
204 };
205 result.map_err(|error| match error {
206 canwu_knowledge::KnowledgeQueryError::ReadCutUnavailable => CanwuError::new(
207 ErrorCode::KnowledgeReadCutUnavailable,
208 "knowledge cursor read cut is no longer available",
209 ),
210 canwu_knowledge::KnowledgeQueryError::InvalidLimit => CanwuError::new(
211 ErrorCode::KnowledgeLimitExceeded,
212 "knowledge query page size is outside the supported range",
213 ),
214 canwu_knowledge::KnowledgeQueryError::InvalidCursor
215 | canwu_knowledge::KnowledgeQueryError::InvalidLedger
216 | canwu_knowledge::KnowledgeQueryError::Encoding => CanwuError::new(
217 ErrorCode::InvalidKnowledgeRecord,
218 "knowledge query, cursor, or ledger is invalid",
219 ),
220 })
221 }
222
223 pub fn command(&self, id: CommandId) -> Result<Option<&CommandRecord>, CanwuError> {
230 self.require_read(&StateKey::core_commands())?;
231 let retained = self.state.evidence().retained_command(id);
232 if retained.is_none()
233 && self
234 .state
235 .evidence()
236 .archived_evidence_receipts
237 .contains_key(&EvidenceRef::Command(id))
238 {
239 return Err(CanwuError::new(
240 ErrorCode::EvidenceContentUnavailable,
241 "command identity is archived; payload inspection requires an archive provider",
242 ));
243 }
244 Ok(retained)
245 }
246
247 pub fn event(&self, id: EventId) -> Result<Option<&SimEvent>, CanwuError> {
254 self.require_read(&StateKey::core_events())?;
255 let retained = self.state.evidence().retained_event(id);
256 if retained.is_none()
257 && self
258 .state
259 .evidence()
260 .archived_evidence_receipts
261 .contains_key(&EvidenceRef::Event(id))
262 {
263 return Err(CanwuError::new(
264 ErrorCode::EvidenceContentUnavailable,
265 "event identity is archived; payload inspection requires an archive provider",
266 ));
267 }
268 Ok(retained)
269 }
270
271 pub fn ingress(&self, id: IngressId) -> Result<Option<&IngressRecord>, CanwuError> {
272 self.require_read(&StateKey::core_ingress())?;
273 if self
274 .allowed_ingress
275 .is_none_or(|allowed| !allowed.contains(&id))
276 {
277 return Ok(None);
278 }
279 let record = self.state.evidence().retained_ingress(id);
280 if record.is_none()
281 && self
282 .state
283 .evidence()
284 .archived_evidence_receipts
285 .contains_key(&EvidenceRef::Ingress(id))
286 {
287 return Err(CanwuError::new(
288 ErrorCode::EvidenceContentUnavailable,
289 "ingress identity is archived; payload inspection requires an archive provider",
290 ));
291 }
292 if let (Some(owner), Some(record)) = (self.ingress_plugin, record)
293 && !matches!(
294 &record.payload,
295 IngressPayload::Plugin { plugin, .. } if plugin == owner
296 )
297 {
298 return Ok(None);
299 }
300 Ok(record)
301 }
302
303 fn ingress_is_pending_or_cancelled(&self, record: &IngressRecord) -> bool {
307 let scheduler = &self.state.runtime().scheduler;
308 scheduler.cancelled_ingress.contains(&record.id)
309 || scheduler
310 .pending_ingress
311 .contains(&IngressQueueKey::from_record(record))
312 }
313
314 pub fn cancellable_plugin_ingress(&self) -> Result<Vec<&IngressRecord>, CanwuError> {
324 self.require_read(&StateKey::core_ingress())?;
325 let Some(owner) = self.ingress_plugin else {
326 return Ok(Vec::new());
327 };
328 let now = self.state.now();
329 let evidence = self.state.evidence();
330 let mut records = Vec::new();
331 for key in &self.state.runtime().scheduler.pending_ingress {
332 if key.due_at <= now {
333 continue;
334 }
335 let Some(record) = evidence.retained_ingress(key.id) else {
336 continue;
337 };
338 let IngressPayload::Plugin { plugin, .. } = &record.payload else {
339 continue;
340 };
341 let issued_by_owner = match &record.cause {
342 Some(CauseRef::Command(_)) => plugin == owner,
343 Some(CauseRef::Boundary(boundary)) => evidence
344 .retained_boundary(*boundary)
345 .is_some_and(|boundary| {
346 boundary.generated_ingress.iter().any(|generation| {
347 generation.ingress == record.id && generation.plugin == owner
348 })
349 }),
350 Some(CauseRef::System(_) | CauseRef::Event(_)) | None => false,
351 };
352 if issued_by_owner {
353 records.push(record);
354 }
355 }
356 records.sort_by_key(|record| record.id);
357 Ok(records)
358 }
359
360 pub fn transition_manifests(
372 &self,
373 ) -> Result<Vec<&super::PendingTransitionManifest>, CanwuError> {
374 self.require_read(&StateKey::core_transitions())?;
375 let (Some(reader), Some(ledger)) = (self.ingress_plugin, self.transitions) else {
376 return Ok(Vec::new());
377 };
378 Ok(ledger
379 .pending()
380 .filter(|manifest| manifest.involves(reader))
381 .collect())
382 }
383
384 pub fn transition_audits(&self) -> Result<Vec<&super::TransitionAuditRecord>, CanwuError> {
396 self.require_read(&StateKey::core_transitions())?;
397 let (Some(reader), Some(ledger)) = (self.ingress_plugin, self.transitions) else {
398 return Ok(Vec::new());
399 };
400 Ok(ledger
401 .audits()
402 .iter()
403 .filter(|audit| audit.involves(reader))
404 .collect())
405 }
406
407 pub fn plugin_ingress_matches(
414 &self,
415 id: IngressId,
416 plugin: &str,
417 packet_type: &str,
418 ) -> Result<bool, CanwuError> {
419 self.require_read(&StateKey::core_ingress())?;
420 let record = self.state.evidence().retained_ingress(id);
421 if record.is_none()
422 && let Some(receipt) = self
423 .state
424 .evidence()
425 .archived_evidence_receipts
426 .get(&EvidenceRef::Ingress(id))
427 {
428 if self
429 .state
430 .runtime()
431 .scheduler
432 .pending_ingress
433 .iter()
434 .any(|key| key.id == id)
435 {
436 return Ok(false);
437 }
438 return Ok(receipt
439 .plugin_ingress_provenance
440 .as_ref()
441 .is_some_and(|provenance| {
442 provenance.plugin == plugin && provenance.packet_type == packet_type
443 }));
444 }
445 let Some(record) = record else {
446 return Ok(false);
447 };
448 if self.ingress_is_pending_or_cancelled(record) {
449 return Ok(false);
450 }
451 if !matches!(
452 &record.payload,
453 IngressPayload::Plugin {
454 plugin: actual_plugin,
455 packet_type: actual_packet_type,
456 ..
457 } if actual_plugin == plugin && actual_packet_type == packet_type
458 ) {
459 return Ok(false);
460 }
461 let Some(CauseRef::Boundary(boundary_id)) = record.cause.as_ref() else {
462 return Ok(false);
463 };
464 let boundary = self.state.evidence().retained_boundary(*boundary_id);
465 if boundary.is_none()
466 && self
467 .state
468 .evidence()
469 .archived_evidence_receipts
470 .contains_key(&EvidenceRef::Boundary(*boundary_id))
471 {
472 return Err(CanwuError::new(
473 ErrorCode::EvidenceContentUnavailable,
474 "ingress producer boundary is archived; provenance inspection requires an archive provider",
475 ));
476 }
477 Ok(boundary.is_some_and(|boundary| {
478 boundary
479 .generated_ingress
480 .iter()
481 .any(|generation| generation.ingress == id && generation.plugin == plugin)
482 }))
483 }
484
485 pub fn plugin_ingress_payload_matches(
491 &self,
492 id: IngressId,
493 plugin: &str,
494 packet_type: &str,
495 occurred_at: SimTime,
496 expected_payload: &Value,
497 ) -> Result<bool, CanwuError> {
498 self.require_read(&StateKey::core_ingress())?;
499 let Some(record) = self.state.evidence().retained_ingress(id) else {
500 if self
501 .state
502 .evidence()
503 .archived_evidence_receipts
504 .contains_key(&EvidenceRef::Ingress(id))
505 {
506 return Err(CanwuError::new(
507 ErrorCode::EvidenceContentUnavailable,
508 "provider ingress payload is archived; exact legal signal binding requires retained content",
509 ));
510 }
511 return Ok(false);
512 };
513 if self.ingress_is_pending_or_cancelled(record) {
514 return Ok(false);
515 }
516 let IngressPayload::Plugin {
517 plugin: actual_plugin,
518 packet_type: actual_packet_type,
519 payload,
520 ..
521 } = &record.payload
522 else {
523 return Ok(false);
524 };
525 if actual_plugin != plugin
526 || actual_packet_type != packet_type
527 || record.due_at != occurred_at
528 || payload != expected_payload
529 {
530 return Ok(false);
531 }
532 let Some(CauseRef::Boundary(boundary_id)) = record.cause.as_ref() else {
533 return Ok(false);
534 };
535 let boundary = self.state.evidence().retained_boundary(*boundary_id);
536 if boundary.is_none()
537 && self
538 .state
539 .evidence()
540 .archived_evidence_receipts
541 .contains_key(&EvidenceRef::Boundary(*boundary_id))
542 {
543 return Err(CanwuError::new(
544 ErrorCode::EvidenceContentUnavailable,
545 "provider ingress producer boundary is archived; exact legal signal binding requires retained content",
546 ));
547 }
548 Ok(boundary.is_some_and(|boundary| {
549 boundary
550 .generated_ingress
551 .iter()
552 .any(|generation| generation.ingress == id && generation.plugin == plugin)
553 }))
554 }
555
556 pub fn decision_attempt(
558 &self,
559 request_id: DecisionRequestId,
560 ) -> Result<Option<&DecisionAttemptRecord>, CanwuError> {
561 self.require_read(&StateKey::core_decisions())?;
562 Ok(self.state.current().decisions.attempt(request_id))
563 }
564
565 pub fn decision_controller(
567 &self,
568 id: &str,
569 ) -> Result<Option<&DecisionControllerBinding>, CanwuError> {
570 self.require_read(&StateKey::core_decisions())?;
571 Ok(self.state.current().decisions.controller(id))
572 }
573
574 pub fn decision_ticket(
576 &self,
577 id: DecisionTicketId,
578 ) -> Result<Option<&DecisionTicket>, CanwuError> {
579 self.require_read(&StateKey::core_decisions())?;
580 Ok(self.state.current().decisions.ticket(id))
581 }
582
583 pub fn domain_record(
584 &self,
585 reference: &DomainRecordRef,
586 ) -> Result<Option<&DomainRecord>, CanwuError> {
587 self.require_domain_record_read(reference)?;
588 Ok(self
589 .record_overlay
590 .and_then(|overlay| overlay.get(reference))
591 .or_else(|| self.state.current().domain_records.get(reference)))
592 }
593
594 pub fn typed_domain_record<T: DomainRecordType>(
595 &self,
596 reference: &TypedDomainRecordRef<T>,
597 ) -> Result<Option<&DomainRecord>, CanwuError> {
598 self.domain_record(reference.as_untyped())
599 }
600
601 pub fn proposed_domain_record(
602 &self,
603 reference: &DomainRecordRef,
604 ) -> Result<Option<&DomainRecord>, CanwuError> {
605 self.require_read(&records::record_state_key(&reference.kind))?;
606 Ok(self
607 .proposed_records
608 .and_then(|records| records.get(reference)))
609 }
610
611 pub fn proposed_typed_domain_record<T: DomainRecordType>(
612 &self,
613 reference: &TypedDomainRecordRef<T>,
614 ) -> Result<Option<&DomainRecord>, CanwuError> {
615 self.proposed_domain_record(reference.as_untyped())
616 }
617
618 pub fn proposed_domain_record_version(
621 &self,
622 reference: &DomainRecordRef,
623 ) -> Result<Option<DomainRecordVersionRef>, CanwuError> {
624 self.require_read(&records::record_state_key(&reference.kind))?;
625 Ok(self.proposed_version_evidence(reference))
626 }
627
628 fn proposed_version_evidence(
631 &self,
632 reference: &DomainRecordRef,
633 ) -> Option<DomainRecordVersionRef> {
634 self.proposal_evidence.and_then(|evidence| {
635 evidence.iter().find_map(|item| match item {
636 EvidenceRef::DomainRecordVersion(version) if version.record == *reference => {
637 Some(version.clone())
638 }
639 _ => None,
640 })
641 })
642 }
643
644 pub fn current_domain_record_version(
656 &self,
657 reference: &DomainRecordRef,
658 ) -> Result<Option<DomainRecordVersionRef>, CanwuError> {
659 self.require_domain_record_read(reference)?;
660 let Some(record) = self.domain_record(reference)? else {
661 return Ok(None);
662 };
663 if let Some(proposed) = self.proposal_evidence.and_then(|evidence| {
664 evidence.iter().find_map(|item| match item {
665 EvidenceRef::DomainRecordVersion(version)
666 if version.record == *reference && version.version == record.version =>
667 {
668 Some(version.clone())
669 }
670 _ => None,
671 })
672 }) {
673 return Ok(Some(proposed));
674 }
675 let current = super::current_domain_record_version(self.state.runtime(), reference)?;
676 if current
677 .as_ref()
678 .is_some_and(|current| current.version != record.version)
679 {
680 return Err(CanwuError::new(
681 ErrorCode::InvalidSnapshot,
682 "visible domain-record version disagrees with the runtime provenance index",
683 ));
684 }
685 Ok(current)
686 }
687
688 pub fn domain_record_version_evidence_exists(
702 &self,
703 reference: &DomainRecordVersionRef,
704 ) -> Result<bool, CanwuError> {
705 self.require_domain_record_read(&reference.record)?;
708 if self
709 .proposed_version_evidence(&reference.record)
710 .is_some_and(|proposed| proposed == *reference)
711 {
712 return Ok(true);
713 }
714 Ok(!matches!(
715 validation::resolve_evidence_reference(
716 &validation::RuntimeValidationContext::new(self.state.runtime()),
717 &EvidenceRef::DomainRecordVersion(reference.clone()),
718 ),
719 validation::EvidenceAvailability::Missing
720 ))
721 }
722
723 pub fn domain_record_version_is_current(
725 &self,
726 reference: &DomainRecordVersionRef,
727 ) -> Result<bool, CanwuError> {
728 self.require_domain_record_read(&reference.record)?;
729 let current = self
730 .record_overlay
731 .and_then(|overlay| overlay.get(&reference.record))
732 .or_else(|| self.state.current().domain_records.get(&reference.record));
733 Ok(
734 current.is_some_and(|record| record.version == reference.version)
735 && self.domain_record_version_evidence_exists(reference)?,
736 )
737 }
738
739 pub fn evidence_exists(&self, reference: &EvidenceRef) -> Result<bool, CanwuError> {
745 match reference {
746 EvidenceRef::Command(_) | EvidenceRef::CommandAttempt(_) => {
747 self.require_read(&StateKey::core_commands())?;
748 }
749 EvidenceRef::Event(_) => self.require_read(&StateKey::core_events())?,
750 EvidenceRef::Ingress(_) => self.require_read(&StateKey::core_ingress())?,
751 EvidenceRef::Boundary(_) | EvidenceRef::RandomDraw(_) => {
752 self.require_read(&StateKey::core_evidence())?;
753 }
754 EvidenceRef::DomainRecordVersion(version) => {
755 return self.domain_record_version_evidence_exists(version);
756 }
757 }
758 Ok(!matches!(
759 validation::resolve_evidence_reference(
760 &validation::RuntimeValidationContext::new(self.state.runtime()),
761 reference,
762 ),
763 validation::EvidenceAvailability::Missing
764 ))
765 }
766
767 pub fn evidence_time(&self, reference: &EvidenceRef) -> Result<Option<SimTime>, CanwuError> {
779 if !self.evidence_exists(reference)? {
780 return Ok(None);
781 }
782 if self
783 .proposal_evidence
784 .is_some_and(|evidence| evidence.contains(reference))
785 {
786 return Ok(Some(self.time()));
787 }
788 Ok(super::retained_evidence_time(
789 self.state.runtime(),
790 reference,
791 ))
792 }
793
794 pub fn domain_record_version(
801 &self,
802 reference: &DomainRecordVersionRef,
803 ) -> Result<Option<DomainRecord>, CanwuError> {
804 self.require_read(&records::record_state_key(&reference.record.kind))?;
805 if let Some(proposed) = self.proposed_domain_record_version(&reference.record)?
806 && proposed == *reference
807 {
808 return Ok(self
809 .proposed_records
810 .and_then(|records| records.get(&reference.record))
811 .or_else(|| {
812 self.record_overlay
813 .and_then(|records| records.get(&reference.record))
814 })
815 .cloned());
816 }
817 Ok(retained_domain_record_version(
818 self.state.runtime(),
819 reference,
820 ))
821 }
822
823 pub fn replay_stable_domain_record_version(
837 &self,
838 reference: &DomainRecordVersionRef,
839 ) -> Result<Option<(DomainRecord, SimTime)>, CanwuError> {
840 self.require_domain_record_read(&reference.record)?;
841 if self.proposal_evidence.is_some_and(|evidence| {
844 evidence.contains(&EvidenceRef::DomainRecordVersion(reference.clone()))
845 }) && let Some(record) = [self.proposed_records, self.record_overlay]
846 .into_iter()
847 .flatten()
848 .find_map(|records| {
849 records
850 .get(&reference.record)
851 .filter(|record| record.version == reference.version)
852 })
853 {
854 return Ok(Some((record.clone(), self.time())));
855 }
856 let runtime = self.state.runtime();
857 if !matches!(
858 reference.established_by,
859 DomainRecordVersionSource::InitialScenario
860 ) && live_current_domain_record_version(runtime, reference).is_none()
861 {
862 return Ok(None);
863 }
864 Ok(
865 retained_domain_record_version(runtime, reference).zip(retained_evidence_time(
866 runtime,
867 &EvidenceRef::DomainRecordVersion(reference.clone()),
868 )),
869 )
870 }
871
872 pub fn domain_records_of_kind(
877 &self,
878 kind: &DomainRecordKind,
879 limit: usize,
880 ) -> Result<Vec<DomainRecord>, CanwuError> {
881 self.domain_records_of_kind_after(kind, None, limit)
882 }
883
884 pub fn domain_records_of_kind_after(
892 &self,
893 kind: &DomainRecordKind,
894 after: Option<&DomainRecordRef>,
895 limit: usize,
896 ) -> Result<Vec<DomainRecord>, CanwuError> {
897 self.require_read(&records::record_state_key(kind))?;
898 validate_domain_record_page_request(kind, after, limit)?;
899
900 let mut records =
901 domain_record_candidates(&self.state.current().domain_records, kind, after, limit);
902 for overlay in [self.record_overlay, self.proposed_records]
903 .into_iter()
904 .flatten()
905 {
906 for (reference, record) in domain_record_candidates(overlay, kind, after, limit) {
907 records.insert(reference, record);
908 }
909 }
910 Ok(records.into_values().take(limit).collect())
911 }
912
913 pub fn knowledge_changes_by_correlation(
918 &self,
919 plugin: &str,
920 producer_correlation: &str,
921 ) -> Result<Vec<BoundaryKnowledgeChange>, CanwuError> {
922 self.require_read(&StateKey::core_knowledge())?;
923 Ok(self
924 .state
925 .evidence()
926 .boundaries
927 .iter()
928 .flat_map(|boundary| &boundary.knowledge_changes)
929 .filter(|change| {
930 change.plugin == plugin
931 && change.producer_correlation.as_deref() == Some(producer_correlation)
932 })
933 .cloned()
934 .collect())
935 }
936
937 pub fn knowledge_changes_by_correlation_prefix(
944 &self,
945 plugin: &str,
946 producer_correlation_prefix: &str,
947 ) -> Result<Vec<BoundaryKnowledgeChange>, CanwuError> {
948 self.require_read(&StateKey::core_knowledge())?;
949 Ok(self
950 .state
951 .evidence()
952 .boundaries
953 .iter()
954 .flat_map(|boundary| &boundary.knowledge_changes)
955 .filter(|change| {
956 change.plugin == plugin
957 && change
958 .producer_correlation
959 .as_deref()
960 .is_some_and(|value| value.starts_with(producer_correlation_prefix))
961 })
962 .cloned()
963 .collect())
964 }
965
966 pub fn reservation(
967 &self,
968 reservation: &ReservationRef,
969 ) -> Result<Option<&ReservationAllocation>, CanwuError> {
970 let reader = self.reader.unwrap_or("unscoped caller");
971 if self
972 .allowed_reservations
973 .is_none_or(|allowed| !allowed.contains(reservation))
974 {
975 return Err(CanwuError::new(
976 ErrorCode::UndeclaredStateRead,
977 format!(
978 "system {reader} did not declare reservation read {}.{}.{}",
979 reservation.plugin, reservation.system, reservation.request
980 ),
981 ));
982 }
983 Ok(self.allocations.and_then(|values| values.get(reservation)))
984 }
985
986 pub fn random_range(
987 &self,
988 stream: &RandomStreamKey,
989 upper_exclusive: u64,
990 purpose: &str,
991 ) -> Result<u64, CanwuError> {
992 let Some(session) = &self.random_session else {
993 return Err(CanwuError::new(
994 ErrorCode::UndeclaredRandomStream,
995 format!(
996 "system {} has no declared random streams",
997 self.reader.unwrap_or("unscoped caller")
998 ),
999 ));
1000 };
1001 session.borrow_mut().range(stream, upper_exclusive, purpose)
1002 }
1003
1004 #[allow(clippy::too_many_arguments)]
1005 pub fn random_range_for_operation(
1006 &self,
1007 stream: &RandomStreamKey,
1008 evidence: EvidenceRef,
1009 operation_kind: &str,
1010 application_operation_id: &str,
1011 target: RandomOperationTarget,
1012 draw_slot: u32,
1013 upper_exclusive: u64,
1014 purpose: &str,
1015 ) -> Result<u64, CanwuError> {
1016 self.random_sample_for_operation(
1017 stream,
1018 evidence,
1019 operation_kind,
1020 application_operation_id,
1021 target,
1022 draw_slot,
1023 upper_exclusive,
1024 purpose,
1025 )
1026 .map(|sample| sample.value)
1027 }
1028
1029 #[allow(clippy::too_many_arguments)]
1030 pub fn random_sample_for_operation(
1031 &self,
1032 stream: &RandomStreamKey,
1033 evidence: EvidenceRef,
1034 operation_kind: &str,
1035 application_operation_id: &str,
1036 target: RandomOperationTarget,
1037 draw_slot: u32,
1038 upper_exclusive: u64,
1039 purpose: &str,
1040 ) -> Result<super::RandomSample, CanwuError> {
1041 let available = self
1042 .proposal_evidence
1043 .is_some_and(|values| values.contains(&evidence))
1044 || validation::resolve_evidence_reference(
1045 &validation::RuntimeValidationContext::new(self.state.runtime()),
1046 &evidence,
1047 ) == validation::EvidenceAvailability::Retained;
1048 if !available {
1049 return Err(CanwuError::new(
1050 ErrorCode::InvalidRandomOperationEvidence,
1051 "operation-keyed random draw references unavailable evidence",
1052 ));
1053 }
1054 let Some(session) = &self.random_session else {
1055 return Err(CanwuError::new(
1056 ErrorCode::UndeclaredRandomStream,
1057 format!(
1058 "system {} has no declared random streams",
1059 self.reader.unwrap_or("unscoped caller")
1060 ),
1061 ));
1062 };
1063 session.borrow_mut().sample_for_operation(
1064 stream,
1065 evidence,
1066 operation_kind,
1067 application_operation_id,
1068 target,
1069 draw_slot,
1070 upper_exclusive,
1071 purpose,
1072 )
1073 }
1074
1075 pub fn component(
1076 &self,
1077 state: &StateKey,
1078 entity: &EntityRef,
1079 component: &str,
1080 ) -> Result<Option<&Value>, CanwuError> {
1081 self.require_read(state)?;
1082 let Some(owner) = self.state_owners.get(state) else {
1083 return Err(CanwuError::new(
1084 ErrorCode::UndeclaredStateRead,
1085 format!(
1086 "state {}.{} has no registered owner",
1087 state.namespace, state.name
1088 ),
1089 ));
1090 };
1091 let key = component_key(owner, state, entity, component);
1092 Ok(self
1093 .component_overlay
1094 .and_then(|overlay| overlay.get(&key))
1095 .or_else(|| self.state.current().plugin_components.get(&key))
1096 .map(|record| &record.value))
1097 }
1098
1099 pub fn proposed_component(
1100 &self,
1101 state: &StateKey,
1102 entity: &EntityRef,
1103 component: &str,
1104 ) -> Result<Option<&Value>, CanwuError> {
1105 self.require_read(state)?;
1106 let Some(owner) = self.state_owners.get(state) else {
1107 return Err(CanwuError::new(
1108 ErrorCode::UndeclaredStateRead,
1109 format!(
1110 "state {}.{} has no registered owner",
1111 state.namespace, state.name
1112 ),
1113 ));
1114 };
1115 let key = component_key(owner, state, entity, component);
1116 Ok(self
1117 .proposed_components
1118 .and_then(|proposals| proposals.get(&key))
1119 .map(|record| &record.value))
1120 }
1121
1122 fn require_read(&self, state: &StateKey) -> Result<(), CanwuError> {
1123 if self
1124 .allowed_reads
1125 .is_some_and(|reads| !reads.contains(state))
1126 {
1127 return Err(CanwuError::new(
1128 ErrorCode::UndeclaredStateRead,
1129 format!(
1130 "{} did not declare read access to {}.{}",
1131 self.reader.unwrap_or("internal system"),
1132 state.namespace,
1133 state.name
1134 ),
1135 ));
1136 }
1137 Ok(())
1138 }
1139
1140 fn require_domain_record_read(&self, reference: &DomainRecordRef) -> Result<(), CanwuError> {
1141 let exact = records::record_state_key(&reference.kind);
1142 if self.allowed_reads.is_some_and(|reads| {
1143 !reads.contains(&exact) && !reads.contains(&StateKey::core_domain_records())
1144 }) {
1145 return Err(CanwuError::new(
1146 ErrorCode::UndeclaredStateRead,
1147 format!(
1148 "{} did not declare read access to {}.{}",
1149 self.reader.unwrap_or("internal system"),
1150 exact.namespace,
1151 exact.name
1152 ),
1153 ));
1154 }
1155 Ok(())
1156 }
1157
1158 pub(super) fn finish_random_session(self) -> Option<random::RandomExecution> {
1159 self.random_session
1160 .map(RefCell::into_inner)
1161 .map(random::RandomSession::finish)
1162 }
1163}
1164
1165impl super::PluginArchiveObjectProvider for SimulationView<'_> {
1166 fn load_plugin_archive_object(
1167 &self,
1168 namespace: &str,
1169 object_id: &str,
1170 ) -> Result<Option<Vec<u8>>, CanwuError> {
1171 self.plugin_archive_object(namespace, object_id)
1172 }
1173}