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