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}
74
75impl SimulationView<'_> {
76    /// Loads a package-owned cold object through the host provider attached to
77    /// this runtime. Package code remains responsible for authenticating the
78    /// bytes against its committed archive root before using them.
79    pub fn plugin_archive_object(
80        &self,
81        namespace: &str,
82        object_id: &str,
83    ) -> Result<Option<Vec<u8>>, CanwuError> {
84        self.plugin_archive_provider
85            .load_plugin_archive_object(namespace, object_id)
86    }
87
88    #[must_use]
89    pub const fn time(&self) -> SimTime {
90        self.state.now()
91    }
92
93    pub fn army(&self, id: ArmyId) -> Result<Option<&Army>, CanwuError> {
94        self.require_read(&StateKey::core_armies())?;
95        Ok(self.state.current().armies.get(&id))
96    }
97
98    pub fn person(&self, id: PersonId) -> Result<Option<&Person>, CanwuError> {
99        self.require_read(&StateKey::core_people())?;
100        Ok(self.state.current().people.get(&id))
101    }
102
103    /// Returns committed core availability for a person after an explicit
104    /// `canwu.core.person_availability` read. `None` means alive and free.
105    pub fn person_availability(
106        &self,
107        id: PersonId,
108    ) -> Result<Option<&PersonAvailability>, CanwuError> {
109        self.require_read(&StateKey::core_person_availability())?;
110        Ok(self.state.current().person_availability.get(&id))
111    }
112
113    /// Finds committed person creations produced by one plugin with an exact
114    /// correlation, so the proposing system can bind engine-allocated IDs at
115    /// a later boundary. The lookup reads the persisted created-person
116    /// registry, so it does not depend on retained evidence.
117    pub fn persons_created_by_correlation(
118        &self,
119        plugin: &str,
120        correlation: &str,
121    ) -> Result<Vec<CreatedPerson>, CanwuError> {
122        self.require_read(&StateKey::core_people())?;
123        Ok(self
124            .state
125            .current()
126            .created_persons
127            .iter()
128            .filter(|created| created.plugin == plugin && created.correlation == correlation)
129            .cloned()
130            .collect())
131    }
132
133    pub fn government(&self, id: GovernmentId) -> Result<Option<&Government>, CanwuError> {
134        self.require_read(&StateKey::core_governments())?;
135        Ok(self.state.current().governments.get(&id))
136    }
137
138    pub fn territory(&self, id: TerritoryId) -> Result<Option<&Territory>, CanwuError> {
139        self.require_read(&StateKey::core_territories())?;
140        Ok(self.state.current().territories.get(&id))
141    }
142
143    pub fn route(&self, id: RouteId) -> Result<Option<&Route>, CanwuError> {
144        self.require_read(&StateKey::core_routes())?;
145        Ok(self.state.current().routes.get(&id))
146    }
147
148    pub fn actor_knowledge(&self, actor: PersonId) -> Result<Option<&ActorKnowledge>, CanwuError> {
149        self.require_read(&StateKey::core_knowledge())?;
150        Ok(self.state.current().knowledge.for_actor(actor))
151    }
152
153    /// Counts records in a knowledge namespace at the current proposal-visible cut.
154    pub fn knowledge_record_count_in_namespace(
155        &self,
156        namespace: &str,
157    ) -> Result<usize, CanwuError> {
158        self.require_read(&StateKey::core_knowledge())?;
159        let settled = self
160            .state
161            .current()
162            .knowledge
163            .record_count_in_namespace(namespace);
164        let proposed = self.knowledge_overlay.map_or(0, |overlay| {
165            overlay
166                .values()
167                .flat_map(BTreeMap::values)
168                .filter(|record| record.schema.kind.namespace == namespace)
169                .count()
170        });
171        settled.checked_add(proposed).ok_or_else(|| {
172            CanwuError::new(
173                ErrorCode::ValueOutOfRange,
174                "knowledge namespace record count overflowed",
175            )
176        })
177    }
178
179    /// Queries holder-relative records for an omniscient plugin system.
180    ///
181    /// This enforces the declared `canwu.core.knowledge` read and returns an
182    /// owned projection. It is not an actor-facing authorization API.
183    pub fn knowledge_records(
184        &self,
185        holder: KnowledgeHolderRef,
186        query: &KnowledgeQuery,
187    ) -> Result<canwu_knowledge::KnowledgeQueryResult, CanwuError> {
188        self.require_read(&StateKey::core_knowledge())?;
189        let result = if let Some(overlay) = self.knowledge_overlay {
190            self.state.current().knowledge.query_with_overlay(
191                holder,
192                query,
193                self.boundary_id,
194                overlay,
195            )
196        } else {
197            self.state
198                .current()
199                .knowledge
200                .query_current(holder, query, self.boundary_id)
201        };
202        result.map_err(|error| match error {
203            canwu_knowledge::KnowledgeQueryError::ReadCutUnavailable => CanwuError::new(
204                ErrorCode::KnowledgeReadCutUnavailable,
205                "knowledge cursor read cut is no longer available",
206            ),
207            canwu_knowledge::KnowledgeQueryError::InvalidLimit => CanwuError::new(
208                ErrorCode::KnowledgeLimitExceeded,
209                "knowledge query page size is outside the supported range",
210            ),
211            canwu_knowledge::KnowledgeQueryError::InvalidCursor
212            | canwu_knowledge::KnowledgeQueryError::InvalidLedger
213            | canwu_knowledge::KnowledgeQueryError::Encoding => CanwuError::new(
214                ErrorCode::InvalidKnowledgeRecord,
215                "knowledge query, cursor, or ledger is invalid",
216            ),
217        })
218    }
219
220    /// Resolves an exact command ID from the retained runtime journal in O(1).
221    ///
222    /// An archived command remains valid identity evidence, but its payload is
223    /// no longer available through this view: lookup returns
224    /// [`ErrorCode::EvidenceContentUnavailable`]. `None` means the ID has
225    /// neither retained content nor a committed archive receipt.
226    pub fn command(&self, id: CommandId) -> Result<Option<&CommandRecord>, CanwuError> {
227        self.require_read(&StateKey::core_commands())?;
228        let retained = self.state.evidence().retained_command(id);
229        if retained.is_none()
230            && self
231                .state
232                .evidence()
233                .archived_evidence_receipts
234                .contains_key(&EvidenceRef::Command(id))
235        {
236            return Err(CanwuError::new(
237                ErrorCode::EvidenceContentUnavailable,
238                "command identity is archived; payload inspection requires an archive provider",
239            ));
240        }
241        Ok(retained)
242    }
243
244    /// Resolves an exact event ID from the retained runtime journal in O(1).
245    ///
246    /// An archived event remains valid identity evidence, but its payload is
247    /// no longer available through this view: lookup returns
248    /// [`ErrorCode::EvidenceContentUnavailable`]. `None` means the ID has
249    /// neither retained content nor a committed archive receipt.
250    pub fn event(&self, id: EventId) -> Result<Option<&SimEvent>, CanwuError> {
251        self.require_read(&StateKey::core_events())?;
252        let retained = self.state.evidence().retained_event(id);
253        if retained.is_none()
254            && self
255                .state
256                .evidence()
257                .archived_evidence_receipts
258                .contains_key(&EvidenceRef::Event(id))
259        {
260            return Err(CanwuError::new(
261                ErrorCode::EvidenceContentUnavailable,
262                "event identity is archived; payload inspection requires an archive provider",
263            ));
264        }
265        Ok(retained)
266    }
267
268    pub fn ingress(&self, id: IngressId) -> Result<Option<&IngressRecord>, CanwuError> {
269        self.require_read(&StateKey::core_ingress())?;
270        if self
271            .allowed_ingress
272            .is_none_or(|allowed| !allowed.contains(&id))
273        {
274            return Ok(None);
275        }
276        let record = self.state.evidence().retained_ingress(id);
277        if record.is_none()
278            && self
279                .state
280                .evidence()
281                .archived_evidence_receipts
282                .contains_key(&EvidenceRef::Ingress(id))
283        {
284            return Err(CanwuError::new(
285                ErrorCode::EvidenceContentUnavailable,
286                "ingress identity is archived; payload inspection requires an archive provider",
287            ));
288        }
289        if let (Some(owner), Some(record)) = (self.ingress_plugin, record)
290            && !matches!(
291                &record.payload,
292                IngressPayload::Plugin { plugin, .. } if plugin == owner
293            )
294        {
295            return Ok(None);
296        }
297        Ok(record)
298    }
299
300    /// A retained record that is still queued, or that a terminal
301    /// cancellation withdrew, was never admitted and is not durable evidence
302    /// of a delivered packet.
303    fn ingress_is_pending_or_cancelled(&self, record: &IngressRecord) -> bool {
304        let scheduler = &self.state.runtime().scheduler;
305        scheduler.cancelled_ingress.contains(&record.id)
306            || scheduler
307                .pending_ingress
308                .contains(&IngressQueueKey::from_record(record))
309    }
310
311    /// Lists the still-pending plugin ingress that this boundary system's own
312    /// plugin scheduled inside the engine and may withdraw with
313    /// [`crate::BoundaryDirective::CancelPluginIngress`], in ingress-ID order.
314    ///
315    /// Only items due strictly after the view time are listed. Items that a
316    /// boundary directive scheduled earlier in the current boundary appear
317    /// from the next boundary; items generated by this plugin's commands are
318    /// listed as soon as they are queued. Views that are not bound to a
319    /// boundary system list nothing.
320    pub fn cancellable_plugin_ingress(&self) -> Result<Vec<&IngressRecord>, CanwuError> {
321        self.require_read(&StateKey::core_ingress())?;
322        let Some(owner) = self.ingress_plugin else {
323            return Ok(Vec::new());
324        };
325        let now = self.state.now();
326        let evidence = self.state.evidence();
327        let mut records = Vec::new();
328        for key in &self.state.runtime().scheduler.pending_ingress {
329            if key.due_at <= now {
330                continue;
331            }
332            let Some(record) = evidence.retained_ingress(key.id) else {
333                continue;
334            };
335            let IngressPayload::Plugin { plugin, .. } = &record.payload else {
336                continue;
337            };
338            let issued_by_owner = match &record.cause {
339                Some(CauseRef::Command(_)) => plugin == owner,
340                Some(CauseRef::Boundary(boundary)) => evidence
341                    .retained_boundary(*boundary)
342                    .is_some_and(|boundary| {
343                        boundary.generated_ingress.iter().any(|generation| {
344                            generation.ingress == record.id && generation.plugin == owner
345                        })
346                    }),
347                Some(CauseRef::System(_) | CauseRef::Event(_)) | None => false,
348            };
349            if issued_by_owner {
350                records.push(record);
351            }
352        }
353        records.sort_by_key(|record| record.id);
354        Ok(records)
355    }
356
357    /// Matches retained, plugin-generated ingress provenance without exposing its payload.
358    ///
359    /// Durable evidence may cite ingress admitted at an earlier boundary, but a
360    /// generated record still waiting in the scheduler is not yet admissible.
361    /// Format-7 archive receipts retain a Merkle-bound compact producer proof,
362    /// so archived payload bytes do not need to return to the hot path.
363    pub fn plugin_ingress_matches(
364        &self,
365        id: IngressId,
366        plugin: &str,
367        packet_type: &str,
368    ) -> Result<bool, CanwuError> {
369        self.require_read(&StateKey::core_ingress())?;
370        let record = self.state.evidence().retained_ingress(id);
371        if record.is_none()
372            && let Some(receipt) = self
373                .state
374                .evidence()
375                .archived_evidence_receipts
376                .get(&EvidenceRef::Ingress(id))
377        {
378            if self
379                .state
380                .runtime()
381                .scheduler
382                .pending_ingress
383                .iter()
384                .any(|key| key.id == id)
385            {
386                return Ok(false);
387            }
388            return Ok(receipt
389                .plugin_ingress_provenance
390                .as_ref()
391                .is_some_and(|provenance| {
392                    provenance.plugin == plugin && provenance.packet_type == packet_type
393                }));
394        }
395        let Some(record) = record else {
396            return Ok(false);
397        };
398        if self.ingress_is_pending_or_cancelled(record) {
399            return Ok(false);
400        }
401        if !matches!(
402            &record.payload,
403            IngressPayload::Plugin {
404                plugin: actual_plugin,
405                packet_type: actual_packet_type,
406                ..
407            } if actual_plugin == plugin && actual_packet_type == packet_type
408        ) {
409            return Ok(false);
410        }
411        let Some(CauseRef::Boundary(boundary_id)) = record.cause.as_ref() else {
412            return Ok(false);
413        };
414        let boundary = self.state.evidence().retained_boundary(*boundary_id);
415        if boundary.is_none()
416            && self
417                .state
418                .evidence()
419                .archived_evidence_receipts
420                .contains_key(&EvidenceRef::Boundary(*boundary_id))
421        {
422            return Err(CanwuError::new(
423                ErrorCode::EvidenceContentUnavailable,
424                "ingress producer boundary is archived; provenance inspection requires an archive provider",
425            ));
426        }
427        Ok(boundary.is_some_and(|boundary| {
428            boundary
429                .generated_ingress
430                .iter()
431                .any(|generation| generation.ingress == id && generation.plugin == plugin)
432        }))
433    }
434
435    /// Matches a retained plugin ingress to an exact provider payload and
436    /// delivery time. Payload inspection is deliberately limited to retained
437    /// records: archived receipts prove producer identity, but cannot safely
438    /// be reused to authorize a different legal proposal without the original
439    /// bytes.
440    pub fn plugin_ingress_payload_matches(
441        &self,
442        id: IngressId,
443        plugin: &str,
444        packet_type: &str,
445        occurred_at: SimTime,
446        expected_payload: &Value,
447    ) -> Result<bool, CanwuError> {
448        self.require_read(&StateKey::core_ingress())?;
449        let Some(record) = self.state.evidence().retained_ingress(id) else {
450            if self
451                .state
452                .evidence()
453                .archived_evidence_receipts
454                .contains_key(&EvidenceRef::Ingress(id))
455            {
456                return Err(CanwuError::new(
457                    ErrorCode::EvidenceContentUnavailable,
458                    "provider ingress payload is archived; exact legal signal binding requires retained content",
459                ));
460            }
461            return Ok(false);
462        };
463        if self.ingress_is_pending_or_cancelled(record) {
464            return Ok(false);
465        }
466        let IngressPayload::Plugin {
467            plugin: actual_plugin,
468            packet_type: actual_packet_type,
469            payload,
470            ..
471        } = &record.payload
472        else {
473            return Ok(false);
474        };
475        if actual_plugin != plugin
476            || actual_packet_type != packet_type
477            || record.due_at != occurred_at
478            || payload != expected_payload
479        {
480            return Ok(false);
481        }
482        let Some(CauseRef::Boundary(boundary_id)) = record.cause.as_ref() else {
483            return Ok(false);
484        };
485        let boundary = self.state.evidence().retained_boundary(*boundary_id);
486        if boundary.is_none()
487            && self
488                .state
489                .evidence()
490                .archived_evidence_receipts
491                .contains_key(&EvidenceRef::Boundary(*boundary_id))
492        {
493            return Err(CanwuError::new(
494                ErrorCode::EvidenceContentUnavailable,
495                "provider ingress producer boundary is archived; exact legal signal binding requires retained content",
496            ));
497        }
498        Ok(boundary.is_some_and(|boundary| {
499            boundary
500                .generated_ingress
501                .iter()
502                .any(|generation| generation.ingress == id && generation.plugin == plugin)
503        }))
504    }
505
506    /// Returns the retained outcome for one exact decision request.
507    pub fn decision_attempt(
508        &self,
509        request_id: DecisionRequestId,
510    ) -> Result<Option<&DecisionAttemptRecord>, CanwuError> {
511        self.require_read(&StateKey::core_decisions())?;
512        Ok(self.state.current().decisions.attempt(request_id))
513    }
514
515    /// Returns one current decision-controller binding after an explicit core read.
516    pub fn decision_controller(
517        &self,
518        id: &str,
519    ) -> Result<Option<&DecisionControllerBinding>, CanwuError> {
520        self.require_read(&StateKey::core_decisions())?;
521        Ok(self.state.current().decisions.controller(id))
522    }
523
524    /// Returns one current decision ticket after an explicit core read.
525    pub fn decision_ticket(
526        &self,
527        id: DecisionTicketId,
528    ) -> Result<Option<&DecisionTicket>, CanwuError> {
529        self.require_read(&StateKey::core_decisions())?;
530        Ok(self.state.current().decisions.ticket(id))
531    }
532
533    pub fn domain_record(
534        &self,
535        reference: &DomainRecordRef,
536    ) -> Result<Option<&DomainRecord>, CanwuError> {
537        self.require_domain_record_read(reference)?;
538        Ok(self
539            .record_overlay
540            .and_then(|overlay| overlay.get(reference))
541            .or_else(|| self.state.current().domain_records.get(reference)))
542    }
543
544    pub fn typed_domain_record<T: DomainRecordType>(
545        &self,
546        reference: &TypedDomainRecordRef<T>,
547    ) -> Result<Option<&DomainRecord>, CanwuError> {
548        self.domain_record(reference.as_untyped())
549    }
550
551    pub fn proposed_domain_record(
552        &self,
553        reference: &DomainRecordRef,
554    ) -> Result<Option<&DomainRecord>, CanwuError> {
555        self.require_read(&records::record_state_key(&reference.kind))?;
556        Ok(self
557            .proposed_records
558            .and_then(|records| records.get(reference)))
559    }
560
561    pub fn proposed_typed_domain_record<T: DomainRecordType>(
562        &self,
563        reference: &TypedDomainRecordRef<T>,
564    ) -> Result<Option<&DomainRecord>, CanwuError> {
565        self.proposed_domain_record(reference.as_untyped())
566    }
567
568    /// Returns the exact evidence reference assigned to a domain-record
569    /// version proposed earlier in the current boundary.
570    pub fn proposed_domain_record_version(
571        &self,
572        reference: &DomainRecordRef,
573    ) -> Result<Option<DomainRecordVersionRef>, CanwuError> {
574        self.require_read(&records::record_state_key(&reference.kind))?;
575        Ok(self.proposed_version_evidence(reference))
576    }
577
578    /// Same-boundary proposal evidence for one record, without a read check;
579    /// callers gate access before using it.
580    fn proposed_version_evidence(
581        &self,
582        reference: &DomainRecordRef,
583    ) -> Option<DomainRecordVersionRef> {
584        self.proposal_evidence.and_then(|evidence| {
585            evidence.iter().find_map(|item| match item {
586                EvidenceRef::DomainRecordVersion(version) if version.record == *reference => {
587                    Some(version.clone())
588                }
589                _ => None,
590            })
591        })
592    }
593
594    /// Returns the exact evidence reference for the currently visible version
595    /// of a domain record.  Strategic aggregation runs after atomic commit,
596    /// therefore it cannot use [`Self::proposed_domain_record_version`].
597    ///
598    /// The current boundary overlay/proposal is preferred, followed by the
599    /// runtime's verified current-record provenance index. The index is
600    /// maintained at commit time and rebuilt from canonical evidence on
601    /// restore, so lookup does not scan retained or archived history.
602    pub fn current_domain_record_version(
603        &self,
604        reference: &DomainRecordRef,
605    ) -> Result<Option<DomainRecordVersionRef>, CanwuError> {
606        self.require_read(&records::record_state_key(&reference.kind))?;
607        let Some(record) = self.domain_record(reference)? else {
608            return Ok(None);
609        };
610        if let Some(proposed) = self.proposal_evidence.and_then(|evidence| {
611            evidence.iter().find_map(|item| match item {
612                EvidenceRef::DomainRecordVersion(version)
613                    if version.record == *reference && version.version == record.version =>
614                {
615                    Some(version.clone())
616                }
617                _ => None,
618            })
619        }) {
620            return Ok(Some(proposed));
621        }
622        let current = super::current_domain_record_version(self.state.runtime(), reference)?;
623        if current
624            .as_ref()
625            .is_some_and(|current| current.version != record.version)
626        {
627            return Err(CanwuError::new(
628                ErrorCode::InvalidSnapshot,
629                "visible domain-record version disagrees with the runtime provenance index",
630            ));
631        }
632        Ok(current)
633    }
634
635    /// Returns whether an exact domain-record version reference is valid at
636    /// this proposal-visible cut.
637    ///
638    /// This validates both the record identity/version and its establishment
639    /// source. Earlier same-boundary proposals are considered before retained
640    /// or archived runtime evidence. Either the exact record-kind read or the
641    /// administrative domain-record read grants access.
642    pub fn domain_record_version_evidence_exists(
643        &self,
644        reference: &DomainRecordVersionRef,
645    ) -> Result<bool, CanwuError> {
646        // The proposal lookup must not re-require the exact kind, or the
647        // administrative read could never resolve evidence.
648        self.require_domain_record_read(&reference.record)?;
649        if self
650            .proposed_version_evidence(&reference.record)
651            .is_some_and(|proposed| proposed == *reference)
652        {
653            return Ok(true);
654        }
655        Ok(!matches!(
656            validation::resolve_evidence_reference(
657                &validation::RuntimeValidationContext::new(self.state.runtime()),
658                &EvidenceRef::DomainRecordVersion(reference.clone()),
659            ),
660            validation::EvidenceAvailability::Missing
661        ))
662    }
663
664    /// Checks that an exact domain-record version is both valid evidence and current.
665    pub fn domain_record_version_is_current(
666        &self,
667        reference: &DomainRecordVersionRef,
668    ) -> Result<bool, CanwuError> {
669        self.require_domain_record_read(&reference.record)?;
670        let current = self
671            .record_overlay
672            .and_then(|overlay| overlay.get(&reference.record))
673            .or_else(|| self.state.current().domain_records.get(&reference.record));
674        Ok(
675            current.is_some_and(|record| record.version == reference.version)
676                && self.domain_record_version_evidence_exists(reference)?,
677        )
678    }
679
680    /// Returns whether a generic evidence identity is retained or archived.
681    ///
682    /// Domain-record versions proposed earlier in this boundary are visible.
683    /// Archived identities count as existing even when their bodies are no
684    /// longer retained.
685    pub fn evidence_exists(&self, reference: &EvidenceRef) -> Result<bool, CanwuError> {
686        match reference {
687            EvidenceRef::Command(_) | EvidenceRef::CommandAttempt(_) => {
688                self.require_read(&StateKey::core_commands())?;
689            }
690            EvidenceRef::Event(_) => self.require_read(&StateKey::core_events())?,
691            EvidenceRef::Ingress(_) => self.require_read(&StateKey::core_ingress())?,
692            EvidenceRef::Boundary(_) | EvidenceRef::RandomDraw(_) => {
693                self.require_read(&StateKey::core_evidence())?;
694            }
695            EvidenceRef::DomainRecordVersion(version) => {
696                return self.domain_record_version_evidence_exists(version);
697            }
698        }
699        Ok(!matches!(
700            validation::resolve_evidence_reference(
701                &validation::RuntimeValidationContext::new(self.state.runtime()),
702                reference,
703            ),
704            validation::EvidenceAvailability::Missing
705        ))
706    }
707
708    /// Returns when retained or earlier same-boundary evidence first became
709    /// authoritative at this proposal-visible cut.
710    ///
711    /// Archived identity receipts do not retain a precise semantic time, so
712    /// they return `None` and callers that require temporal ordering must fail
713    /// closed or load the archived evidence body.
714    pub fn evidence_time(&self, reference: &EvidenceRef) -> Result<Option<SimTime>, CanwuError> {
715        if !self.evidence_exists(reference)? {
716            return Ok(None);
717        }
718        if self
719            .proposal_evidence
720            .is_some_and(|evidence| evidence.contains(reference))
721        {
722            return Ok(Some(self.time()));
723        }
724        Ok(super::retained_evidence_time(
725            self.state.runtime(),
726            reference,
727        ))
728    }
729
730    /// Resolves the retained record body for one exact domain-record version.
731    ///
732    /// Archived receipts prove that a version existed but do not contain its
733    /// body, so this returns `None` when the corresponding evidence segment is
734    /// not live in the runtime.
735    pub fn domain_record_version(
736        &self,
737        reference: &DomainRecordVersionRef,
738    ) -> Result<Option<DomainRecord>, CanwuError> {
739        self.require_read(&records::record_state_key(&reference.record.kind))?;
740        if let Some(proposed) = self.proposed_domain_record_version(&reference.record)?
741            && proposed == *reference
742        {
743            return Ok(self
744                .proposed_records
745                .and_then(|records| records.get(&reference.record))
746                .or_else(|| {
747                    self.record_overlay
748                        .and_then(|records| records.get(&reference.record))
749                })
750                .cloned());
751        }
752        Ok(retained_domain_record_version(
753            self.state.runtime(),
754            reference,
755        ))
756    }
757
758    /// Returns a bounded, deterministic projection of records of one kind.
759    ///
760    /// Same-boundary overlays take precedence over current state. Records are
761    /// ordered by their canonical reference, so result order is replay stable.
762    pub fn domain_records_of_kind(
763        &self,
764        kind: &DomainRecordKind,
765        limit: usize,
766    ) -> Result<Vec<DomainRecord>, CanwuError> {
767        self.domain_records_of_kind_after(kind, None, limit)
768    }
769
770    /// Returns one bounded deterministic page of records after a canonical
771    /// record-reference cursor.
772    ///
773    /// The cursor is exclusive and must name the same kind. This keeps plugin
774    /// scans bounded without imposing a 10,000-record lifetime ceiling on a
775    /// domain kind. Same-boundary overlays retain the same precedence as
776    /// [`Self::domain_records_of_kind`].
777    pub fn domain_records_of_kind_after(
778        &self,
779        kind: &DomainRecordKind,
780        after: Option<&DomainRecordRef>,
781        limit: usize,
782    ) -> Result<Vec<DomainRecord>, CanwuError> {
783        self.require_read(&records::record_state_key(kind))?;
784        validate_domain_record_page_request(kind, after, limit)?;
785
786        let mut records =
787            domain_record_candidates(&self.state.current().domain_records, kind, after, limit);
788        for overlay in [self.record_overlay, self.proposed_records]
789            .into_iter()
790            .flatten()
791        {
792            for (reference, record) in domain_record_candidates(overlay, kind, after, limit) {
793                records.insert(reference, record);
794            }
795        }
796        Ok(records.into_values().take(limit).collect())
797    }
798
799    /// Finds committed knowledge changes produced with an exact correlation.
800    ///
801    /// This supports next-boundary operation finalization without granting a
802    /// plugin unrestricted access to unrelated knowledge payloads.
803    pub fn knowledge_changes_by_correlation(
804        &self,
805        plugin: &str,
806        producer_correlation: &str,
807    ) -> Result<Vec<BoundaryKnowledgeChange>, CanwuError> {
808        self.require_read(&StateKey::core_knowledge())?;
809        Ok(self
810            .state
811            .evidence()
812            .boundaries
813            .iter()
814            .flat_map(|boundary| &boundary.knowledge_changes)
815            .filter(|change| {
816                change.plugin == plugin
817                    && change.producer_correlation.as_deref() == Some(producer_correlation)
818            })
819            .cloned()
820            .collect())
821    }
822
823    /// Finds committed knowledge changes whose producer correlation begins
824    /// with a deterministic operation prefix.
825    ///
826    /// The prefix remains plugin-scoped. This is intended for bounded
827    /// multi-holder operation finalization where every holder batch must keep
828    /// a unique full correlation value.
829    pub fn knowledge_changes_by_correlation_prefix(
830        &self,
831        plugin: &str,
832        producer_correlation_prefix: &str,
833    ) -> Result<Vec<BoundaryKnowledgeChange>, CanwuError> {
834        self.require_read(&StateKey::core_knowledge())?;
835        Ok(self
836            .state
837            .evidence()
838            .boundaries
839            .iter()
840            .flat_map(|boundary| &boundary.knowledge_changes)
841            .filter(|change| {
842                change.plugin == plugin
843                    && change
844                        .producer_correlation
845                        .as_deref()
846                        .is_some_and(|value| value.starts_with(producer_correlation_prefix))
847            })
848            .cloned()
849            .collect())
850    }
851
852    pub fn reservation(
853        &self,
854        reservation: &ReservationRef,
855    ) -> Result<Option<&ReservationAllocation>, CanwuError> {
856        let reader = self.reader.unwrap_or("unscoped caller");
857        if self
858            .allowed_reservations
859            .is_none_or(|allowed| !allowed.contains(reservation))
860        {
861            return Err(CanwuError::new(
862                ErrorCode::UndeclaredStateRead,
863                format!(
864                    "system {reader} did not declare reservation read {}.{}.{}",
865                    reservation.plugin, reservation.system, reservation.request
866                ),
867            ));
868        }
869        Ok(self.allocations.and_then(|values| values.get(reservation)))
870    }
871
872    pub fn random_range(
873        &self,
874        stream: &RandomStreamKey,
875        upper_exclusive: u64,
876        purpose: &str,
877    ) -> Result<u64, CanwuError> {
878        let Some(session) = &self.random_session else {
879            return Err(CanwuError::new(
880                ErrorCode::UndeclaredRandomStream,
881                format!(
882                    "system {} has no declared random streams",
883                    self.reader.unwrap_or("unscoped caller")
884                ),
885            ));
886        };
887        session.borrow_mut().range(stream, upper_exclusive, purpose)
888    }
889
890    #[allow(clippy::too_many_arguments)]
891    pub fn random_range_for_operation(
892        &self,
893        stream: &RandomStreamKey,
894        evidence: EvidenceRef,
895        operation_kind: &str,
896        application_operation_id: &str,
897        target: RandomOperationTarget,
898        draw_slot: u32,
899        upper_exclusive: u64,
900        purpose: &str,
901    ) -> Result<u64, CanwuError> {
902        self.random_sample_for_operation(
903            stream,
904            evidence,
905            operation_kind,
906            application_operation_id,
907            target,
908            draw_slot,
909            upper_exclusive,
910            purpose,
911        )
912        .map(|sample| sample.value)
913    }
914
915    #[allow(clippy::too_many_arguments)]
916    pub fn random_sample_for_operation(
917        &self,
918        stream: &RandomStreamKey,
919        evidence: EvidenceRef,
920        operation_kind: &str,
921        application_operation_id: &str,
922        target: RandomOperationTarget,
923        draw_slot: u32,
924        upper_exclusive: u64,
925        purpose: &str,
926    ) -> Result<super::RandomSample, CanwuError> {
927        let available = self
928            .proposal_evidence
929            .is_some_and(|values| values.contains(&evidence))
930            || validation::resolve_evidence_reference(
931                &validation::RuntimeValidationContext::new(self.state.runtime()),
932                &evidence,
933            ) == validation::EvidenceAvailability::Retained;
934        if !available {
935            return Err(CanwuError::new(
936                ErrorCode::InvalidRandomOperationEvidence,
937                "operation-keyed random draw references unavailable evidence",
938            ));
939        }
940        let Some(session) = &self.random_session else {
941            return Err(CanwuError::new(
942                ErrorCode::UndeclaredRandomStream,
943                format!(
944                    "system {} has no declared random streams",
945                    self.reader.unwrap_or("unscoped caller")
946                ),
947            ));
948        };
949        session.borrow_mut().sample_for_operation(
950            stream,
951            evidence,
952            operation_kind,
953            application_operation_id,
954            target,
955            draw_slot,
956            upper_exclusive,
957            purpose,
958        )
959    }
960
961    pub fn component(
962        &self,
963        state: &StateKey,
964        entity: &EntityRef,
965        component: &str,
966    ) -> Result<Option<&Value>, CanwuError> {
967        self.require_read(state)?;
968        let Some(owner) = self.state_owners.get(state) else {
969            return Err(CanwuError::new(
970                ErrorCode::UndeclaredStateRead,
971                format!(
972                    "state {}.{} has no registered owner",
973                    state.namespace, state.name
974                ),
975            ));
976        };
977        let key = component_key(owner, state, entity, component);
978        Ok(self
979            .component_overlay
980            .and_then(|overlay| overlay.get(&key))
981            .or_else(|| self.state.current().plugin_components.get(&key))
982            .map(|record| &record.value))
983    }
984
985    pub fn proposed_component(
986        &self,
987        state: &StateKey,
988        entity: &EntityRef,
989        component: &str,
990    ) -> Result<Option<&Value>, CanwuError> {
991        self.require_read(state)?;
992        let Some(owner) = self.state_owners.get(state) else {
993            return Err(CanwuError::new(
994                ErrorCode::UndeclaredStateRead,
995                format!(
996                    "state {}.{} has no registered owner",
997                    state.namespace, state.name
998                ),
999            ));
1000        };
1001        let key = component_key(owner, state, entity, component);
1002        Ok(self
1003            .proposed_components
1004            .and_then(|proposals| proposals.get(&key))
1005            .map(|record| &record.value))
1006    }
1007
1008    fn require_read(&self, state: &StateKey) -> Result<(), CanwuError> {
1009        if self
1010            .allowed_reads
1011            .is_some_and(|reads| !reads.contains(state))
1012        {
1013            return Err(CanwuError::new(
1014                ErrorCode::UndeclaredStateRead,
1015                format!(
1016                    "{} did not declare read access to {}.{}",
1017                    self.reader.unwrap_or("internal system"),
1018                    state.namespace,
1019                    state.name
1020                ),
1021            ));
1022        }
1023        Ok(())
1024    }
1025
1026    fn require_domain_record_read(&self, reference: &DomainRecordRef) -> Result<(), CanwuError> {
1027        let exact = records::record_state_key(&reference.kind);
1028        if self.allowed_reads.is_some_and(|reads| {
1029            !reads.contains(&exact) && !reads.contains(&StateKey::core_domain_records())
1030        }) {
1031            return Err(CanwuError::new(
1032                ErrorCode::UndeclaredStateRead,
1033                format!(
1034                    "{} did not declare read access to {}.{}",
1035                    self.reader.unwrap_or("internal system"),
1036                    exact.namespace,
1037                    exact.name
1038                ),
1039            ));
1040        }
1041        Ok(())
1042    }
1043
1044    pub(super) fn finish_random_session(self) -> Option<random::RandomExecution> {
1045        self.random_session
1046            .map(RefCell::into_inner)
1047            .map(random::RandomSession::finish)
1048    }
1049}
1050
1051impl super::PluginArchiveObjectProvider for SimulationView<'_> {
1052    fn load_plugin_archive_object(
1053        &self,
1054        namespace: &str,
1055        object_id: &str,
1056    ) -> Result<Option<Vec<u8>>, CanwuError> {
1057        self.plugin_archive_object(namespace, object_id)
1058    }
1059}