Skip to main content

canwu_sim/
settlement.rs

1use super::{
2    AssertUnwindSafe, BTreeMap, BTreeSet, BoundaryChange, BoundaryContext, BoundaryDirective,
3    BoundaryEmission, BoundaryEmissionKind, BoundaryId, BoundaryIngressGeneration, BoundaryPhase,
4    BoundaryProposal, BoundaryReceipt, BoundaryRecord, BoundaryRequest, BoundaryStateHashFormat,
5    BoundarySystemContract, BoundaryTransactionCheckpoint, CanwuError, CauseRef, CommandIngress,
6    CommandRequest, CommitmentDomains, DomainRecord, DomainRecordChange, DomainRecordRef,
7    EntityRef, ErrorCode, EventKind, GENESIS_BOUNDARY_HASH, HashSet, IngressPayload,
8    PluginComponentKey, PluginComponentRecord, PluginRegistry, RefCell, ReservationAllocation,
9    ReservationDisposition, ReservationOffer, ReservationOfferRecord, ReservationPoolKey,
10    ReservationRef, ReservationRequest, ReservationRequestRecord, RunConfigurationSnapshot,
11    RuntimeCurrentState, ScheduleKey, ScheduledAction, SimTime, Simulation, SimulationView,
12    SimulationViewState, StateKey, StateVisibility, SystemCadence, SystemDirective, canonical_text,
13    catch_unwind, claim_counter, component_key, compute_boundary_hash, invalid_snapshot_error,
14    is_domain_record_state, proposal_entity_exists, proposal_entity_identity_exists, random,
15    record_change_affected_entities, records, runtime_current_entity_exists, runtime_entity_exists,
16    runtime_entity_exists_with_record_overlay, runtime_entity_identity_exists,
17    validate_domain_dependents_with_records, validate_runtime_domain_dependents,
18};
19
20impl Simulation {
21    pub fn settle_boundary(
22        &mut self,
23        request: BoundaryRequest,
24    ) -> Result<BoundaryReceipt, CanwuError> {
25        self.settle_boundary_with_state_hash_format(request, BoundaryStateHashFormat::CommitmentsV1)
26    }
27
28    pub(super) fn settle_boundary_with_state_hash_format(
29        &mut self,
30        mut request: BoundaryRequest,
31        state_hash_format: BoundaryStateHashFormat,
32    ) -> Result<BoundaryReceipt, CanwuError> {
33        self.ensure_runtime_ready()?;
34        if request.at < self.state.scheduler.now {
35            return Err(CanwuError::new(
36                ErrorCode::InvalidBoundary,
37                "a settlement boundary cannot precede committed simulation time",
38            ));
39        }
40        if self
41            .state
42            .scheduler
43            .pending_ingress
44            .first()
45            .is_some_and(|key| key.due_at < request.at)
46        {
47            return Err(CanwuError::new(
48                ErrorCode::InvalidBoundary,
49                "a settlement boundary cannot step past earlier canonical ingress",
50            ));
51        }
52        if request.cadences.contains(&SystemCadence::EventDriven) {
53            return Err(CanwuError::new(
54                ErrorCode::InvalidBoundary,
55                "event-driven cadence is derived from admitted events, not caller supplied",
56            ));
57        }
58        request.cadences.sort();
59        request.cadences.dedup();
60
61        let transaction = BoundaryTransactionCheckpoint::capture(&self.state);
62        match self.settle_boundary_inner(request, state_hash_format) {
63            Ok(receipt) => Ok(receipt),
64            Err(error) => {
65                transaction.restore(&mut self.state);
66                Err(error)
67            }
68        }
69    }
70
71    fn settle_boundary_inner(
72        &mut self,
73        mut request: BoundaryRequest,
74        state_hash_format: BoundaryStateHashFormat,
75    ) -> Result<BoundaryReceipt, CanwuError> {
76        self.advance_to_before_boundary(request.at)?;
77
78        let admitted_ingress = self.take_due_ingress(request.at);
79        let admitted_ingress_index: HashSet<_> = admitted_ingress.iter().copied().collect();
80        for ingress_id in &admitted_ingress {
81            let record = self
82                .state
83                .evidence
84                .retained_ingress(*ingress_id)
85                .cloned()
86                .ok_or_else(|| {
87                    CanwuError::new(
88                        ErrorCode::InvalidSnapshot,
89                        "pending ingress references an unknown record",
90                    )
91                })?;
92            match record.payload {
93                IngressPayload::Command { request: command } => {
94                    let CommandRequest {
95                        request_id,
96                        expected_revision,
97                        envelope,
98                    } = *command;
99                    self.admit_command(
100                        Some(request_id),
101                        Some(expected_revision),
102                        envelope,
103                        CommandIngress::LiveRequest,
104                        true,
105                    )?;
106                }
107                IngressPayload::Calendar { cadences } => request.cadences.extend(cadences),
108                IngressPayload::Plugin { .. } => {}
109            }
110        }
111        self.execute_scheduled_at(request.at)?;
112        request.cadences.sort();
113        request.cadences.dedup();
114
115        let admitted_attempt_count = self
116            .state
117            .evidence
118            .archived
119            .command_attempt_count
120            .checked_add(
121                u64::try_from(self.state.evidence.command_attempts.len()).map_err(|_| {
122                    invalid_snapshot_error("attempt journal exceeds admission cursor range")
123                })?,
124            )
125            .ok_or_else(|| invalid_snapshot_error("attempt journal cursor is exhausted"))?;
126        let admitted_command_count = self
127            .state
128            .evidence
129            .archived
130            .command_count
131            .checked_add(
132                u64::try_from(self.state.evidence.commands.len()).map_err(|_| {
133                    invalid_snapshot_error("command journal exceeds admission cursor range")
134                })?,
135            )
136            .ok_or_else(|| invalid_snapshot_error("command journal cursor is exhausted"))?;
137        let admitted_event_count = self
138            .state
139            .evidence
140            .archived
141            .event_count
142            .checked_add(
143                u64::try_from(self.state.evidence.events.len()).map_err(|_| {
144                    invalid_snapshot_error("event journal exceeds admission cursor range")
145                })?,
146            )
147            .ok_or_else(|| invalid_snapshot_error("event journal cursor is exhausted"))?;
148        let admitted_attempt_start = self
149            .state
150            .counters
151            .admitted_attempt_count
152            .checked_sub(self.state.evidence.archived.command_attempt_count)
153            .ok_or_else(|| {
154                invalid_snapshot_error("runtime attempt admission cursor precedes live evidence")
155            })?;
156        let admitted_attempt_start = usize::try_from(admitted_attempt_start).map_err(|_| {
157            invalid_snapshot_error("runtime attempt admission cursor exceeds platform range")
158        })?;
159        let admitted_command_start = self
160            .state
161            .counters
162            .admitted_command_count
163            .checked_sub(self.state.evidence.archived.command_count)
164            .ok_or_else(|| {
165                invalid_snapshot_error("runtime command admission cursor precedes live evidence")
166            })?;
167        let admitted_command_start = usize::try_from(admitted_command_start).map_err(|_| {
168            invalid_snapshot_error("runtime command admission cursor exceeds platform range")
169        })?;
170        let admitted_event_start = self
171            .state
172            .counters
173            .admitted_event_count
174            .checked_sub(self.state.evidence.archived.event_count)
175            .ok_or_else(|| {
176                invalid_snapshot_error("runtime event admission cursor precedes live evidence")
177            })?;
178        let admitted_event_start = usize::try_from(admitted_event_start).map_err(|_| {
179            invalid_snapshot_error("runtime event admission cursor exceeds platform range")
180        })?;
181        let admitted_attempts: Vec<_> = self
182            .state
183            .evidence
184            .command_attempts
185            .get(admitted_attempt_start..)
186            .ok_or_else(|| {
187                invalid_snapshot_error("runtime attempt admission cursor exceeds its journal")
188            })?
189            .iter()
190            .map(|record| record.id)
191            .collect();
192        let admitted_commands: Vec<_> = self
193            .state
194            .evidence
195            .commands
196            .get(admitted_command_start..)
197            .ok_or_else(|| {
198                invalid_snapshot_error("runtime command admission cursor exceeds its journal")
199            })?
200            .iter()
201            .map(|record| record.id)
202            .collect();
203        let admitted_events: Vec<_> = self
204            .state
205            .evidence
206            .events
207            .get(admitted_event_start..)
208            .ok_or_else(|| {
209                invalid_snapshot_error("runtime event admission cursor exceeds its journal")
210            })?
211            .iter()
212            .map(|event| event.id)
213            .collect();
214
215        let (boundary_id_value, next_boundary_id) =
216            claim_counter(self.state.counters.next_boundary_id, "boundary ID")?;
217        let (correlation_id, next_correlation_id) = claim_counter(
218            self.state.counters.next_correlation_id,
219            "boundary correlation ID",
220        )?;
221        self.state.counters.next_boundary_id = next_boundary_id;
222        self.state.counters.next_correlation_id = next_correlation_id;
223        let boundary_id = BoundaryId::new(boundary_id_value);
224
225        let boundary_snapshot = self.state.current.clone();
226        let boundary_time = self.state.scheduler.now;
227        let systems = self.plugins.boundary_systems.clone();
228        let state_owners = self.plugins.state_owners.clone();
229        let record_schemas = self.plugins.record_schemas.clone();
230        let mut allocations = BTreeMap::new();
231        let mut allocation_records = Vec::new();
232        let mut reservation_offer_records = Vec::new();
233        let mut reservation_request_records = Vec::new();
234        let mut offers = Vec::new();
235        let mut requests = Vec::new();
236        let mut random_overlay = boundary_snapshot.random_streams.clone();
237        let mut pending_random_draws = Vec::new();
238        let mut visible_overlay = BTreeMap::new();
239        let mut candidate_overlay = BTreeMap::new();
240        let mut visible_record_overlay = BTreeMap::new();
241        let mut candidate_record_overlay = BTreeMap::new();
242        let mut ordinary = Vec::new();
243        let mut transitions = Vec::new();
244        let mut deferred = Vec::new();
245        let mut evidence = PendingBoundaryEvidence::default();
246
247        for phase in BoundaryPhase::ALL {
248            match phase {
249                BoundaryPhase::AtomicDomainCommit => {
250                    let (same_boundary, next_boundary) =
251                        partition_boundary_visibility(std::mem::take(&mut ordinary));
252                    self.apply_boundary_stage(
253                        boundary_id,
254                        correlation_id,
255                        same_boundary,
256                        &mut evidence,
257                    )?;
258                    deferred.extend(next_boundary);
259                    visible_overlay.clear();
260                    candidate_overlay.clear();
261                    visible_record_overlay.clear();
262                    candidate_record_overlay.clear();
263                }
264                BoundaryPhase::ConditionalTransitionCommit => {
265                    let (same_boundary, next_boundary) =
266                        partition_boundary_visibility(std::mem::take(&mut transitions));
267                    self.apply_boundary_stage(
268                        boundary_id,
269                        correlation_id,
270                        same_boundary,
271                        &mut evidence,
272                    )?;
273                    deferred.extend(next_boundary);
274                    visible_overlay.clear();
275                    visible_record_overlay.clear();
276                }
277                _ => {}
278            }
279
280            let mut phase_directives = Vec::new();
281            for registered in systems.iter().filter(|registered| {
282                registered.contract.phase == phase
283                    && boundary_system_due(
284                        &registered.contract,
285                        &request.cadences,
286                        !admitted_events.is_empty() || !admitted_ingress.is_empty(),
287                    )
288            }) {
289                let reader = format!("{}.{}", registered.plugin, registered.contract.name);
290                let (view_current, view_now) = if phase <= BoundaryPhase::InvariantValidation {
291                    (&boundary_snapshot, boundary_time)
292                } else {
293                    (&self.state.current, self.state.scheduler.now)
294                };
295                let random_session = random::RandomSession::new(
296                    &random_overlay,
297                    &registered.contract.random_streams,
298                )?;
299                let view = SimulationView {
300                    state: SimulationViewState::Boundary {
301                        current: view_current,
302                        now: view_now,
303                        evidence: &self.state.evidence,
304                    },
305                    state_owners: &state_owners,
306                    reader: Some(&reader),
307                    allowed_reads: Some(&registered.contract.reads),
308                    allowed_ingress: Some(&admitted_ingress_index),
309                    ingress_plugin: Some(&registered.plugin),
310                    component_overlay: Some(&visible_overlay),
311                    proposed_components: (phase == BoundaryPhase::InvariantValidation)
312                        .then_some(&candidate_overlay),
313                    record_overlay: Some(&visible_record_overlay),
314                    proposed_records: (phase == BoundaryPhase::InvariantValidation)
315                        .then_some(&candidate_record_overlay),
316                    allocations: Some(&allocations),
317                    allowed_reservations: Some(&registered.contract.reservation_reads),
318                    random_session: Some(RefCell::new(random_session)),
319                };
320                let context = BoundaryContext {
321                    boundary_id,
322                    at: request.at,
323                    phase,
324                    plugin: registered.plugin.clone(),
325                    system: registered.contract.name.clone(),
326                    admitted_attempts: admitted_attempts.clone(),
327                    admitted_commands: admitted_commands.clone(),
328                    admitted_ingress: admitted_ingress.clone(),
329                    admitted_events: admitted_events.clone(),
330                    emitted_events: evidence
331                        .emissions
332                        .iter()
333                        .map(|emission| emission.event)
334                        .collect(),
335                };
336                let proposal =
337                    catch_unwind(AssertUnwindSafe(|| (registered.handler)(&view, &context)))
338                        .map_err(|_| {
339                            CanwuError::new(
340                                ErrorCode::PluginPanicked,
341                                format!(
342                                    "boundary system {}.{} panicked",
343                                    registered.plugin, registered.contract.name
344                                ),
345                            )
346                        })??;
347                validate_boundary_proposal(
348                    &registered.plugin,
349                    &registered.contract,
350                    view_current,
351                    view_now,
352                    &self.plugins,
353                    &visible_record_overlay,
354                    &proposal,
355                )?;
356                let random_execution = view
357                    .finish_random_session()
358                    .expect("boundary views always have a random session");
359                random_overlay.extend(random_execution.states);
360                pending_random_draws.extend(random_execution.draws.into_iter().map(|draw| {
361                    PendingBoundaryRandomDraw {
362                        plugin: registered.plugin.clone(),
363                        system: registered.contract.name.clone(),
364                        draw,
365                    }
366                }));
367                offers.extend(
368                    proposal
369                        .offers
370                        .into_iter()
371                        .map(|offer| PendingReservationOffer {
372                            plugin: registered.plugin.clone(),
373                            system: registered.contract.name.clone(),
374                            offer,
375                        }),
376                );
377                requests.extend(proposal.requests.into_iter().map(|request| {
378                    PendingReservationRequest {
379                        reservation: ReservationRef::new(
380                            &registered.plugin,
381                            &registered.contract.name,
382                            &request.request,
383                        ),
384                        request,
385                    }
386                }));
387                phase_directives.extend(proposal.directives.into_iter().map(|directive| {
388                    StagedBoundaryDirective {
389                        plugin: registered.plugin.clone(),
390                        system: registered.contract.name.clone(),
391                        phase,
392                        visibility: registered.contract.visibility,
393                        directive,
394                    }
395                }));
396            }
397
398            match phase {
399                BoundaryPhase::ReservationAndAllocation => {
400                    let result = allocate_reservations(
401                        std::mem::take(&mut offers),
402                        std::mem::take(&mut requests),
403                    )?;
404                    allocations = result.by_reservation;
405                    allocation_records = result.records;
406                    reservation_offer_records = result.offers;
407                    reservation_request_records = result.requests;
408                }
409                BoundaryPhase::DomainDeltaProposal => {
410                    let record_context = BoundaryRecordOverlayContext {
411                        current: &boundary_snapshot,
412                        now: boundary_time,
413                        scheduled_actions: &self.state.scheduler.actions,
414                        run_configuration: &self.state.metadata.run_configuration,
415                        schemas: &record_schemas,
416                    };
417                    extend_boundary_record_candidate_overlay(
418                        &record_context,
419                        &mut candidate_record_overlay,
420                        &phase_directives,
421                    )?;
422                    extend_boundary_candidate_overlay(
423                        &boundary_snapshot,
424                        &candidate_record_overlay,
425                        &mut candidate_overlay,
426                        &phase_directives,
427                    )?;
428                    extend_boundary_record_overlay(
429                        &record_context,
430                        &mut visible_record_overlay,
431                        &phase_directives,
432                    )?;
433                    extend_boundary_overlay(
434                        &boundary_snapshot,
435                        &visible_record_overlay,
436                        &mut visible_overlay,
437                        &phase_directives,
438                    )?;
439                    ordinary.extend(phase_directives);
440                }
441                BoundaryPhase::HistoricalCandidateEvaluation => {
442                    let record_context = BoundaryRecordOverlayContext {
443                        current: &self.state.current,
444                        now: self.state.scheduler.now,
445                        scheduled_actions: &self.state.scheduler.actions,
446                        run_configuration: &self.state.metadata.run_configuration,
447                        schemas: &record_schemas,
448                    };
449                    extend_boundary_record_overlay(
450                        &record_context,
451                        &mut visible_record_overlay,
452                        &phase_directives,
453                    )?;
454                    extend_boundary_overlay(
455                        &self.state.current,
456                        &visible_record_overlay,
457                        &mut visible_overlay,
458                        &phase_directives,
459                    )?;
460                    transitions.extend(phase_directives);
461                }
462                BoundaryPhase::StrategicAggregation
463                | BoundaryPhase::PerspectiveAndReportMaterialization => {
464                    let (same_boundary, next_boundary) =
465                        partition_boundary_visibility(phase_directives);
466                    self.apply_boundary_stage(
467                        boundary_id,
468                        correlation_id,
469                        same_boundary,
470                        &mut evidence,
471                    )?;
472                    deferred.extend(next_boundary);
473                }
474                _ if !phase_directives.is_empty() => {
475                    return Err(CanwuError::new(
476                        ErrorCode::InvalidBoundary,
477                        format!("boundary phase {phase:?} cannot produce state directives"),
478                    ));
479                }
480                _ => {}
481            }
482        }
483
484        self.apply_boundary_stage(boundary_id, correlation_id, deferred, &mut evidence)?;
485        let PendingBoundaryEvidence {
486            changes,
487            record_changes,
488            emissions,
489            generated_ingress,
490        } = evidence;
491        self.invalidate_commitments(CommitmentDomains::RANDOM_STREAMS);
492        self.state.current.random_streams = random_overlay;
493        let random_draws =
494            self.append_boundary_random_draws(boundary_id, correlation_id, pending_random_draws)?;
495        self.state.metadata.plugin_registration_closed = true;
496        let state_hash = self.compute_boundary_state_hash_for(state_hash_format)?;
497        let previous_hash = self
498            .state
499            .evidence
500            .boundary_head_hash()
501            .map_or_else(|| GENESIS_BOUNDARY_HASH.to_owned(), str::to_owned);
502        let mut record = BoundaryRecord {
503            id: boundary_id,
504            at: request.at,
505            correlation_id,
506            cadences: request.cadences,
507            admitted_attempts,
508            admitted_commands,
509            admitted_ingress,
510            generated_ingress: generated_ingress.clone(),
511            admitted_events,
512            reservation_offers: reservation_offer_records,
513            reservation_requests: reservation_request_records,
514            allocations: allocation_records.clone(),
515            random_draws: random_draws.clone(),
516            changes: changes.clone(),
517            record_changes: record_changes.clone(),
518            emissions: emissions.clone(),
519            state_hash: Some(state_hash),
520            previous_hash,
521            hash: String::new(),
522        };
523        record.hash = compute_boundary_hash(&record)?;
524        let boundary_hash = record.hash.clone();
525        self.state.evidence.boundaries.push(record);
526        self.state.counters.admitted_attempt_count = admitted_attempt_count;
527        self.state.counters.admitted_command_count = admitted_command_count;
528        self.state.counters.admitted_event_count = admitted_event_count;
529        self.advance_state_revision()?;
530        self.refresh_checkpoint_hash()?;
531        Ok(BoundaryReceipt {
532            boundary_id,
533            settled_at: request.at,
534            emitted_events: emissions
535                .into_iter()
536                .map(|emission| emission.event)
537                .collect(),
538            generated_ingress: generated_ingress
539                .into_iter()
540                .map(|generation| generation.ingress)
541                .collect(),
542            random_draws,
543            boundary_hash,
544            change_count: changes.len(),
545            record_change_count: record_changes.len(),
546            allocations: allocation_records,
547        })
548    }
549
550    fn apply_boundary_stage(
551        &mut self,
552        boundary_id: BoundaryId,
553        correlation_id: u64,
554        directives: Vec<StagedBoundaryDirective>,
555        evidence: &mut PendingBoundaryEvidence,
556    ) -> Result<(), CanwuError> {
557        let changes = &mut evidence.changes;
558        let record_changes = &mut evidence.record_changes;
559        let emissions = &mut evidence.emissions;
560        let generated_ingress = &mut evidence.generated_ingress;
561        let mutation_requests: Vec<_> = directives
562            .iter()
563            .filter_map(|staged| match &staged.directive {
564                BoundaryDirective::MutateRecord { mutation, summary } => {
565                    Some(records::DomainMutationRequest {
566                        plugin: &staged.plugin,
567                        system: &staged.system,
568                        visibility: staged.visibility,
569                        mutation,
570                        summary,
571                    })
572                }
573                BoundaryDirective::SetComponent { .. }
574                | BoundaryDirective::Emit { .. }
575                | BoundaryDirective::ScheduleIngress { .. } => None,
576            })
577            .collect();
578        let mut stage_record_changes = BTreeMap::new();
579        if !mutation_requests.is_empty() {
580            let (next_records, applied) = records::apply_mutation_bundle(
581                &self.state.current.domain_records,
582                &self.plugins.record_schemas,
583                self.state.scheduler.now,
584                &|entity| runtime_entity_exists(&self.state, entity),
585                mutation_requests,
586            )?;
587            let first_index = record_changes.len();
588            for (offset, change) in applied.iter().enumerate() {
589                let index = first_index.checked_add(offset).ok_or_else(|| {
590                    CanwuError::new(
591                        ErrorCode::IdentifierExhausted,
592                        "boundary record-change index exceeds the persistent identifier space",
593                    )
594                })?;
595                let index = u64::try_from(index).map_err(|_| {
596                    CanwuError::new(
597                        ErrorCode::IdentifierExhausted,
598                        "boundary record-change index exceeds the persistent identifier space",
599                    )
600                })?;
601                stage_record_changes
602                    .insert(change.current.reference.clone(), (index, change.clone()));
603            }
604            self.invalidate_commitments(CommitmentDomains::DOMAIN_RECORDS);
605            self.state.current.domain_records = next_records;
606            record_changes.extend(applied);
607        }
608
609        for staged in &directives {
610            let unavailable = match &staged.directive {
611                BoundaryDirective::SetComponent { entity, .. } => {
612                    (!runtime_entity_exists(&self.state, entity)).then_some(entity)
613                }
614                BoundaryDirective::Emit { affected, .. } => affected
615                    .iter()
616                    .find(|entity| !runtime_entity_exists(&self.state, entity)),
617                BoundaryDirective::ScheduleIngress { affected, .. } => affected
618                    .iter()
619                    .find(|entity| !runtime_entity_identity_exists(&self.state, entity)),
620                BoundaryDirective::MutateRecord { .. } => None,
621            };
622            if let Some(entity) = unavailable {
623                return Err(CanwuError::new(
624                    ErrorCode::EntityNotFound,
625                    format!(
626                        "boundary stage {}.{} references unavailable entity {entity}",
627                        staged.plugin, staged.system
628                    ),
629                )
630                .with_entity(entity.clone()));
631            }
632        }
633
634        for staged in directives {
635            match staged.directive {
636                BoundaryDirective::SetComponent {
637                    state,
638                    entity,
639                    component,
640                    value,
641                    summary,
642                } => {
643                    let key = component_key(&staged.plugin, &state, &entity, &component);
644                    self.invalidate_commitments(CommitmentDomains::PLUGIN_COMPONENTS);
645                    let previous = self
646                        .state
647                        .current
648                        .plugin_components
649                        .get(&key)
650                        .map(|record| record.value.clone());
651                    self.state.current.plugin_components.insert(
652                        key,
653                        PluginComponentRecord {
654                            plugin: staged.plugin.clone(),
655                            state: state.clone(),
656                            entity: entity.clone(),
657                            component: component.clone(),
658                            value: value.clone(),
659                        },
660                    );
661                    let change_index = u64::try_from(changes.len()).map_err(|_| {
662                        CanwuError::new(
663                            ErrorCode::IdentifierExhausted,
664                            "boundary change index exceeds the persistent identifier space",
665                        )
666                    })?;
667                    changes.push(BoundaryChange {
668                        plugin: staged.plugin.clone(),
669                        system: staged.system.clone(),
670                        state,
671                        entity: entity.clone(),
672                        component: component.clone(),
673                        previous,
674                        value,
675                        visibility: staged.visibility,
676                        summary: summary.clone(),
677                    });
678                    let event = self.append_event(
679                        EventKind::Plugin {
680                            plugin: staged.plugin.clone(),
681                            event_type: format!("{component}_changed"),
682                        },
683                        vec![entity],
684                        summary,
685                        Some(CauseRef::Boundary(boundary_id)),
686                        correlation_id,
687                    )?;
688                    emissions.push(BoundaryEmission {
689                        plugin: staged.plugin,
690                        system: staged.system,
691                        event: event.id,
692                        kind: BoundaryEmissionKind::Change { change_index },
693                    });
694                }
695                BoundaryDirective::MutateRecord { mutation, .. } => {
696                    let Some((change_index, change)) = stage_record_changes.get(mutation.target())
697                    else {
698                        return Err(CanwuError::new(
699                            ErrorCode::InvalidBoundary,
700                            "record mutation is missing its committed change evidence",
701                        ));
702                    };
703                    let event = self.append_event(
704                        EventKind::Plugin {
705                            plugin: staged.plugin.clone(),
706                            event_type: change.operation.event_type().to_owned(),
707                        },
708                        record_change_affected_entities(change),
709                        change.summary.clone(),
710                        Some(CauseRef::Boundary(boundary_id)),
711                        correlation_id,
712                    )?;
713                    emissions.push(BoundaryEmission {
714                        plugin: staged.plugin,
715                        system: staged.system,
716                        event: event.id,
717                        kind: BoundaryEmissionKind::RecordChange {
718                            change_index: *change_index,
719                        },
720                    });
721                }
722                BoundaryDirective::Emit {
723                    event_type,
724                    summary,
725                    affected,
726                } => {
727                    let event = self.append_event(
728                        EventKind::Plugin {
729                            plugin: staged.plugin.clone(),
730                            event_type,
731                        },
732                        affected,
733                        summary,
734                        Some(CauseRef::Boundary(boundary_id)),
735                        correlation_id,
736                    )?;
737                    emissions.push(BoundaryEmission {
738                        plugin: staged.plugin,
739                        system: staged.system,
740                        event: event.id,
741                        kind: BoundaryEmissionKind::Explicit,
742                    });
743                }
744                BoundaryDirective::ScheduleIngress {
745                    after,
746                    packet_type,
747                    priority,
748                    payload,
749                    mut affected,
750                } => {
751                    self.ensure_canonical_ingress_can_start()?;
752                    let descriptor = self
753                        .plugins
754                        .ingress
755                        .get(&(staged.plugin.clone(), packet_type.clone()))
756                        .ok_or_else(|| {
757                            CanwuError::new(
758                                ErrorCode::InvalidPayload,
759                                format!(
760                                    "boundary system {}.{} scheduled undeclared ingress type {packet_type}",
761                                    staged.plugin, staged.system
762                                ),
763                            )
764                        })?
765                        .clone();
766                    descriptor.payload_schema.validate(&payload)?;
767                    affected.sort();
768                    affected.dedup();
769                    let due_at = self.state.scheduler.now.checked_add(after).ok_or_else(|| {
770                        CanwuError::new(
771                            ErrorCode::InvalidDuration,
772                            "boundary-generated ingress exceeds the supported time range",
773                        )
774                    })?;
775                    let receipt = self.append_ingress(
776                        due_at,
777                        descriptor.class,
778                        priority,
779                        IngressPayload::Plugin {
780                            plugin: staged.plugin.clone(),
781                            packet_type,
782                            payload,
783                            affected_entities: affected,
784                        },
785                        Some(CauseRef::Boundary(boundary_id)),
786                        true,
787                    )?;
788                    generated_ingress.push(BoundaryIngressGeneration {
789                        ingress: receipt.ingress_id,
790                        plugin: staged.plugin,
791                        system: staged.system,
792                        phase: staged.phase,
793                        visibility: staged.visibility,
794                    });
795                }
796            }
797        }
798        validate_runtime_domain_dependents(&self.state)?;
799        Ok(())
800    }
801
802    pub(super) fn apply_directives(
803        &mut self,
804        plugin: &str,
805        directives: Vec<SystemDirective>,
806        allowed_writes: &[StateKey],
807        cause: &CauseRef,
808        correlation_id: u64,
809    ) -> Result<(), CanwuError> {
810        for directive in directives {
811            match directive {
812                SystemDirective::SetComponent {
813                    state,
814                    entity,
815                    component,
816                    value,
817                    summary,
818                } => {
819                    let key = component_key(plugin, &state, &entity, &component);
820                    self.invalidate_commitments(CommitmentDomains::PLUGIN_COMPONENTS);
821                    self.state.current.plugin_components.insert(
822                        key,
823                        PluginComponentRecord {
824                            plugin: plugin.to_owned(),
825                            state,
826                            entity: entity.clone(),
827                            component: component.clone(),
828                            value,
829                        },
830                    );
831                    self.emit(
832                        EventKind::Plugin {
833                            plugin: plugin.to_owned(),
834                            event_type: format!("{component}_changed"),
835                        },
836                        vec![entity],
837                        summary,
838                        Some(cause.clone()),
839                        correlation_id,
840                    )?;
841                }
842                SystemDirective::Emit {
843                    event_type,
844                    summary,
845                    affected,
846                } => {
847                    self.emit(
848                        EventKind::Plugin {
849                            plugin: plugin.to_owned(),
850                            event_type,
851                        },
852                        affected,
853                        summary,
854                        Some(cause.clone()),
855                        correlation_id,
856                    )?;
857                }
858                SystemDirective::Schedule { after, directive } => {
859                    let at = self.state.scheduler.now.checked_add(after).ok_or_else(|| {
860                        CanwuError::new(
861                            ErrorCode::InvalidDuration,
862                            "plugin scheduled time exceeds the supported range",
863                        )
864                    })?;
865                    self.schedule_at(
866                        at,
867                        ScheduledAction::PluginDirective {
868                            plugin: plugin.to_owned(),
869                            directive,
870                            allowed_writes: allowed_writes.to_vec(),
871                            cause: cause.clone(),
872                            correlation_id,
873                        },
874                    )?;
875                }
876            }
877        }
878        Ok(())
879    }
880}
881
882struct PendingReservationOffer {
883    plugin: String,
884    system: String,
885    offer: ReservationOffer,
886}
887
888struct PendingReservationRequest {
889    reservation: ReservationRef,
890    request: ReservationRequest,
891}
892
893struct ReservationAllocationResult {
894    by_reservation: BTreeMap<ReservationRef, ReservationAllocation>,
895    offers: Vec<ReservationOfferRecord>,
896    requests: Vec<ReservationRequestRecord>,
897    records: Vec<ReservationAllocation>,
898}
899
900struct StagedBoundaryDirective {
901    plugin: String,
902    system: String,
903    phase: BoundaryPhase,
904    visibility: StateVisibility,
905    directive: BoundaryDirective,
906}
907
908#[derive(Default)]
909struct PendingBoundaryEvidence {
910    changes: Vec<BoundaryChange>,
911    record_changes: Vec<DomainRecordChange>,
912    emissions: Vec<BoundaryEmission>,
913    generated_ingress: Vec<BoundaryIngressGeneration>,
914}
915
916pub(super) struct PendingBoundaryRandomDraw {
917    pub(super) plugin: String,
918    pub(super) system: String,
919    pub(super) draw: random::PendingRandomDraw,
920}
921
922pub(super) fn boundary_system_due(
923    contract: &BoundarySystemContract,
924    cadences: &[SystemCadence],
925    has_admitted_events: bool,
926) -> bool {
927    match contract.cadence {
928        SystemCadence::EventDriven => has_admitted_events,
929        _ => cadences.contains(&contract.cadence),
930    }
931}
932
933pub(super) fn boundary_has_event_ingress(record: &BoundaryRecord) -> bool {
934    !record.admitted_events.is_empty() || !record.admitted_ingress.is_empty()
935}
936
937fn validate_boundary_proposal(
938    plugin: &str,
939    contract: &BoundarySystemContract,
940    current: &RuntimeCurrentState,
941    now: SimTime,
942    plugins: &PluginRegistry,
943    record_overlay: &BTreeMap<DomainRecordRef, DomainRecord>,
944    proposal: &BoundaryProposal,
945) -> Result<(), CanwuError> {
946    if contract.phase != BoundaryPhase::ReservationAndAllocation
947        && (!proposal.offers.is_empty() || !proposal.requests.is_empty())
948    {
949        return Err(CanwuError::new(
950            ErrorCode::InvalidBoundary,
951            format!(
952                "boundary system {plugin}.{} proposed reservations in phase {:?}",
953                contract.name, contract.phase
954            ),
955        ));
956    }
957
958    let entity_exists = |entity: &EntityRef| {
959        proposal_entity_exists(
960            current,
961            &plugins.record_schemas,
962            record_overlay,
963            proposal,
964            entity,
965        )
966    };
967    let mut offered_pools = BTreeSet::new();
968    for offer in &proposal.offers {
969        validate_reservation_pool(&offer.pool, &entity_exists)?;
970        if !contract.reservation_offers.contains(&offer.pool.state)
971            || plugins
972                .state_owners
973                .get(&offer.pool.state)
974                .is_none_or(|owner| owner != plugin)
975        {
976            return Err(CanwuError::new(
977                ErrorCode::InvalidBoundary,
978                format!(
979                    "boundary system {plugin}.{} offered undeclared state {}.{}",
980                    contract.name, offer.pool.state.namespace, offer.pool.state.name
981                ),
982            ));
983        }
984        if !offered_pools.insert(&offer.pool) {
985            return Err(CanwuError::new(
986                ErrorCode::InvalidBoundary,
987                format!(
988                    "boundary system {plugin}.{} offered the same reservation pool twice",
989                    contract.name
990                ),
991            ));
992        }
993    }
994
995    let mut request_names = BTreeSet::new();
996    for request in &proposal.requests {
997        validate_reservation_pool(&request.pool, &entity_exists)?;
998        if request.request.trim().is_empty()
999            || request.request != request.request.trim()
1000            || request.tie_break.trim().is_empty()
1001            || request.tie_break != request.tie_break.trim()
1002            || request.quantity == 0
1003            || !request_names.insert(&request.request)
1004            || !contract.reservation_requests.contains(&request.pool.state)
1005        {
1006            return Err(CanwuError::new(
1007                ErrorCode::InvalidBoundary,
1008                format!(
1009                    "boundary system {plugin}.{} produced an invalid reservation request",
1010                    contract.name
1011                ),
1012            ));
1013        }
1014    }
1015
1016    let mut component_keys = BTreeSet::new();
1017    let mut record_targets = BTreeSet::new();
1018    for directive in &proposal.directives {
1019        match directive {
1020            BoundaryDirective::SetComponent {
1021                state: state_key,
1022                entity,
1023                component,
1024                ..
1025            } => {
1026                if component.trim().is_empty()
1027                    || component != component.trim()
1028                    || !contract.writes.contains(state_key)
1029                    || plugins
1030                        .state_owners
1031                        .get(state_key)
1032                        .is_none_or(|owner| owner != plugin)
1033                    || is_domain_record_state(&plugins.record_schemas, state_key)
1034                {
1035                    return Err(CanwuError::new(
1036                        ErrorCode::UndeclaredStateWrite,
1037                        format!(
1038                            "boundary system {plugin}.{} produced an undeclared component write",
1039                            contract.name
1040                        ),
1041                    ));
1042                }
1043                if !entity_exists(entity) {
1044                    return Err(CanwuError::new(
1045                        ErrorCode::EntityNotFound,
1046                        format!(
1047                            "boundary system {plugin}.{} targeted missing entity {entity}",
1048                            contract.name
1049                        ),
1050                    )
1051                    .with_entity(entity.clone()));
1052                }
1053                let key = component_key(plugin, state_key, entity, component);
1054                if !component_keys.insert(key) {
1055                    return Err(CanwuError::new(
1056                        ErrorCode::InvalidBoundary,
1057                        format!(
1058                            "boundary system {plugin}.{} wrote the same component twice",
1059                            contract.name
1060                        ),
1061                    ));
1062                }
1063            }
1064            BoundaryDirective::MutateRecord { mutation, summary } => {
1065                let target = mutation.target();
1066                let state_key = records::record_state_key(&target.kind);
1067                if !canonical_text(summary)
1068                    || !contract.writes.contains(&state_key)
1069                    || plugins
1070                        .state_owners
1071                        .get(&state_key)
1072                        .is_none_or(|owner| owner != plugin)
1073                    || plugins
1074                        .record_schemas
1075                        .get(&target.kind)
1076                        .is_none_or(|(owner, _)| owner != plugin)
1077                {
1078                    return Err(CanwuError::new(
1079                        ErrorCode::UndeclaredStateWrite,
1080                        format!(
1081                            "boundary system {plugin}.{} produced an undeclared record mutation",
1082                            contract.name
1083                        ),
1084                    ));
1085                }
1086                if !record_targets.insert(target.clone()) {
1087                    return Err(CanwuError::new(
1088                        ErrorCode::InvalidBoundary,
1089                        format!(
1090                            "boundary system {plugin}.{} mutated the same record twice",
1091                            contract.name
1092                        ),
1093                    ));
1094                }
1095            }
1096            BoundaryDirective::Emit {
1097                event_type,
1098                affected,
1099                ..
1100            } => {
1101                if event_type.trim().is_empty()
1102                    || event_type != event_type.trim()
1103                    || !contract.emits.contains(event_type)
1104                {
1105                    return Err(CanwuError::new(
1106                        ErrorCode::InvalidBoundary,
1107                        format!(
1108                            "boundary system {plugin}.{} emitted an undeclared event type",
1109                            contract.name
1110                        ),
1111                    ));
1112                }
1113                if affected.iter().any(|entity| !entity_exists(entity)) {
1114                    return Err(CanwuError::new(
1115                        ErrorCode::EntityNotFound,
1116                        format!(
1117                            "boundary system {plugin}.{} emitted an event for a missing entity",
1118                            contract.name
1119                        ),
1120                    ));
1121                }
1122            }
1123            BoundaryDirective::ScheduleIngress {
1124                after,
1125                packet_type,
1126                payload,
1127                affected,
1128                ..
1129            } => {
1130                let descriptor = plugins
1131                    .ingress
1132                    .get(&(plugin.to_owned(), packet_type.clone()))
1133                    .ok_or_else(|| {
1134                        CanwuError::new(
1135                            ErrorCode::InvalidPayload,
1136                            format!(
1137                                "boundary system {plugin}.{} scheduled undeclared ingress type {packet_type}",
1138                                contract.name
1139                            ),
1140                        )
1141                    })?;
1142                if after.is_negative() || now.checked_add(*after).is_none() {
1143                    return Err(CanwuError::new(
1144                        ErrorCode::InvalidDuration,
1145                        "boundary-generated ingress requires a nonnegative supported delay",
1146                    ));
1147                }
1148                descriptor.payload_schema.validate(payload)?;
1149                if affected.iter().any(|entity| {
1150                    !proposal_entity_identity_exists(
1151                        current,
1152                        &plugins.record_schemas,
1153                        proposal,
1154                        entity,
1155                    )
1156                }) {
1157                    return Err(CanwuError::new(
1158                        ErrorCode::EntityNotFound,
1159                        format!(
1160                            "boundary system {plugin}.{} scheduled ingress for an unknown entity identity",
1161                            contract.name
1162                        ),
1163                    ));
1164                }
1165            }
1166        }
1167    }
1168    Ok(())
1169}
1170
1171fn validate_reservation_pool(
1172    pool: &ReservationPoolKey,
1173    entity_exists: &dyn Fn(&EntityRef) -> bool,
1174) -> Result<(), CanwuError> {
1175    if pool.resource.trim().is_empty()
1176        || pool.resource != pool.resource.trim()
1177        || !entity_exists(&pool.entity)
1178    {
1179        return Err(CanwuError::new(
1180            ErrorCode::InvalidBoundary,
1181            "reservation pools require a canonical resource and an existing entity",
1182        ));
1183    }
1184    Ok(())
1185}
1186
1187fn extend_boundary_overlay(
1188    current: &RuntimeCurrentState,
1189    record_overlay: &BTreeMap<DomainRecordRef, DomainRecord>,
1190    overlay: &mut BTreeMap<PluginComponentKey, PluginComponentRecord>,
1191    directives: &[StagedBoundaryDirective],
1192) -> Result<(), CanwuError> {
1193    extend_boundary_component_overlay(current, record_overlay, overlay, directives, false)
1194}
1195
1196fn extend_boundary_candidate_overlay(
1197    current: &RuntimeCurrentState,
1198    record_overlay: &BTreeMap<DomainRecordRef, DomainRecord>,
1199    overlay: &mut BTreeMap<PluginComponentKey, PluginComponentRecord>,
1200    directives: &[StagedBoundaryDirective],
1201) -> Result<(), CanwuError> {
1202    extend_boundary_component_overlay(current, record_overlay, overlay, directives, true)
1203}
1204
1205fn extend_boundary_component_overlay(
1206    current: &RuntimeCurrentState,
1207    record_overlay: &BTreeMap<DomainRecordRef, DomainRecord>,
1208    overlay: &mut BTreeMap<PluginComponentKey, PluginComponentRecord>,
1209    directives: &[StagedBoundaryDirective],
1210    include_next_boundary: bool,
1211) -> Result<(), CanwuError> {
1212    for staged in directives.iter().filter(|staged| {
1213        include_next_boundary || staged.visibility == StateVisibility::SameBoundary
1214    }) {
1215        if let BoundaryDirective::SetComponent {
1216            state: state_key,
1217            entity,
1218            component,
1219            value,
1220            ..
1221        } = &staged.directive
1222        {
1223            let key = component_key(&staged.plugin, state_key, entity, component);
1224            if overlay.contains_key(&key) {
1225                return Err(CanwuError::new(
1226                    ErrorCode::InvalidBoundary,
1227                    "multiple boundary proposals target the same component",
1228                ));
1229            }
1230            if !runtime_entity_exists_with_record_overlay(current, record_overlay, entity) {
1231                return Err(CanwuError::new(
1232                    ErrorCode::EntityNotFound,
1233                    format!("boundary proposal targeted missing entity {entity}"),
1234                ));
1235            }
1236            overlay.insert(
1237                key,
1238                PluginComponentRecord {
1239                    plugin: staged.plugin.clone(),
1240                    state: state_key.clone(),
1241                    entity: entity.clone(),
1242                    component: component.clone(),
1243                    value: value.clone(),
1244                },
1245            );
1246        }
1247    }
1248    Ok(())
1249}
1250
1251fn extend_boundary_record_overlay(
1252    context: &BoundaryRecordOverlayContext<'_>,
1253    overlay: &mut BTreeMap<DomainRecordRef, DomainRecord>,
1254    directives: &[StagedBoundaryDirective],
1255) -> Result<(), CanwuError> {
1256    extend_boundary_domain_record_overlay(context, overlay, directives, false)
1257}
1258
1259fn extend_boundary_record_candidate_overlay(
1260    context: &BoundaryRecordOverlayContext<'_>,
1261    overlay: &mut BTreeMap<DomainRecordRef, DomainRecord>,
1262    directives: &[StagedBoundaryDirective],
1263) -> Result<(), CanwuError> {
1264    extend_boundary_domain_record_overlay(context, overlay, directives, true)
1265}
1266
1267struct BoundaryRecordOverlayContext<'a> {
1268    current: &'a RuntimeCurrentState,
1269    now: SimTime,
1270    scheduled_actions: &'a BTreeMap<ScheduleKey, ScheduledAction>,
1271    run_configuration: &'a RunConfigurationSnapshot,
1272    schemas: &'a records::DomainRecordSchemas,
1273}
1274
1275fn extend_boundary_domain_record_overlay(
1276    context: &BoundaryRecordOverlayContext<'_>,
1277    overlay: &mut BTreeMap<DomainRecordRef, DomainRecord>,
1278    directives: &[StagedBoundaryDirective],
1279    include_next_boundary: bool,
1280) -> Result<(), CanwuError> {
1281    let mut base = context.current.domain_records.clone();
1282    base.extend(
1283        overlay
1284            .iter()
1285            .map(|(reference, record)| (reference.clone(), record.clone())),
1286    );
1287    let requests: Vec<_> = directives
1288        .iter()
1289        .filter(|staged| {
1290            include_next_boundary || staged.visibility == StateVisibility::SameBoundary
1291        })
1292        .filter_map(|staged| match &staged.directive {
1293            BoundaryDirective::MutateRecord { mutation, summary } => {
1294                Some(records::DomainMutationRequest {
1295                    plugin: &staged.plugin,
1296                    system: &staged.system,
1297                    visibility: staged.visibility,
1298                    mutation,
1299                    summary,
1300                })
1301            }
1302            BoundaryDirective::SetComponent { .. }
1303            | BoundaryDirective::Emit { .. }
1304            | BoundaryDirective::ScheduleIngress { .. } => None,
1305        })
1306        .collect();
1307    if requests.is_empty() {
1308        return Ok(());
1309    }
1310    let (next, changes) = records::apply_mutation_bundle(
1311        &base,
1312        context.schemas,
1313        context.now,
1314        &|entity| runtime_current_entity_exists(context.current, entity),
1315        requests,
1316    )?;
1317    validate_domain_dependents_with_records(
1318        &context.current.plugin_components,
1319        context.scheduled_actions,
1320        context.run_configuration,
1321        &next,
1322    )?;
1323    for change in changes {
1324        overlay.insert(change.current.reference.clone(), change.current);
1325    }
1326    Ok(())
1327}
1328
1329fn partition_boundary_visibility(
1330    directives: Vec<StagedBoundaryDirective>,
1331) -> (Vec<StagedBoundaryDirective>, Vec<StagedBoundaryDirective>) {
1332    directives
1333        .into_iter()
1334        .partition(|staged| staged.visibility == StateVisibility::SameBoundary)
1335}
1336
1337fn allocate_reservations(
1338    mut offers: Vec<PendingReservationOffer>,
1339    mut requests: Vec<PendingReservationRequest>,
1340) -> Result<ReservationAllocationResult, CanwuError> {
1341    offers.sort_by(|left, right| {
1342        left.offer
1343            .pool
1344            .cmp(&right.offer.pool)
1345            .then_with(|| left.plugin.cmp(&right.plugin))
1346            .then_with(|| left.system.cmp(&right.system))
1347    });
1348    let mut remaining = BTreeMap::new();
1349    let mut offer_records = Vec::new();
1350    for pending in offers {
1351        if remaining
1352            .insert(pending.offer.pool.clone(), pending.offer.capacity)
1353            .is_some()
1354        {
1355            return Err(CanwuError::new(
1356                ErrorCode::InvalidBoundary,
1357                format!(
1358                    "reservation pool was offered more than once, including by {}.{}",
1359                    pending.plugin, pending.system
1360                ),
1361            ));
1362        }
1363        offer_records.push(ReservationOfferRecord {
1364            plugin: pending.plugin,
1365            system: pending.system,
1366            offer: pending.offer,
1367        });
1368    }
1369    requests.sort_by(|left, right| {
1370        left.request
1371            .pool
1372            .cmp(&right.request.pool)
1373            .then_with(|| right.request.priority.cmp(&left.request.priority))
1374            .then_with(|| left.request.tie_break.cmp(&right.request.tie_break))
1375            .then_with(|| left.reservation.cmp(&right.reservation))
1376    });
1377    let mut seen = BTreeSet::new();
1378    let mut by_reservation = BTreeMap::new();
1379    let mut request_records = Vec::new();
1380    let mut records = Vec::new();
1381    for pending in requests {
1382        if !seen.insert(pending.reservation.clone()) {
1383            return Err(CanwuError::new(
1384                ErrorCode::InvalidBoundary,
1385                "reservation request identity is duplicated",
1386            ));
1387        }
1388        request_records.push(ReservationRequestRecord {
1389            reservation: pending.reservation.clone(),
1390            request: pending.request.clone(),
1391        });
1392        let available = remaining.entry(pending.request.pool.clone()).or_default();
1393        let granted = pending.request.quantity.min(*available);
1394        *available -= granted;
1395        let disposition = if granted == pending.request.quantity {
1396            ReservationDisposition::Fulfilled
1397        } else if granted == 0 {
1398            ReservationDisposition::Rejected
1399        } else {
1400            ReservationDisposition::Partial
1401        };
1402        let allocation = ReservationAllocation {
1403            reservation: pending.reservation.clone(),
1404            pool: pending.request.pool,
1405            requested: pending.request.quantity,
1406            granted,
1407            remaining_after: *available,
1408            disposition,
1409        };
1410        by_reservation.insert(pending.reservation, allocation.clone());
1411        records.push(allocation);
1412    }
1413    Ok(ReservationAllocationResult {
1414        by_reservation,
1415        offers: offer_records,
1416        requests: request_records,
1417        records,
1418    })
1419}