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 ®istered.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 ®istered.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(®istered.contract.reads),
308 allowed_ingress: Some(&admitted_ingress_index),
309 ingress_plugin: Some(®istered.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(®istered.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 ®istered.plugin,
349 ®istered.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 ®istered.plugin,
381 ®istered.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}