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