Skip to main content

canwu_sim/runtime/
view.rs

1use super::{
2    ActorKnowledge, Army, ArmyId, BoundaryId, BoundaryKnowledgeChange, CanwuError, CommandId,
3    CommandRecord, DomainRecord, DomainRecordKind, DomainRecordRef, DomainRecordType,
4    DomainRecordVersionRef, EntityRef, ErrorCode, EventId, EvidenceRef, Government, GovernmentId,
5    HashSet, IngressId, IngressPayload, IngressRecord, KnowledgeHolderRef, KnowledgeQuery,
6    KnowledgeRecord, KnowledgeRecordId, Person, PersonId, PluginComponentKey,
7    PluginComponentRecord, RandomOperationTarget, RandomStreamKey, RefCell, ReservationAllocation,
8    ReservationRef, Route, RouteId, RuntimeCurrentState, RuntimeEvidence, RuntimeState, SimEvent,
9    SimTime, StateKey, Territory, TerritoryId, TypedDomainRecordRef, Value, component_key,
10    domain_record_candidates, random, records, retained_domain_record_version,
11    validate_domain_record_page_request, validation,
12};
13use std::collections::{BTreeMap, BTreeSet};
14
15pub(super) enum SimulationViewState<'a> {
16    Runtime(&'a RuntimeState),
17    Boundary {
18        current: &'a RuntimeCurrentState,
19        now: SimTime,
20        runtime: &'a RuntimeState,
21    },
22}
23
24impl SimulationViewState<'_> {
25    const fn current(&self) -> &RuntimeCurrentState {
26        match self {
27            Self::Runtime(state) => &state.current,
28            Self::Boundary { current, .. } => current,
29        }
30    }
31
32    const fn now(&self) -> SimTime {
33        match self {
34            Self::Runtime(state) => state.scheduler.now,
35            Self::Boundary { now, .. } => *now,
36        }
37    }
38
39    const fn evidence(&self) -> &RuntimeEvidence {
40        match self {
41            Self::Runtime(state) => &state.evidence,
42            Self::Boundary { runtime, .. } => &runtime.evidence,
43        }
44    }
45
46    const fn runtime(&self) -> &RuntimeState {
47        match self {
48            Self::Runtime(state) | Self::Boundary { runtime: state, .. } => state,
49        }
50    }
51}
52
53pub struct SimulationView<'a> {
54    pub(super) state: SimulationViewState<'a>,
55    pub(super) state_owners: &'a BTreeMap<StateKey, String>,
56    pub(super) reader: Option<&'a str>,
57    pub(super) allowed_reads: Option<&'a [StateKey]>,
58    pub(super) allowed_ingress: Option<&'a HashSet<IngressId>>,
59    pub(super) ingress_plugin: Option<&'a str>,
60    pub(super) component_overlay: Option<&'a BTreeMap<PluginComponentKey, PluginComponentRecord>>,
61    pub(super) proposed_components: Option<&'a BTreeMap<PluginComponentKey, PluginComponentRecord>>,
62    pub(super) record_overlay: Option<&'a BTreeMap<DomainRecordRef, DomainRecord>>,
63    pub(super) proposed_records: Option<&'a BTreeMap<DomainRecordRef, DomainRecord>>,
64    pub(super) boundary_id: Option<BoundaryId>,
65    pub(super) proposal_evidence: Option<&'a BTreeSet<EvidenceRef>>,
66    pub(super) knowledge_overlay:
67        Option<&'a BTreeMap<KnowledgeHolderRef, BTreeMap<KnowledgeRecordId, KnowledgeRecord>>>,
68    pub(super) allocations: Option<&'a BTreeMap<ReservationRef, ReservationAllocation>>,
69    pub(super) allowed_reservations: Option<&'a [ReservationRef]>,
70    pub(super) random_session: Option<RefCell<random::RandomSession>>,
71}
72
73impl SimulationView<'_> {
74    #[must_use]
75    pub const fn time(&self) -> SimTime {
76        self.state.now()
77    }
78
79    pub fn army(&self, id: ArmyId) -> Result<Option<&Army>, CanwuError> {
80        self.require_read(&StateKey::core_armies())?;
81        Ok(self.state.current().armies.get(&id))
82    }
83
84    pub fn person(&self, id: PersonId) -> Result<Option<&Person>, CanwuError> {
85        self.require_read(&StateKey::core_people())?;
86        Ok(self.state.current().people.get(&id))
87    }
88
89    pub fn government(&self, id: GovernmentId) -> Result<Option<&Government>, CanwuError> {
90        self.require_read(&StateKey::core_governments())?;
91        Ok(self.state.current().governments.get(&id))
92    }
93
94    pub fn territory(&self, id: TerritoryId) -> Result<Option<&Territory>, CanwuError> {
95        self.require_read(&StateKey::core_territories())?;
96        Ok(self.state.current().territories.get(&id))
97    }
98
99    pub fn route(&self, id: RouteId) -> Result<Option<&Route>, CanwuError> {
100        self.require_read(&StateKey::core_routes())?;
101        Ok(self.state.current().routes.get(&id))
102    }
103
104    pub fn actor_knowledge(&self, actor: PersonId) -> Result<Option<&ActorKnowledge>, CanwuError> {
105        self.require_read(&StateKey::core_knowledge())?;
106        Ok(self.state.current().knowledge.for_actor(actor))
107    }
108
109    /// Counts records in a knowledge namespace at the current proposal-visible cut.
110    pub fn knowledge_record_count_in_namespace(
111        &self,
112        namespace: &str,
113    ) -> Result<usize, CanwuError> {
114        self.require_read(&StateKey::core_knowledge())?;
115        let settled = self
116            .state
117            .current()
118            .knowledge
119            .record_count_in_namespace(namespace);
120        let proposed = self.knowledge_overlay.map_or(0, |overlay| {
121            overlay
122                .values()
123                .flat_map(BTreeMap::values)
124                .filter(|record| record.schema.kind.namespace == namespace)
125                .count()
126        });
127        settled.checked_add(proposed).ok_or_else(|| {
128            CanwuError::new(
129                ErrorCode::ValueOutOfRange,
130                "knowledge namespace record count overflowed",
131            )
132        })
133    }
134
135    /// Queries holder-relative records for an omniscient plugin system.
136    ///
137    /// This enforces the declared `canwu.core.knowledge` read and returns an
138    /// owned projection. It is not an actor-facing authorization API.
139    pub fn knowledge_records(
140        &self,
141        holder: KnowledgeHolderRef,
142        query: &KnowledgeQuery,
143    ) -> Result<canwu_knowledge::KnowledgeQueryResult, CanwuError> {
144        self.require_read(&StateKey::core_knowledge())?;
145        let result = if let Some(overlay) = self.knowledge_overlay {
146            self.state.current().knowledge.query_with_overlay(
147                holder,
148                query,
149                self.boundary_id,
150                overlay,
151            )
152        } else {
153            self.state
154                .current()
155                .knowledge
156                .query_current(holder, query, self.boundary_id)
157        };
158        result.map_err(|error| match error {
159            canwu_knowledge::KnowledgeQueryError::ReadCutUnavailable => CanwuError::new(
160                ErrorCode::KnowledgeReadCutUnavailable,
161                "knowledge cursor read cut is no longer available",
162            ),
163            canwu_knowledge::KnowledgeQueryError::InvalidLimit => CanwuError::new(
164                ErrorCode::KnowledgeLimitExceeded,
165                "knowledge query page size is outside the supported range",
166            ),
167            canwu_knowledge::KnowledgeQueryError::InvalidCursor
168            | canwu_knowledge::KnowledgeQueryError::InvalidLedger
169            | canwu_knowledge::KnowledgeQueryError::Encoding => CanwuError::new(
170                ErrorCode::InvalidKnowledgeRecord,
171                "knowledge query, cursor, or ledger is invalid",
172            ),
173        })
174    }
175
176    /// Resolves an exact command ID from the retained runtime journal in O(1).
177    ///
178    /// An archived command remains valid identity evidence, but its payload is
179    /// no longer available through this view: lookup returns
180    /// [`ErrorCode::EvidenceContentUnavailable`]. `None` means the ID has
181    /// neither retained content nor a committed archive receipt.
182    pub fn command(&self, id: CommandId) -> Result<Option<&CommandRecord>, CanwuError> {
183        self.require_read(&StateKey::core_commands())?;
184        let retained = self.state.evidence().retained_command(id);
185        if retained.is_none()
186            && self
187                .state
188                .evidence()
189                .archived_evidence_receipts
190                .contains_key(&EvidenceRef::Command(id))
191        {
192            return Err(CanwuError::new(
193                ErrorCode::EvidenceContentUnavailable,
194                "command identity is archived; payload inspection requires an archive provider",
195            ));
196        }
197        Ok(retained)
198    }
199
200    /// Resolves an exact event ID from the retained runtime journal in O(1).
201    ///
202    /// An archived event remains valid identity evidence, but its payload is
203    /// no longer available through this view: lookup returns
204    /// [`ErrorCode::EvidenceContentUnavailable`]. `None` means the ID has
205    /// neither retained content nor a committed archive receipt.
206    pub fn event(&self, id: EventId) -> Result<Option<&SimEvent>, CanwuError> {
207        self.require_read(&StateKey::core_events())?;
208        let retained = self.state.evidence().retained_event(id);
209        if retained.is_none()
210            && self
211                .state
212                .evidence()
213                .archived_evidence_receipts
214                .contains_key(&EvidenceRef::Event(id))
215        {
216            return Err(CanwuError::new(
217                ErrorCode::EvidenceContentUnavailable,
218                "event identity is archived; payload inspection requires an archive provider",
219            ));
220        }
221        Ok(retained)
222    }
223
224    pub fn ingress(&self, id: IngressId) -> Result<Option<&IngressRecord>, CanwuError> {
225        self.require_read(&StateKey::core_ingress())?;
226        if self
227            .allowed_ingress
228            .is_none_or(|allowed| !allowed.contains(&id))
229        {
230            return Ok(None);
231        }
232        let record = self.state.evidence().retained_ingress(id);
233        if record.is_none()
234            && self
235                .state
236                .evidence()
237                .archived_evidence_receipts
238                .contains_key(&EvidenceRef::Ingress(id))
239        {
240            return Err(CanwuError::new(
241                ErrorCode::EvidenceContentUnavailable,
242                "ingress identity is archived; payload inspection requires an archive provider",
243            ));
244        }
245        if let (Some(owner), Some(record)) = (self.ingress_plugin, record)
246            && !matches!(
247                &record.payload,
248                IngressPayload::Plugin { plugin, .. } if plugin == owner
249            )
250        {
251            return Ok(None);
252        }
253        Ok(record)
254    }
255
256    pub fn domain_record(
257        &self,
258        reference: &DomainRecordRef,
259    ) -> Result<Option<&DomainRecord>, CanwuError> {
260        self.require_read(&records::record_state_key(&reference.kind))?;
261        Ok(self
262            .record_overlay
263            .and_then(|overlay| overlay.get(reference))
264            .or_else(|| self.state.current().domain_records.get(reference)))
265    }
266
267    pub fn typed_domain_record<T: DomainRecordType>(
268        &self,
269        reference: &TypedDomainRecordRef<T>,
270    ) -> Result<Option<&DomainRecord>, CanwuError> {
271        self.domain_record(reference.as_untyped())
272    }
273
274    pub fn proposed_domain_record(
275        &self,
276        reference: &DomainRecordRef,
277    ) -> Result<Option<&DomainRecord>, CanwuError> {
278        self.require_read(&records::record_state_key(&reference.kind))?;
279        Ok(self
280            .proposed_records
281            .and_then(|records| records.get(reference)))
282    }
283
284    pub fn proposed_typed_domain_record<T: DomainRecordType>(
285        &self,
286        reference: &TypedDomainRecordRef<T>,
287    ) -> Result<Option<&DomainRecord>, CanwuError> {
288        self.proposed_domain_record(reference.as_untyped())
289    }
290
291    /// Returns the exact evidence reference assigned to a domain-record
292    /// version proposed earlier in the current boundary.
293    pub fn proposed_domain_record_version(
294        &self,
295        reference: &DomainRecordRef,
296    ) -> Result<Option<DomainRecordVersionRef>, CanwuError> {
297        self.require_read(&records::record_state_key(&reference.kind))?;
298        Ok(self.proposal_evidence.and_then(|evidence| {
299            evidence.iter().find_map(|item| match item {
300                EvidenceRef::DomainRecordVersion(version) if version.record == *reference => {
301                    Some(version.clone())
302                }
303                _ => None,
304            })
305        }))
306    }
307
308    /// Returns whether an exact domain-record version reference is valid at
309    /// this proposal-visible cut.
310    ///
311    /// This validates both the record identity/version and its establishment
312    /// source. Earlier same-boundary proposals are considered before retained
313    /// or archived runtime evidence.
314    pub fn domain_record_version_evidence_exists(
315        &self,
316        reference: &DomainRecordVersionRef,
317    ) -> Result<bool, CanwuError> {
318        self.require_read(&records::record_state_key(&reference.record.kind))?;
319        if self
320            .proposed_domain_record_version(&reference.record)?
321            .is_some_and(|proposed| proposed == *reference)
322        {
323            return Ok(true);
324        }
325        Ok(!matches!(
326            validation::resolve_evidence_reference(
327                &validation::RuntimeValidationContext::new(self.state.runtime()),
328                &EvidenceRef::DomainRecordVersion(reference.clone()),
329            ),
330            validation::EvidenceAvailability::Missing
331        ))
332    }
333
334    /// Returns whether a generic evidence identity is retained or archived.
335    ///
336    /// Domain-record versions proposed earlier in this boundary are visible.
337    /// Archived identities count as existing even when their bodies are no
338    /// longer retained.
339    pub fn evidence_exists(&self, reference: &EvidenceRef) -> Result<bool, CanwuError> {
340        match reference {
341            EvidenceRef::Command(_) | EvidenceRef::CommandAttempt(_) => {
342                self.require_read(&StateKey::core_commands())?;
343            }
344            EvidenceRef::Event(_) => self.require_read(&StateKey::core_events())?,
345            EvidenceRef::Ingress(_) => self.require_read(&StateKey::core_ingress())?,
346            EvidenceRef::Boundary(_) | EvidenceRef::RandomDraw(_) => {
347                self.require_read(&StateKey::core_evidence())?;
348            }
349            EvidenceRef::DomainRecordVersion(version) => {
350                return self.domain_record_version_evidence_exists(version);
351            }
352        }
353        Ok(!matches!(
354            validation::resolve_evidence_reference(
355                &validation::RuntimeValidationContext::new(self.state.runtime()),
356                reference,
357            ),
358            validation::EvidenceAvailability::Missing
359        ))
360    }
361
362    /// Returns when retained or earlier same-boundary evidence first became
363    /// authoritative at this proposal-visible cut.
364    ///
365    /// Archived identity receipts do not retain a precise semantic time, so
366    /// they return `None` and callers that require temporal ordering must fail
367    /// closed or load the archived evidence body.
368    pub fn evidence_time(&self, reference: &EvidenceRef) -> Result<Option<SimTime>, CanwuError> {
369        if !self.evidence_exists(reference)? {
370            return Ok(None);
371        }
372        if self
373            .proposal_evidence
374            .is_some_and(|evidence| evidence.contains(reference))
375        {
376            return Ok(Some(self.time()));
377        }
378        Ok(super::retained_evidence_time(
379            self.state.runtime(),
380            reference,
381        ))
382    }
383
384    /// Resolves the retained record body for one exact domain-record version.
385    ///
386    /// Archived receipts prove that a version existed but do not contain its
387    /// body, so this returns `None` when the corresponding evidence segment is
388    /// not live in the runtime.
389    pub fn domain_record_version(
390        &self,
391        reference: &DomainRecordVersionRef,
392    ) -> Result<Option<DomainRecord>, CanwuError> {
393        self.require_read(&records::record_state_key(&reference.record.kind))?;
394        if let Some(proposed) = self.proposed_domain_record_version(&reference.record)?
395            && proposed == *reference
396        {
397            return Ok(self
398                .proposed_records
399                .and_then(|records| records.get(&reference.record))
400                .or_else(|| {
401                    self.record_overlay
402                        .and_then(|records| records.get(&reference.record))
403                })
404                .cloned());
405        }
406        Ok(retained_domain_record_version(
407            self.state.runtime(),
408            reference,
409        ))
410    }
411
412    /// Returns a bounded, deterministic projection of records of one kind.
413    ///
414    /// Same-boundary overlays take precedence over current state. Records are
415    /// ordered by their canonical reference, so result order is replay stable.
416    pub fn domain_records_of_kind(
417        &self,
418        kind: &DomainRecordKind,
419        limit: usize,
420    ) -> Result<Vec<DomainRecord>, CanwuError> {
421        self.domain_records_of_kind_after(kind, None, limit)
422    }
423
424    /// Returns one bounded deterministic page of records after a canonical
425    /// record-reference cursor.
426    ///
427    /// The cursor is exclusive and must name the same kind. This keeps plugin
428    /// scans bounded without imposing a 10,000-record lifetime ceiling on a
429    /// domain kind. Same-boundary overlays retain the same precedence as
430    /// [`Self::domain_records_of_kind`].
431    pub fn domain_records_of_kind_after(
432        &self,
433        kind: &DomainRecordKind,
434        after: Option<&DomainRecordRef>,
435        limit: usize,
436    ) -> Result<Vec<DomainRecord>, CanwuError> {
437        self.require_read(&records::record_state_key(kind))?;
438        validate_domain_record_page_request(kind, after, limit)?;
439
440        let mut records =
441            domain_record_candidates(&self.state.current().domain_records, kind, after, limit);
442        for overlay in [self.record_overlay, self.proposed_records]
443            .into_iter()
444            .flatten()
445        {
446            for (reference, record) in domain_record_candidates(overlay, kind, after, limit) {
447                records.insert(reference, record);
448            }
449        }
450        Ok(records.into_values().take(limit).collect())
451    }
452
453    /// Finds committed knowledge changes produced with an exact correlation.
454    ///
455    /// This supports next-boundary operation finalization without granting a
456    /// plugin unrestricted access to unrelated knowledge payloads.
457    pub fn knowledge_changes_by_correlation(
458        &self,
459        plugin: &str,
460        producer_correlation: &str,
461    ) -> Result<Vec<BoundaryKnowledgeChange>, CanwuError> {
462        self.require_read(&StateKey::core_knowledge())?;
463        Ok(self
464            .state
465            .evidence()
466            .boundaries
467            .iter()
468            .flat_map(|boundary| &boundary.knowledge_changes)
469            .filter(|change| {
470                change.plugin == plugin
471                    && change.producer_correlation.as_deref() == Some(producer_correlation)
472            })
473            .cloned()
474            .collect())
475    }
476
477    /// Finds committed knowledge changes whose producer correlation begins
478    /// with a deterministic operation prefix.
479    ///
480    /// The prefix remains plugin-scoped. This is intended for bounded
481    /// multi-holder operation finalization where every holder batch must keep
482    /// a unique full correlation value.
483    pub fn knowledge_changes_by_correlation_prefix(
484        &self,
485        plugin: &str,
486        producer_correlation_prefix: &str,
487    ) -> Result<Vec<BoundaryKnowledgeChange>, CanwuError> {
488        self.require_read(&StateKey::core_knowledge())?;
489        Ok(self
490            .state
491            .evidence()
492            .boundaries
493            .iter()
494            .flat_map(|boundary| &boundary.knowledge_changes)
495            .filter(|change| {
496                change.plugin == plugin
497                    && change
498                        .producer_correlation
499                        .as_deref()
500                        .is_some_and(|value| value.starts_with(producer_correlation_prefix))
501            })
502            .cloned()
503            .collect())
504    }
505
506    pub fn reservation(
507        &self,
508        reservation: &ReservationRef,
509    ) -> Result<Option<&ReservationAllocation>, CanwuError> {
510        let reader = self.reader.unwrap_or("unscoped caller");
511        if self
512            .allowed_reservations
513            .is_none_or(|allowed| !allowed.contains(reservation))
514        {
515            return Err(CanwuError::new(
516                ErrorCode::UndeclaredStateRead,
517                format!(
518                    "system {reader} did not declare reservation read {}.{}.{}",
519                    reservation.plugin, reservation.system, reservation.request
520                ),
521            ));
522        }
523        Ok(self.allocations.and_then(|values| values.get(reservation)))
524    }
525
526    pub fn random_range(
527        &self,
528        stream: &RandomStreamKey,
529        upper_exclusive: u64,
530        purpose: &str,
531    ) -> Result<u64, CanwuError> {
532        let Some(session) = &self.random_session else {
533            return Err(CanwuError::new(
534                ErrorCode::UndeclaredRandomStream,
535                format!(
536                    "system {} has no declared random streams",
537                    self.reader.unwrap_or("unscoped caller")
538                ),
539            ));
540        };
541        session.borrow_mut().range(stream, upper_exclusive, purpose)
542    }
543
544    #[allow(clippy::too_many_arguments)]
545    pub fn random_range_for_operation(
546        &self,
547        stream: &RandomStreamKey,
548        evidence: EvidenceRef,
549        operation_kind: &str,
550        application_operation_id: &str,
551        target: RandomOperationTarget,
552        draw_slot: u32,
553        upper_exclusive: u64,
554        purpose: &str,
555    ) -> Result<u64, CanwuError> {
556        let available = self
557            .proposal_evidence
558            .is_some_and(|values| values.contains(&evidence))
559            || validation::resolve_evidence_reference(
560                &validation::RuntimeValidationContext::new(self.state.runtime()),
561                &evidence,
562            ) == validation::EvidenceAvailability::Retained;
563        if !available {
564            return Err(CanwuError::new(
565                ErrorCode::InvalidRandomOperationEvidence,
566                "operation-keyed random draw references unavailable evidence",
567            ));
568        }
569        let Some(session) = &self.random_session else {
570            return Err(CanwuError::new(
571                ErrorCode::UndeclaredRandomStream,
572                format!(
573                    "system {} has no declared random streams",
574                    self.reader.unwrap_or("unscoped caller")
575                ),
576            ));
577        };
578        session.borrow_mut().range_for_operation(
579            stream,
580            evidence,
581            operation_kind,
582            application_operation_id,
583            target,
584            draw_slot,
585            upper_exclusive,
586            purpose,
587        )
588    }
589
590    pub fn component(
591        &self,
592        state: &StateKey,
593        entity: &EntityRef,
594        component: &str,
595    ) -> Result<Option<&Value>, CanwuError> {
596        self.require_read(state)?;
597        let Some(owner) = self.state_owners.get(state) else {
598            return Err(CanwuError::new(
599                ErrorCode::UndeclaredStateRead,
600                format!(
601                    "state {}.{} has no registered owner",
602                    state.namespace, state.name
603                ),
604            ));
605        };
606        let key = component_key(owner, state, entity, component);
607        Ok(self
608            .component_overlay
609            .and_then(|overlay| overlay.get(&key))
610            .or_else(|| self.state.current().plugin_components.get(&key))
611            .map(|record| &record.value))
612    }
613
614    pub fn proposed_component(
615        &self,
616        state: &StateKey,
617        entity: &EntityRef,
618        component: &str,
619    ) -> Result<Option<&Value>, CanwuError> {
620        self.require_read(state)?;
621        let Some(owner) = self.state_owners.get(state) else {
622            return Err(CanwuError::new(
623                ErrorCode::UndeclaredStateRead,
624                format!(
625                    "state {}.{} has no registered owner",
626                    state.namespace, state.name
627                ),
628            ));
629        };
630        let key = component_key(owner, state, entity, component);
631        Ok(self
632            .proposed_components
633            .and_then(|proposals| proposals.get(&key))
634            .map(|record| &record.value))
635    }
636
637    fn require_read(&self, state: &StateKey) -> Result<(), CanwuError> {
638        if self
639            .allowed_reads
640            .is_some_and(|reads| !reads.contains(state))
641        {
642            return Err(CanwuError::new(
643                ErrorCode::UndeclaredStateRead,
644                format!(
645                    "{} did not declare read access to {}.{}",
646                    self.reader.unwrap_or("internal system"),
647                    state.namespace,
648                    state.name
649                ),
650            ));
651        }
652        Ok(())
653    }
654
655    pub(super) fn finish_random_session(self) -> Option<random::RandomExecution> {
656        self.random_session
657            .map(RefCell::into_inner)
658            .map(random::RandomSession::finish)
659    }
660}