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, 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    /// Loads a package-owned cold object through the host provider attached to
78    /// this runtime. Package code remains responsible for authenticating the
79    /// bytes against its committed archive root before using them.
80    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    /// Returns committed core availability for a person after an explicit
105    /// `canwu.core.person_availability` read. `None` means alive and free.
106    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    /// Finds committed person creations produced by one plugin with an exact
115    /// correlation, so the proposing system can bind engine-allocated IDs at
116    /// a later boundary. The lookup reads the persisted created-person
117    /// registry, so it does not depend on retained evidence.
118    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    /// Counts records in a knowledge namespace at the current proposal-visible cut.
155    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    /// Queries holder-relative records for an omniscient plugin system.
181    ///
182    /// This enforces the declared `canwu.core.knowledge` read and returns an
183    /// owned projection. It is not an actor-facing authorization API.
184    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    /// Resolves an exact command ID from the retained runtime journal in O(1).
222    ///
223    /// An archived command remains valid identity evidence, but its payload is
224    /// no longer available through this view: lookup returns
225    /// [`ErrorCode::EvidenceContentUnavailable`]. `None` means the ID has
226    /// neither retained content nor a committed archive receipt.
227    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    /// Resolves an exact event ID from the retained runtime journal in O(1).
246    ///
247    /// An archived event remains valid identity evidence, but its payload is
248    /// no longer available through this view: lookup returns
249    /// [`ErrorCode::EvidenceContentUnavailable`]. `None` means the ID has
250    /// neither retained content nor a committed archive receipt.
251    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    /// A retained record that is still queued, or that a terminal
302    /// cancellation withdrew, was never admitted and is not durable evidence
303    /// of a delivered packet.
304    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    /// Lists the still-pending plugin ingress that this boundary system's own
313    /// plugin scheduled inside the engine and may withdraw with
314    /// [`crate::BoundaryDirective::CancelPluginIngress`], in ingress-ID order.
315    ///
316    /// Only items due strictly after the view time are listed. Items that a
317    /// boundary directive scheduled earlier in the current boundary appear
318    /// from the next boundary; items generated by this plugin's commands are
319    /// listed as soon as they are queued. Views that are not bound to a
320    /// boundary system list nothing.
321    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    /// Lists the pending transition manifests that this boundary system's
359    /// plugin coordinates or participates in, in manifest-ID order, after an
360    /// explicit `canwu.core.transitions` read.
361    ///
362    /// A manifest's expected versions name other plugins' records, so the
363    /// list is relative to the reading plugin; other plugins' manifests and
364    /// views not bound to a boundary system list nothing. Inside a boundary,
365    /// a manifest registered by an earlier phase is listed from the next phase
366    /// on, and a manifest that settled at phase 11 is no longer listed. A
367    /// participant stages for the manifests whose `ready_at` is the current
368    /// boundary and that list its plugin.
369    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    /// Returns the audits of the transition manifests that settled at phase 11
383    /// of the current boundary and that this boundary system's plugin
384    /// coordinates or participates in, in manifest-ID order, after an
385    /// explicit `canwu.core.transitions` read.
386    ///
387    /// A failed audit fails the boundary instead of leaving a record, so every
388    /// outcome is [`crate::TransitionAuditOutcome::Committed`] or
389    /// [`crate::TransitionAuditOutcome::Expired`]. Earlier phases and views
390    /// not bound to a boundary system see no audits; hosts read settled
391    /// audits from [`crate::BoundaryRecord::transition_audits`] and
392    /// [`crate::BoundaryReceipt::transition_audits`].
393    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    /// Matches retained, plugin-generated ingress provenance without exposing its payload.
406    ///
407    /// Durable evidence may cite ingress admitted at an earlier boundary, but a
408    /// generated record still waiting in the scheduler is not yet admissible.
409    /// Format-7 archive receipts retain a Merkle-bound compact producer proof,
410    /// so archived payload bytes do not need to return to the hot path.
411    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    /// Matches a retained plugin ingress to an exact provider payload and
484    /// delivery time. Payload inspection is deliberately limited to retained
485    /// records: archived receipts prove producer identity, but cannot safely
486    /// be reused to authorize a different legal proposal without the original
487    /// bytes.
488    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    /// Returns the retained outcome for one exact decision request.
555    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    /// Returns one current decision-controller binding after an explicit core read.
564    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    /// Returns one current decision ticket after an explicit core read.
573    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    /// Returns the exact evidence reference assigned to a domain-record
617    /// version proposed earlier in the current boundary.
618    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    /// Same-boundary proposal evidence for one record, without a read check;
627    /// callers gate access before using it.
628    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    /// Returns the exact evidence reference for the currently visible version
643    /// of a domain record.  Strategic aggregation runs after atomic commit,
644    /// therefore it cannot use [`Self::proposed_domain_record_version`].
645    ///
646    /// The current boundary overlay/proposal is preferred, followed by the
647    /// runtime's verified current-record provenance index. The index is
648    /// maintained at commit time and rebuilt from canonical evidence on
649    /// restore, so lookup does not scan retained or archived history.
650    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    /// Returns whether an exact domain-record version reference is valid at
684    /// this proposal-visible cut.
685    ///
686    /// This validates both the record identity/version and its establishment
687    /// source. Earlier same-boundary proposals are considered before retained
688    /// or archived runtime evidence. Either the exact record-kind read or the
689    /// administrative domain-record read grants access.
690    pub fn domain_record_version_evidence_exists(
691        &self,
692        reference: &DomainRecordVersionRef,
693    ) -> Result<bool, CanwuError> {
694        // The proposal lookup must not re-require the exact kind, or the
695        // administrative read could never resolve evidence.
696        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    /// Checks that an exact domain-record version is both valid evidence and current.
713    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    /// Returns whether a generic evidence identity is retained or archived.
729    ///
730    /// Domain-record versions proposed earlier in this boundary are visible.
731    /// Archived identities count as existing even when their bodies are no
732    /// longer retained.
733    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    /// Returns when retained or earlier same-boundary evidence first became
757    /// authoritative at this proposal-visible cut.
758    ///
759    /// Archived identity receipts do not retain a precise semantic time, so
760    /// they return `None` and callers that require temporal ordering must fail
761    /// closed or load the archived evidence body.
762    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    /// Resolves the retained record body for one exact domain-record version.
779    ///
780    /// Archived receipts prove that a version existed but do not contain its
781    /// body, so this returns `None` when the corresponding evidence segment is
782    /// not live in the runtime.
783    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    /// Returns a bounded, deterministic projection of records of one kind.
807    ///
808    /// Same-boundary overlays take precedence over current state. Records are
809    /// ordered by their canonical reference, so result order is replay stable.
810    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    /// Returns one bounded deterministic page of records after a canonical
819    /// record-reference cursor.
820    ///
821    /// The cursor is exclusive and must name the same kind. This keeps plugin
822    /// scans bounded without imposing a 10,000-record lifetime ceiling on a
823    /// domain kind. Same-boundary overlays retain the same precedence as
824    /// [`Self::domain_records_of_kind`].
825    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    /// Finds committed knowledge changes produced with an exact correlation.
848    ///
849    /// This supports next-boundary operation finalization without granting a
850    /// plugin unrestricted access to unrelated knowledge payloads.
851    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    /// Finds committed knowledge changes whose producer correlation begins
872    /// with a deterministic operation prefix.
873    ///
874    /// The prefix remains plugin-scoped. This is intended for bounded
875    /// multi-holder operation finalization where every holder batch must keep
876    /// a unique full correlation value.
877    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}