Skip to main content

canwu_sim/runtime/
view.rs

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    /// Loads a package-owned cold object through the host provider attached to
80    /// this runtime. Package code remains responsible for authenticating the
81    /// bytes against its committed archive root before using them.
82    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    /// Returns committed core availability for a person after an explicit
107    /// `canwu.core.person_availability` read. `None` means alive and free.
108    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    /// Finds committed person creations produced by one plugin with an exact
117    /// correlation, so the proposing system can bind engine-allocated IDs at
118    /// a later boundary. The lookup reads the persisted created-person
119    /// registry, so it does not depend on retained evidence.
120    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    /// Counts records in a knowledge namespace at the current proposal-visible cut.
157    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    /// Queries holder-relative records for an omniscient plugin system.
183    ///
184    /// This enforces the declared `canwu.core.knowledge` read and returns an
185    /// owned projection. It is not an actor-facing authorization API.
186    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    /// Resolves an exact command ID from the retained runtime journal in O(1).
224    ///
225    /// An archived command remains valid identity evidence, but its payload is
226    /// no longer available through this view: lookup returns
227    /// [`ErrorCode::EvidenceContentUnavailable`]. `None` means the ID has
228    /// neither retained content nor a committed archive receipt.
229    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    /// Resolves an exact event ID from the retained runtime journal in O(1).
248    ///
249    /// An archived event remains valid identity evidence, but its payload is
250    /// no longer available through this view: lookup returns
251    /// [`ErrorCode::EvidenceContentUnavailable`]. `None` means the ID has
252    /// neither retained content nor a committed archive receipt.
253    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    /// A retained record that is still queued, or that a terminal
304    /// cancellation withdrew, was never admitted and is not durable evidence
305    /// of a delivered packet.
306    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    /// Lists the still-pending plugin ingress that this boundary system's own
315    /// plugin scheduled inside the engine and may withdraw with
316    /// [`crate::BoundaryDirective::CancelPluginIngress`], in ingress-ID order.
317    ///
318    /// Only items due strictly after the view time are listed. Items that a
319    /// boundary directive scheduled earlier in the current boundary appear
320    /// from the next boundary; items generated by this plugin's commands are
321    /// listed as soon as they are queued. Views that are not bound to a
322    /// boundary system list nothing.
323    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    /// Lists the pending transition manifests that this boundary system's
361    /// plugin coordinates or participates in, in manifest-ID order, after an
362    /// explicit `canwu.core.transitions` read.
363    ///
364    /// A manifest's expected versions name other plugins' records, so the
365    /// list is relative to the reading plugin; other plugins' manifests and
366    /// views not bound to a boundary system list nothing. Inside a boundary,
367    /// a manifest registered by an earlier phase is listed from the next phase
368    /// on, and a manifest that settled at phase 11 is no longer listed. A
369    /// participant stages for the manifests whose `ready_at` is the current
370    /// boundary and that list its plugin.
371    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    /// Returns the audits of the transition manifests that settled at phase 11
385    /// of the current boundary and that this boundary system's plugin
386    /// coordinates or participates in, in manifest-ID order, after an
387    /// explicit `canwu.core.transitions` read.
388    ///
389    /// A failed audit fails the boundary instead of leaving a record, so every
390    /// outcome is [`crate::TransitionAuditOutcome::Committed`] or
391    /// [`crate::TransitionAuditOutcome::Expired`]. Earlier phases and views
392    /// not bound to a boundary system see no audits; hosts read settled
393    /// audits from [`crate::BoundaryRecord::transition_audits`] and
394    /// [`crate::BoundaryReceipt::transition_audits`].
395    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    /// Matches retained, plugin-generated ingress provenance without exposing its payload.
408    ///
409    /// Durable evidence may cite ingress admitted at an earlier boundary, but a
410    /// generated record still waiting in the scheduler is not yet admissible.
411    /// Format-7 archive receipts retain a Merkle-bound compact producer proof,
412    /// so archived payload bytes do not need to return to the hot path.
413    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    /// Matches a retained plugin ingress to an exact provider payload and
486    /// delivery time. Payload inspection is deliberately limited to retained
487    /// records: archived receipts prove producer identity, but cannot safely
488    /// be reused to authorize a different legal proposal without the original
489    /// bytes.
490    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    /// Returns the retained outcome for one exact decision request.
557    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    /// Returns one current decision-controller binding after an explicit core read.
566    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    /// Returns one current decision ticket after an explicit core read.
575    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    /// Returns the exact evidence reference assigned to a domain-record
619    /// version proposed earlier in the current boundary.
620    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    /// Same-boundary proposal evidence for one record, without a read check;
629    /// callers gate access before using it.
630    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    /// Returns the exact evidence reference for the currently visible version
645    /// of a domain record.  Strategic aggregation runs after atomic commit,
646    /// therefore it cannot use [`Self::proposed_domain_record_version`].
647    ///
648    /// The current boundary overlay/proposal is preferred, followed by the
649    /// runtime's verified current-record provenance index. The index is
650    /// maintained at commit time and rebuilt from canonical evidence on
651    /// restore, so lookup does not scan retained or archived history, and its
652    /// answer does not depend on how much evidence a seal kept. Either the
653    /// exact record-kind read or the administrative domain-record read grants
654    /// access.
655    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    /// Returns whether an exact domain-record version reference is valid at
689    /// this proposal-visible cut.
690    ///
691    /// This validates both the record identity/version and its establishment
692    /// source. Earlier same-boundary proposals are considered before retained
693    /// or archived runtime evidence. An archived version resolves through its
694    /// own kept receipt. Each seal keeps the receipt of every record's current
695    /// version, including retired and deleted records, and of any other version
696    /// only if a live dependency declares it at that seal. Exact replay always
697    /// finds that other version, so a rule whose outcome must match exact
698    /// replay should use [`Self::replay_stable_domain_record_version`]
699    /// instead. Either the exact record-kind read or the administrative
700    /// domain-record read grants access.
701    pub fn domain_record_version_evidence_exists(
702        &self,
703        reference: &DomainRecordVersionRef,
704    ) -> Result<bool, CanwuError> {
705        // The proposal lookup must not re-require the exact kind, or the
706        // administrative read could never resolve evidence.
707        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    /// Checks that an exact domain-record version is both valid evidence and current.
724    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    /// Returns whether a generic evidence identity is retained or archived.
740    ///
741    /// Domain-record versions proposed earlier in this boundary are visible.
742    /// Archived identities count as existing even when their bodies are no
743    /// longer retained.
744    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    /// Returns when retained or earlier same-boundary evidence first became
768    /// authoritative at this proposal-visible cut.
769    ///
770    /// Archived identity receipts do not retain a precise semantic time, so
771    /// they return `None` where exact replay, which retains every item, returns
772    /// the time. A record's committed current version is the exception: it
773    /// keeps its establishment time after its establishing boundary is sealed.
774    /// A rule whose outcome must match exact replay therefore reads an exact
775    /// version's time through [`Self::replay_stable_domain_record_version`].
776    /// Generic evidence has no such read: after a seal, both its existence
777    /// and its time depend on what the seal kept.
778    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    /// Resolves the retained record body for one exact domain-record version.
795    ///
796    /// Archived receipts prove that a version existed but do not contain its
797    /// body, so this returns `None` when the corresponding evidence segment is
798    /// not live in the runtime. The committed current version resolves to its
799    /// live body even after its establishing boundary is sealed.
800    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    /// Resolves an exact domain-record version whose body and establishment
824    /// time do not depend on how much evidence a seal kept.
825    ///
826    /// Returns them for the committed current version of a record, for a
827    /// version proposed earlier in this boundary and not superseded since, and
828    /// for an initial-scenario version. Any other version returns `None` even while it is retained,
829    /// because `seal_evidence` can remove its existence proof, body, and time
830    /// while exact replay keeps them. A rule whose outcome must match exact
831    /// replay relies on this instead of [`Self::domain_record_version`] and
832    /// [`Self::evidence_time`]. The committed current version still resolves
833    /// after an earlier system in this boundary proposes a newer one. Either
834    /// the exact record-kind read or the administrative domain-record read
835    /// grants access.
836    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        // A proposal still visible in an overlay resolves from it. A change
842        // this boundary has already committed is the current version below.
843        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    /// Returns a bounded, deterministic projection of records of one kind.
873    ///
874    /// Same-boundary overlays take precedence over current state. Records are
875    /// ordered by their canonical reference, so result order is replay stable.
876    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    /// Returns one bounded deterministic page of records after a canonical
885    /// record-reference cursor.
886    ///
887    /// The cursor is exclusive and must name the same kind. This keeps plugin
888    /// scans bounded without imposing a 10,000-record lifetime ceiling on a
889    /// domain kind. Same-boundary overlays retain the same precedence as
890    /// [`Self::domain_records_of_kind`].
891    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    /// Finds committed knowledge changes produced with an exact correlation.
914    ///
915    /// This supports next-boundary operation finalization without granting a
916    /// plugin unrestricted access to unrelated knowledge payloads.
917    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    /// Finds committed knowledge changes whose producer correlation begins
938    /// with a deterministic operation prefix.
939    ///
940    /// The prefix remains plugin-scoped. This is intended for bounded
941    /// multi-holder operation finalization where every holder batch must keep
942    /// a unique full correlation value.
943    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}