1use super::event_payloads::{KnowledgePublished, RuntimeEventPayload};
2use super::ingress::{PluginIngressCancellationProof, valid_ingress_cancellation_reason};
3use super::validation::{
4 EvidenceAvailability, RuntimeValidationContext, resolve_evidence_reference,
5};
6use super::{
7 AssertUnwindSafe, BTreeMap, BTreeSet, BoundaryChange, BoundaryContext, BoundaryDirective,
8 BoundaryEmission, BoundaryEmissionKind, BoundaryId, BoundaryIngressGeneration,
9 BoundaryKnowledgeChange, BoundaryPhase, BoundaryProposal, BoundaryReceipt, BoundaryRecord,
10 BoundaryRequest, BoundaryStateHashFormat, BoundarySystemContract,
11 BoundaryTransactionCheckpoint, CanwuError, CauseRef, Command, CommandEnvelope, CommandIngress,
12 CommandRequest, CommitmentDomains, DecisionAction, DecisionIngressRequest, DecisionMutation,
13 DecisionOutcome, DecisionPolicyKind, DecisionRandomEvidence, DecisionStage, DomainRecord,
14 DomainRecordChange, DomainRecordRef, DomainRecordVersionSource, EntityRef, ErrorCode,
15 EventKind, EvidenceRef, GENESIS_BOUNDARY_HASH, HashSet, IngressCancellationAuthority,
16 IngressPayload, KnowledgeHolderRef, KnowledgeRecord, KnowledgeRecordId, PluginComponentKey,
17 PluginComponentRecord, PluginRegistry, PolicyDecision, RandomDrawAddress, RandomDrawOutcome,
18 RandomOperationTarget, RefCell, ReservationAllocation, ReservationDisposition,
19 ReservationOffer, ReservationOfferRecord, ReservationPoolKey, ReservationRef,
20 ReservationRequest, ReservationRequestRecord, RunConfigurationSnapshot, RuntimeCurrentState,
21 RuntimeState, ScheduleKey, ScheduledAction, SimTime, Simulation, SimulationView,
22 SimulationViewState, StateKey, StateVisibility, SystemCadence, SystemDirective, canonical_hash,
23 canonical_text, catch_unwind, claim_counter, component_key, compute_boundary_hash,
24 invalid_snapshot_error, is_domain_record_state, proposal_entity_exists,
25 proposal_entity_identity_exists, random, record_change_affected_entities, records,
26 runtime_current_entity_exists, runtime_entity_exists,
27 runtime_entity_exists_with_record_overlay, runtime_entity_identity_exists,
28 validate_domain_dependents_with_records, validate_runtime_domain_dependents,
29};
30
31impl Simulation {
32 pub fn settle_boundary(
33 &mut self,
34 request: BoundaryRequest,
35 ) -> Result<BoundaryReceipt, CanwuError> {
36 self.settle_boundary_with_state_hash_format(request, BoundaryStateHashFormat::CommitmentsV1)
37 }
38
39 pub(super) fn settle_boundary_with_state_hash_format(
40 &mut self,
41 mut request: BoundaryRequest,
42 state_hash_format: BoundaryStateHashFormat,
43 ) -> Result<BoundaryReceipt, CanwuError> {
44 self.ensure_runtime_ready()?;
45 if request.at < self.state.scheduler.now {
46 return Err(CanwuError::new(
47 ErrorCode::InvalidBoundary,
48 "a settlement boundary cannot precede committed simulation time",
49 ));
50 }
51 if self
52 .state
53 .scheduler
54 .pending_ingress
55 .first()
56 .is_some_and(|key| key.due_at < request.at)
57 {
58 return Err(CanwuError::new(
59 ErrorCode::InvalidBoundary,
60 "a settlement boundary cannot step past earlier canonical ingress",
61 ));
62 }
63 if request.cadences.contains(&SystemCadence::EventDriven) {
64 return Err(CanwuError::new(
65 ErrorCode::InvalidBoundary,
66 "event-driven cadence is derived from admitted events, not caller supplied",
67 ));
68 }
69 request.cadences.sort();
70 request.cadences.dedup();
71
72 let transaction = BoundaryTransactionCheckpoint::capture(&self.state);
73 match self.settle_boundary_inner(request, state_hash_format) {
74 Ok(receipt) => Ok(receipt),
75 Err(error) => {
76 transaction.restore(&mut self.state);
77 Err(error)
78 }
79 }
80 }
81
82 fn settle_boundary_inner(
83 &mut self,
84 mut request: BoundaryRequest,
85 state_hash_format: BoundaryStateHashFormat,
86 ) -> Result<BoundaryReceipt, CanwuError> {
87 self.advance_to_before_boundary(request.at)?;
88
89 let admitted_ingress = self.take_due_ingress(request.at);
90 let admitted_ingress_index: HashSet<_> = admitted_ingress.iter().copied().collect();
91 let mut maintenance_changes = Vec::new();
92 let mut maintenance_record_changes = Vec::new();
93 for ingress_id in &admitted_ingress {
94 let record = self
95 .state
96 .evidence
97 .retained_ingress(*ingress_id)
98 .cloned()
99 .ok_or_else(|| {
100 CanwuError::new(
101 ErrorCode::InvalidSnapshot,
102 "pending ingress references an unknown record",
103 )
104 })?;
105 match record.payload {
106 IngressPayload::Command { request: command } => {
107 let CommandRequest {
108 request_id,
109 expected_revision,
110 envelope,
111 } = *command;
112 self.admit_command(
113 Some(request_id),
114 Some(expected_revision),
115 envelope,
116 CommandIngress::LiveRequest,
117 None,
118 true,
119 )?;
120 }
121 IngressPayload::Calendar { cadences } => request.cadences.extend(cadences),
122 IngressPayload::Plugin { .. } => {}
123 IngressPayload::PluginCancellation { .. } => {
124 return Err(CanwuError::new(
125 ErrorCode::InvalidSnapshot,
126 "a terminal ingress cancellation cannot be admitted",
127 ));
128 }
129 IngressPayload::Decision { request } => {
130 self.apply_decision_request(*request)?;
131 }
132 IngressPayload::Maintenance { request } => {
133 let (change, record_changes) = self.apply_maintenance_request(*request)?;
134 maintenance_changes.push(change);
135 maintenance_record_changes.extend(record_changes);
136 }
137 }
138 }
139 self.state
140 .current
141 .decisions
142 .advance_time(request.at)
143 .map_err(super::decision::decision_error)?;
144 self.invalidate_commitments(CommitmentDomains::DECISIONS);
145 self.execute_scheduled_at(request.at)?;
146 request.cadences.sort();
147 request.cadences.dedup();
148
149 let admitted_attempt_count = self
150 .state
151 .evidence
152 .archived
153 .command_attempt_count
154 .checked_add(
155 u64::try_from(self.state.evidence.command_attempts.len()).map_err(|_| {
156 invalid_snapshot_error("attempt journal exceeds admission cursor range")
157 })?,
158 )
159 .ok_or_else(|| invalid_snapshot_error("attempt journal cursor is exhausted"))?;
160 let admitted_command_count = self
161 .state
162 .evidence
163 .archived
164 .command_count
165 .checked_add(
166 u64::try_from(self.state.evidence.commands.len()).map_err(|_| {
167 invalid_snapshot_error("command journal exceeds admission cursor range")
168 })?,
169 )
170 .ok_or_else(|| invalid_snapshot_error("command journal cursor is exhausted"))?;
171 let admitted_event_count = self
172 .state
173 .evidence
174 .archived
175 .event_count
176 .checked_add(
177 u64::try_from(self.state.evidence.events.len()).map_err(|_| {
178 invalid_snapshot_error("event journal exceeds admission cursor range")
179 })?,
180 )
181 .ok_or_else(|| invalid_snapshot_error("event journal cursor is exhausted"))?;
182 let admitted_attempt_start = self
183 .state
184 .counters
185 .admitted_attempt_count
186 .checked_sub(self.state.evidence.archived.command_attempt_count)
187 .ok_or_else(|| {
188 invalid_snapshot_error("runtime attempt admission cursor precedes live evidence")
189 })?;
190 let admitted_attempt_start = usize::try_from(admitted_attempt_start).map_err(|_| {
191 invalid_snapshot_error("runtime attempt admission cursor exceeds platform range")
192 })?;
193 let admitted_command_start = self
194 .state
195 .counters
196 .admitted_command_count
197 .checked_sub(self.state.evidence.archived.command_count)
198 .ok_or_else(|| {
199 invalid_snapshot_error("runtime command admission cursor precedes live evidence")
200 })?;
201 let admitted_command_start = usize::try_from(admitted_command_start).map_err(|_| {
202 invalid_snapshot_error("runtime command admission cursor exceeds platform range")
203 })?;
204 let admitted_event_start = self
205 .state
206 .counters
207 .admitted_event_count
208 .checked_sub(self.state.evidence.archived.event_count)
209 .ok_or_else(|| {
210 invalid_snapshot_error("runtime event admission cursor precedes live evidence")
211 })?;
212 let admitted_event_start = usize::try_from(admitted_event_start).map_err(|_| {
213 invalid_snapshot_error("runtime event admission cursor exceeds platform range")
214 })?;
215 let admitted_attempts: Vec<_> = self
216 .state
217 .evidence
218 .command_attempts
219 .get(admitted_attempt_start..)
220 .ok_or_else(|| {
221 invalid_snapshot_error("runtime attempt admission cursor exceeds its journal")
222 })?
223 .iter()
224 .map(|record| record.id)
225 .collect();
226 let admitted_commands: Vec<_> = self
227 .state
228 .evidence
229 .commands
230 .get(admitted_command_start..)
231 .ok_or_else(|| {
232 invalid_snapshot_error("runtime command admission cursor exceeds its journal")
233 })?
234 .iter()
235 .map(|record| record.id)
236 .collect();
237 let admitted_events: Vec<_> = self
238 .state
239 .evidence
240 .events
241 .get(admitted_event_start..)
242 .ok_or_else(|| {
243 invalid_snapshot_error("runtime event admission cursor exceeds its journal")
244 })?
245 .iter()
246 .map(|event| event.id)
247 .collect();
248
249 let (boundary_id_value, next_boundary_id) =
250 claim_counter(self.state.counters.next_boundary_id, "boundary ID")?;
251 let (correlation_id, next_correlation_id) = claim_counter(
252 self.state.counters.next_correlation_id,
253 "boundary correlation ID",
254 )?;
255 self.state.counters.next_boundary_id = next_boundary_id;
256 self.state.counters.next_correlation_id = next_correlation_id;
257 let boundary_id = BoundaryId::new(boundary_id_value);
258
259 let boundary_snapshot = self.state.current.clone();
260 let boundary_time = self.state.scheduler.now;
261 let systems = self.plugins.boundary_systems.clone();
262 let state_owners = self.plugins.state_owners.clone();
263 let record_schemas = self.plugins.record_schemas.clone();
264 let mut allocations = BTreeMap::new();
265 let mut allocation_records = Vec::new();
266 let mut reservation_offer_records = Vec::new();
267 let mut reservation_request_records = Vec::new();
268 let mut offers = Vec::new();
269 let mut requests = Vec::new();
270 let mut random_overlay = boundary_snapshot.random_streams.clone();
271 let mut pending_random_draws = Vec::new();
272 let mut keyed_random_draws = random::keyed_draws_with_reservations(
273 &self.state.evidence.random_draws,
274 &self.state.evidence.keyed_draw_reservations,
275 )?;
276 let mut visible_overlay = BTreeMap::new();
277 let mut candidate_overlay = BTreeMap::new();
278 let mut visible_record_overlay = BTreeMap::new();
279 let mut candidate_record_overlay = BTreeMap::new();
280 let mut visible_knowledge_overlay = BTreeMap::new();
281 let mut pending_knowledge_changes = Vec::new();
282 let mut knowledge_correlations = BTreeSet::new();
283 let mut ordinary = Vec::new();
284 let mut transitions = Vec::new();
285 let mut deferred = Vec::new();
286 let mut evidence = PendingBoundaryEvidence::default();
287 let mut person_writes = super::persons::BoundaryPersonWrites::default();
288 let evaluation_limits = self.state.metadata.run_configuration.evaluation_limits();
289 let mut evaluation_traces = Vec::new();
290 let mut transition_ledger = super::transitions::BoundaryTransitionLedger::new(
291 boundary_id,
292 &self.state.scheduler.transition_manifests,
293 );
294 for change in maintenance_record_changes {
295 let change_index = u64::try_from(evidence.record_changes.len()).map_err(|_| {
296 CanwuError::new(
297 ErrorCode::IdentifierExhausted,
298 "boundary record-change index exceeds the persistent identifier space",
299 )
300 })?;
301 index_current_domain_record_version(
302 &mut self.state,
303 boundary_id,
304 change_index,
305 &change,
306 );
307 let event = self.append_event(
308 EventKind::plugin(change.plugin.clone(), change.operation.event_type()),
309 record_change_affected_entities(&change),
310 change.summary.clone(),
311 Some(CauseRef::Boundary(boundary_id)),
312 correlation_id,
313 )?;
314 evidence.emissions.push(BoundaryEmission {
315 plugin: change.plugin.clone(),
316 system: change.system.clone(),
317 event: event.id,
318 kind: BoundaryEmissionKind::RecordChange { change_index },
319 });
320 evidence.record_changes.push(change);
321 }
322 for phase in BoundaryPhase::ALL {
323 match phase {
324 BoundaryPhase::AtomicDomainCommit => {
325 let (same_boundary, next_boundary) =
326 partition_boundary_visibility(std::mem::take(&mut ordinary));
327 self.apply_boundary_stage(
328 boundary_id,
329 correlation_id,
330 same_boundary,
331 &mut evidence,
332 )?;
333 deferred.extend(next_boundary);
334 visible_overlay.clear();
335 candidate_overlay.clear();
336 visible_record_overlay.clear();
337 candidate_record_overlay.clear();
338 }
339 BoundaryPhase::ConditionalTransitionCommit => {
340 transition_ledger.settle_ready(&self.state)?;
343 let (same_boundary, next_boundary) =
344 partition_boundary_visibility(std::mem::take(&mut transitions));
345 self.apply_boundary_stage(
346 boundary_id,
347 correlation_id,
348 same_boundary,
349 &mut evidence,
350 )?;
351 deferred.extend(next_boundary);
352 if transition_ledger.requires_post_check() {
353 self.check_transition_post_versions(
354 &mut transition_ledger,
355 &record_schemas,
356 &deferred,
357 )?;
358 }
359 visible_overlay.clear();
360 visible_record_overlay.clear();
361 }
362 _ => {}
363 }
364
365 let mut phase_directives = Vec::new();
366 for registered in systems.iter().filter(|registered| {
367 registered.contract.phase == phase
368 && boundary_system_due(
369 ®istered.contract,
370 &request.cadences,
371 !admitted_events.is_empty() || !admitted_ingress.is_empty(),
372 )
373 }) {
374 let reader = format!("{}.{}", registered.plugin, registered.contract.name);
375 let (view_current, view_now) = if phase <= BoundaryPhase::InvariantValidation {
376 (&boundary_snapshot, boundary_time)
377 } else {
378 (&self.state.current, self.state.scheduler.now)
379 };
380 let random_session = random::RandomSession::new(
381 &random_overlay,
382 ®istered.contract.random_streams,
383 boundary_snapshot.root_seed,
384 ®istered.plugin,
385 &keyed_random_draws,
386 )?;
387 let proposal_evidence = proposal_evidence_refs(boundary_id, &evidence);
388 let view = SimulationView {
389 state: SimulationViewState::Boundary {
390 current: view_current,
391 now: view_now,
392 runtime: &self.state,
393 },
394 state_owners: &state_owners,
395 reader: Some(&reader),
396 allowed_reads: Some(®istered.contract.reads),
397 allowed_ingress: Some(&admitted_ingress_index),
398 ingress_plugin: Some(®istered.plugin),
399 component_overlay: Some(&visible_overlay),
400 proposed_components: (phase == BoundaryPhase::InvariantValidation)
401 .then_some(&candidate_overlay),
402 record_overlay: Some(&visible_record_overlay),
403 proposed_records: (phase == BoundaryPhase::InvariantValidation)
404 .then_some(&candidate_record_overlay),
405 boundary_id: Some(boundary_id),
406 proposal_evidence: Some(&proposal_evidence),
407 knowledge_overlay: Some(&visible_knowledge_overlay),
408 allocations: Some(&allocations),
409 allowed_reservations: Some(®istered.contract.reservation_reads),
410 random_session: Some(RefCell::new(random_session)),
411 plugin_archive_provider: self.plugin_archive_provider.as_ref(),
412 transitions: Some(&transition_ledger),
413 };
414 let context = BoundaryContext {
415 boundary_id,
416 at: request.at,
417 phase,
418 plugin: registered.plugin.clone(),
419 system: registered.contract.name.clone(),
420 admitted_attempts: admitted_attempts.clone(),
421 admitted_commands: admitted_commands.clone(),
422 admitted_ingress: admitted_ingress.clone(),
423 admitted_events: admitted_events.clone(),
424 emitted_events: evidence
425 .emissions
426 .iter()
427 .map(|emission| emission.event)
428 .collect(),
429 };
430 let proposal =
431 catch_unwind(AssertUnwindSafe(|| (registered.handler)(&view, &context)))
432 .map_err(|_| {
433 CanwuError::new(
434 ErrorCode::PluginPanicked,
435 format!(
436 "boundary system {}.{} panicked",
437 registered.plugin, registered.contract.name
438 ),
439 )
440 })??;
441 let random_execution = view
442 .finish_random_session()
443 .expect("boundary views always have a random session");
444 let proposal = transition_ledger.admit(
447 ®istered.plugin,
448 ®istered.contract,
449 &self.plugins,
450 proposal,
451 )?;
452 super::evaluation::check_trace_budget(
453 evaluation_traces.len(),
454 &proposal.directives,
455 evaluation_limits,
456 )?;
457 validate_boundary_proposal(
458 ®istered.plugin,
459 ®istered.contract,
460 view_current,
461 &boundary_snapshot.person_availability,
462 view_now,
463 &self.state,
464 boundary_id,
465 &evidence,
466 &self.plugins,
467 &visible_record_overlay,
468 &visible_knowledge_overlay,
469 &proposal,
470 &random_execution.draws,
471 )?;
472 person_writes.stage(
473 ®istered.plugin,
474 ®istered.contract.name,
475 &proposal.directives,
476 )?;
477 random::extend_keyed_draws(&mut keyed_random_draws, &random_execution.draws)?;
478 random_overlay.extend(random_execution.states);
479 pending_random_draws.extend(random_execution.draws.into_iter().map(|draw| {
480 PendingBoundaryRandomDraw {
481 plugin: registered.plugin.clone(),
482 system: registered.contract.name.clone(),
483 draw,
484 }
485 }));
486 offers.extend(
487 proposal
488 .offers
489 .into_iter()
490 .map(|offer| PendingReservationOffer {
491 plugin: registered.plugin.clone(),
492 system: registered.contract.name.clone(),
493 offer,
494 }),
495 );
496 requests.extend(proposal.requests.into_iter().map(|request| {
497 PendingReservationRequest {
498 reservation: ReservationRef::new(
499 ®istered.plugin,
500 ®istered.contract.name,
501 &request.request,
502 ),
503 request,
504 }
505 }));
506 let mut state_directives = Vec::with_capacity(proposal.directives.len());
509 let mut traces = Vec::new();
510 for directive in proposal.directives {
511 match directive {
512 BoundaryDirective::RecordEvaluationTrace { trace } => traces.push(trace),
513 directive => state_directives.push(directive),
514 }
515 }
516 super::evaluation::stage_traces(
517 &mut evaluation_traces,
518 ®istered.plugin,
519 ®istered.contract,
520 traces,
521 );
522 phase_directives.extend(state_directives.into_iter().map(|directive| {
523 let visibility = if matches!(directive, BoundaryDirective::CreatePerson { .. })
526 {
527 StateVisibility::NextBoundary
528 } else {
529 registered.contract.visibility
530 };
531 StagedBoundaryDirective {
532 plugin: registered.plugin.clone(),
533 system: registered.contract.name.clone(),
534 phase,
535 visibility,
536 directive,
537 }
538 }));
539 }
540 transition_ledger.close_phase();
541
542 let (knowledge_directives, phase_directives) =
543 partition_knowledge_directives(phase_directives);
544 if !knowledge_directives.is_empty()
545 && !matches!(
546 phase,
547 BoundaryPhase::PerceptionAndAttentionRefresh
548 | BoundaryPhase::PerspectiveAndReportMaterialization
549 )
550 {
551 return Err(CanwuError::new(
552 ErrorCode::UndeclaredKnowledgeWrite,
553 "knowledge publication is allowed only in phases 4 and 13",
554 ));
555 }
556
557 match phase {
558 BoundaryPhase::PerceptionAndAttentionRefresh => {
559 if !phase_directives.is_empty() {
560 return Err(CanwuError::new(
561 ErrorCode::InvalidBoundary,
562 "phase 4 accepts knowledge publications but no ordinary directives",
563 ));
564 }
565 self.stage_knowledge_publications(
566 phase,
567 knowledge_directives,
568 &mut visible_knowledge_overlay,
569 &mut pending_knowledge_changes,
570 &mut knowledge_correlations,
571 )?;
572 }
573 BoundaryPhase::ReservationAndAllocation => {
574 let result = allocate_reservations(
575 std::mem::take(&mut offers),
576 std::mem::take(&mut requests),
577 )?;
578 allocations = result.by_reservation;
579 allocation_records = result.records;
580 reservation_offer_records = result.offers;
581 reservation_request_records = result.requests;
582 }
583 BoundaryPhase::DomainDeltaProposal => {
584 let record_context = BoundaryRecordOverlayContext {
585 current: &boundary_snapshot,
586 now: boundary_time,
587 scheduled_actions: &self.state.scheduler.actions,
588 run_configuration: &self.state.metadata.run_configuration,
589 schemas: &record_schemas,
590 };
591 extend_boundary_record_candidate_overlay(
592 &record_context,
593 &mut candidate_record_overlay,
594 &phase_directives,
595 )?;
596 extend_boundary_candidate_overlay(
597 &boundary_snapshot,
598 &candidate_record_overlay,
599 &mut candidate_overlay,
600 &phase_directives,
601 )?;
602 extend_boundary_record_overlay(
603 &record_context,
604 &mut visible_record_overlay,
605 &phase_directives,
606 )?;
607 extend_boundary_overlay(
608 &boundary_snapshot,
609 &visible_record_overlay,
610 &mut visible_overlay,
611 &phase_directives,
612 )?;
613 ordinary.extend(phase_directives);
614 }
615 BoundaryPhase::HistoricalCandidateEvaluation => {
616 let record_context = BoundaryRecordOverlayContext {
617 current: &self.state.current,
618 now: self.state.scheduler.now,
619 scheduled_actions: &self.state.scheduler.actions,
620 run_configuration: &self.state.metadata.run_configuration,
621 schemas: &record_schemas,
622 };
623 extend_boundary_record_overlay(
624 &record_context,
625 &mut visible_record_overlay,
626 &phase_directives,
627 )?;
628 extend_boundary_overlay(
629 &self.state.current,
630 &visible_record_overlay,
631 &mut visible_overlay,
632 &phase_directives,
633 )?;
634 transitions.extend(phase_directives);
635 }
636 BoundaryPhase::StrategicAggregation
637 | BoundaryPhase::PerspectiveAndReportMaterialization => {
638 if phase == BoundaryPhase::PerspectiveAndReportMaterialization {
639 self.stage_knowledge_publications(
640 phase,
641 knowledge_directives,
642 &mut visible_knowledge_overlay,
643 &mut pending_knowledge_changes,
644 &mut knowledge_correlations,
645 )?;
646 }
647 let (same_boundary, next_boundary) =
648 partition_boundary_visibility(phase_directives);
649 self.apply_boundary_stage(
650 boundary_id,
651 correlation_id,
652 same_boundary,
653 &mut evidence,
654 )?;
655 deferred.extend(next_boundary);
656 }
657 _ if !phase_directives.is_empty() => {
658 return Err(CanwuError::new(
659 ErrorCode::InvalidBoundary,
660 format!("boundary phase {phase:?} cannot produce state directives"),
661 ));
662 }
663 _ => {}
664 }
665 }
666
667 self.apply_boundary_stage(boundary_id, correlation_id, deferred, &mut evidence)?;
668 self.commit_knowledge_publications(
669 boundary_id,
670 correlation_id,
671 &pending_knowledge_changes,
672 &mut evidence.emissions,
673 )?;
674 let PendingBoundaryEvidence {
675 changes,
676 record_changes,
677 emissions,
678 mut generated_ingress,
679 random_decisions,
680 mut person_availability_changes,
681 created_persons,
682 } = evidence;
683 self.invalidate_commitments(CommitmentDomains::RANDOM_STREAMS);
684 self.state.current.random_streams = random_overlay;
685 let mut random_outcomes = BTreeMap::new();
686 for pending in &random_decisions {
687 let ticket = self
688 .state
689 .current
690 .decisions
691 .ticket(pending.resolution.ticket_id)
692 .ok_or_else(|| {
693 CanwuError::new(
694 ErrorCode::InvalidDecision,
695 "random decision ticket disappeared before boundary commit",
696 )
697 })?;
698 let option_id = random_resolution_selection(ticket, &pending.resolution)?;
699 let previous = random_outcomes.insert(
700 (
701 pending.resolution.sample.stream.clone(),
702 pending.resolution.sample.address.clone(),
703 ),
704 RandomDrawOutcome::DecisionSelection {
705 ticket_id: ticket.id,
706 ticket_version: ticket.version,
707 option_id,
708 },
709 );
710 if previous.is_some() {
711 return Err(CanwuError::new(
712 ErrorCode::InvalidRandomDraw,
713 "one random draw cannot resolve multiple decisions",
714 ));
715 }
716 }
717 let committed_random_draws = self.append_boundary_random_draws(
718 boundary_id,
719 correlation_id,
720 pending_random_draws,
721 random_outcomes,
722 )?;
723 self.materialize_boundary_random_decisions(
724 boundary_id,
725 &random_decisions,
726 &committed_random_draws,
727 &mut generated_ingress,
728 )?;
729 self.cancel_unavailable_person_tickets(&mut person_availability_changes)?;
730 let transition_evidence = transition_ledger.finish();
731 if !transition_evidence.registered.is_empty() || !transition_evidence.audits.is_empty() {
732 self.invalidate_commitments(CommitmentDomains::SCHEDULER);
733 self.state.scheduler.transition_manifests = transition_evidence.pending;
734 }
735 let random_draws = committed_random_draws
736 .iter()
737 .map(|draw| draw.id)
738 .collect::<Vec<_>>();
739 self.state.metadata.plugin_registration_closed = true;
740 let state_hash = self.compute_boundary_state_hash_for(state_hash_format)?;
741 let previous_hash = self
742 .state
743 .evidence
744 .boundary_head_hash()
745 .map_or_else(|| GENESIS_BOUNDARY_HASH.to_owned(), str::to_owned);
746 let previous_maintenance_root = self
747 .state
748 .evidence
749 .boundaries
750 .last()
751 .and_then(|boundary| boundary.maintenance_terminal_root.as_deref());
752 let maintenance_terminal_root =
753 if previous_maintenance_root.is_some() || !maintenance_changes.is_empty() {
754 Some(canonical_hash(
755 "canwu.maintenance.terminal-root.v1",
756 &(
757 previous_maintenance_root.unwrap_or(GENESIS_BOUNDARY_HASH),
758 &maintenance_changes,
759 ),
760 )?)
761 } else {
762 None
763 };
764 let mut record = BoundaryRecord {
765 id: boundary_id,
766 at: request.at,
767 correlation_id,
768 cadences: request.cadences,
769 admitted_attempts,
770 admitted_commands,
771 admitted_ingress,
772 generated_ingress: generated_ingress.clone(),
773 admitted_events,
774 reservation_offers: reservation_offer_records,
775 reservation_requests: reservation_request_records,
776 allocations: allocation_records.clone(),
777 random_draws: random_draws.clone(),
778 changes: changes.clone(),
779 record_changes: record_changes.clone(),
780 knowledge_changes: pending_knowledge_changes.clone(),
781 maintenance_changes,
782 maintenance_terminal_root,
783 person_availability_changes,
784 created_persons: created_persons.clone(),
785 evaluation_traces,
786 transition_manifests: transition_evidence.registered,
787 transition_audits: transition_evidence.audits.clone(),
788 emissions: emissions.clone(),
789 state_hash: Some(state_hash),
790 previous_hash,
791 hash: String::new(),
792 };
793 record.hash = compute_boundary_hash(&record)?;
794 let boundary_hash = record.hash.clone();
795 self.state.evidence.boundaries.push(record);
796 self.state.counters.admitted_attempt_count = admitted_attempt_count;
797 self.state.counters.admitted_command_count = admitted_command_count;
798 self.state.counters.admitted_event_count = admitted_event_count;
799 self.advance_state_revision()?;
800 self.refresh_checkpoint_hash()?;
801 Ok(BoundaryReceipt {
802 boundary_id,
803 settled_at: request.at,
804 emitted_events: emissions
805 .into_iter()
806 .map(|emission| emission.event)
807 .collect(),
808 generated_ingress: generated_ingress
809 .into_iter()
810 .map(|generation| generation.ingress)
811 .collect(),
812 random_draws,
813 boundary_hash,
814 change_count: changes.len(),
815 record_change_count: record_changes.len(),
816 knowledge_batch_count: pending_knowledge_changes.len(),
817 knowledge_record_count: pending_knowledge_changes
818 .iter()
819 .map(|change| change.records.len())
820 .sum(),
821 allocations: allocation_records,
822 created_persons: created_persons
823 .iter()
824 .map(super::CreatedPerson::from)
825 .collect(),
826 transition_audits: transition_evidence.audits,
827 })
828 }
829
830 fn check_transition_post_versions(
837 &self,
838 ledger: &mut super::transitions::BoundaryTransitionLedger,
839 record_schemas: &records::DomainRecordSchemas,
840 deferred: &[StagedBoundaryDirective],
841 ) -> Result<(), CanwuError> {
842 let mut candidate = BTreeMap::new();
843 extend_boundary_record_candidate_overlay(
844 &BoundaryRecordOverlayContext {
845 current: &self.state.current,
846 now: self.state.scheduler.now,
847 scheduled_actions: &self.state.scheduler.actions,
848 run_configuration: &self.state.metadata.run_configuration,
849 schemas: record_schemas,
850 },
851 &mut candidate,
852 deferred,
853 )?;
854 ledger.check_expected_post(&|record| {
855 candidate
856 .get(record)
857 .or_else(|| self.state.current.domain_records.get(record))
858 .map(|record| record.version)
859 })
860 }
861
862 fn apply_boundary_stage(
863 &mut self,
864 boundary_id: BoundaryId,
865 correlation_id: u64,
866 directives: Vec<StagedBoundaryDirective>,
867 evidence: &mut PendingBoundaryEvidence,
868 ) -> Result<(), CanwuError> {
869 let random_decision_count = directives
870 .iter()
871 .filter(|staged| {
872 matches!(
873 &staged.directive,
874 BoundaryDirective::ResolveDecisionRandomly { .. }
875 )
876 })
877 .count();
878 if evidence
879 .random_decisions
880 .len()
881 .checked_add(random_decision_count)
882 .is_none_or(|count| count > 1)
883 {
884 return Err(CanwuError::new(
885 ErrorCode::InvalidDecision,
886 "one boundary may generate at most one random decision resolution",
887 ));
888 }
889 let changes = &mut evidence.changes;
890 let record_changes = &mut evidence.record_changes;
891 let emissions = &mut evidence.emissions;
892 let generated_ingress = &mut evidence.generated_ingress;
893 let mutation_requests: Vec<_> = directives
894 .iter()
895 .filter_map(|staged| match &staged.directive {
896 BoundaryDirective::MutateRecord { mutation, summary } => {
897 Some(records::DomainMutationRequest {
898 plugin: &staged.plugin,
899 system: &staged.system,
900 visibility: staged.visibility,
901 mutation,
902 summary,
903 })
904 }
905 BoundaryDirective::SetComponent { .. }
906 | BoundaryDirective::Emit { .. }
907 | BoundaryDirective::ScheduleIngress { .. }
908 | BoundaryDirective::SchedulePluginIngress { .. }
909 | BoundaryDirective::ResolveDecisionRandomly { .. }
910 | BoundaryDirective::PublishKnowledge { .. }
911 | BoundaryDirective::SetPersonAvailability { .. }
912 | BoundaryDirective::CreatePerson { .. }
913 | BoundaryDirective::CancelPluginIngress { .. }
914 | BoundaryDirective::RecordEvaluationTrace { .. }
915 | BoundaryDirective::RegisterTransitionManifest { .. }
916 | BoundaryDirective::StageTransitionWrite { .. } => None,
917 })
918 .collect();
919 let mut stage_record_changes = BTreeMap::new();
920 if !mutation_requests.is_empty() {
921 let (next_records, applied) = records::apply_mutation_bundle_cow(
922 &self.state.current.domain_records,
923 &self.plugins.record_schemas,
924 self.state.scheduler.now,
925 &|entity| runtime_entity_exists(&self.state, entity),
926 mutation_requests,
927 )?;
928 let first_index = record_changes.len();
929 for (offset, change) in applied.iter().enumerate() {
930 let index = first_index.checked_add(offset).ok_or_else(|| {
931 CanwuError::new(
932 ErrorCode::IdentifierExhausted,
933 "boundary record-change index exceeds the persistent identifier space",
934 )
935 })?;
936 let index = u64::try_from(index).map_err(|_| {
937 CanwuError::new(
938 ErrorCode::IdentifierExhausted,
939 "boundary record-change index exceeds the persistent identifier space",
940 )
941 })?;
942 index_current_domain_record_version(&mut self.state, boundary_id, index, change);
943 stage_record_changes
944 .insert(change.current.reference.clone(), (index, change.clone()));
945 }
946 self.invalidate_commitments(CommitmentDomains::DOMAIN_RECORDS);
947 self.state.current.domain_records = next_records;
948 record_changes.extend(applied);
949 }
950
951 for staged in &directives {
952 let unavailable = match &staged.directive {
953 BoundaryDirective::SetComponent { entity, .. } => {
954 (!runtime_entity_exists(&self.state, entity)).then_some(entity)
955 }
956 BoundaryDirective::Emit { affected, .. } => affected
957 .iter()
958 .find(|entity| !runtime_entity_exists(&self.state, entity)),
959 BoundaryDirective::ScheduleIngress { affected, .. } => affected
960 .iter()
961 .find(|entity| !runtime_entity_identity_exists(&self.state, entity)),
962 BoundaryDirective::SchedulePluginIngress { affected, .. } => affected
963 .iter()
964 .find(|entity| !runtime_entity_identity_exists(&self.state, entity)),
965 BoundaryDirective::MutateRecord { .. }
966 | BoundaryDirective::ResolveDecisionRandomly { .. }
967 | BoundaryDirective::PublishKnowledge { .. }
968 | BoundaryDirective::SetPersonAvailability { .. }
969 | BoundaryDirective::CreatePerson { .. }
970 | BoundaryDirective::CancelPluginIngress { .. }
971 | BoundaryDirective::RecordEvaluationTrace { .. }
972 | BoundaryDirective::RegisterTransitionManifest { .. }
973 | BoundaryDirective::StageTransitionWrite { .. } => None,
974 };
975 if let Some(entity) = unavailable {
976 return Err(CanwuError::new(
977 ErrorCode::EntityNotFound,
978 format!(
979 "boundary stage {}.{} references unavailable entity {entity}",
980 staged.plugin, staged.system
981 ),
982 )
983 .with_entity(entity.clone()));
984 }
985 }
986
987 for staged in directives {
988 match staged.directive {
989 BoundaryDirective::SetComponent {
990 state,
991 entity,
992 component,
993 value,
994 summary,
995 } => {
996 let key = component_key(&staged.plugin, &state, &entity, &component);
997 self.invalidate_commitments(CommitmentDomains::PLUGIN_COMPONENTS);
998 let previous = self
999 .state
1000 .current
1001 .plugin_components
1002 .get(&key)
1003 .map(|record| record.value.clone());
1004 self.state.current.plugin_components.insert(
1005 key,
1006 PluginComponentRecord {
1007 plugin: staged.plugin.clone(),
1008 state: state.clone(),
1009 entity: entity.clone(),
1010 component: component.clone(),
1011 value: value.clone(),
1012 },
1013 );
1014 let change_index = u64::try_from(changes.len()).map_err(|_| {
1015 CanwuError::new(
1016 ErrorCode::IdentifierExhausted,
1017 "boundary change index exceeds the persistent identifier space",
1018 )
1019 })?;
1020 changes.push(BoundaryChange {
1021 plugin: staged.plugin.clone(),
1022 system: staged.system.clone(),
1023 state,
1024 entity: entity.clone(),
1025 component: component.clone(),
1026 previous,
1027 value,
1028 visibility: staged.visibility,
1029 summary: summary.clone(),
1030 });
1031 let event = self.append_event(
1032 EventKind::plugin(staged.plugin.clone(), format!("{component}_changed")),
1033 vec![entity],
1034 summary,
1035 Some(CauseRef::Boundary(boundary_id)),
1036 correlation_id,
1037 )?;
1038 emissions.push(BoundaryEmission {
1039 plugin: staged.plugin,
1040 system: staged.system,
1041 event: event.id,
1042 kind: BoundaryEmissionKind::Change { change_index },
1043 });
1044 }
1045 BoundaryDirective::MutateRecord { mutation, .. } => {
1046 let Some((change_index, change)) = stage_record_changes.get(mutation.target())
1047 else {
1048 return Err(CanwuError::new(
1049 ErrorCode::InvalidBoundary,
1050 "record mutation is missing its committed change evidence",
1051 ));
1052 };
1053 let event = self.append_event(
1054 EventKind::plugin(staged.plugin.clone(), change.operation.event_type()),
1055 record_change_affected_entities(change),
1056 change.summary.clone(),
1057 Some(CauseRef::Boundary(boundary_id)),
1058 correlation_id,
1059 )?;
1060 emissions.push(BoundaryEmission {
1061 plugin: staged.plugin,
1062 system: staged.system,
1063 event: event.id,
1064 kind: BoundaryEmissionKind::RecordChange {
1065 change_index: *change_index,
1066 },
1067 });
1068 }
1069 BoundaryDirective::Emit {
1070 event_type,
1071 summary,
1072 affected,
1073 } => {
1074 let event = self.append_event(
1075 EventKind::plugin(staged.plugin.clone(), event_type),
1076 affected,
1077 summary,
1078 Some(CauseRef::Boundary(boundary_id)),
1079 correlation_id,
1080 )?;
1081 emissions.push(BoundaryEmission {
1082 plugin: staged.plugin,
1083 system: staged.system,
1084 event: event.id,
1085 kind: BoundaryEmissionKind::Explicit,
1086 });
1087 }
1088 BoundaryDirective::PublishKnowledge { .. } => {
1089 return Err(CanwuError::new(
1090 ErrorCode::InvalidBoundary,
1091 "knowledge publication execution is not enabled in this runtime slice",
1092 ));
1093 }
1094 BoundaryDirective::ScheduleIngress {
1095 after,
1096 packet_type,
1097 priority,
1098 payload,
1099 mut affected,
1100 } => {
1101 self.ensure_canonical_ingress_can_start()?;
1102 let descriptor = self
1103 .plugins
1104 .ingress
1105 .get(&(staged.plugin.clone(), packet_type.clone()))
1106 .ok_or_else(|| {
1107 CanwuError::new(
1108 ErrorCode::InvalidPayload,
1109 format!(
1110 "boundary system {}.{} scheduled undeclared ingress type {packet_type}",
1111 staged.plugin, staged.system
1112 ),
1113 )
1114 })?
1115 .clone();
1116 descriptor.payload_schema.validate(&payload)?;
1117 affected.sort();
1118 affected.dedup();
1119 let due_at = self.state.scheduler.now.checked_add(after).ok_or_else(|| {
1120 CanwuError::new(
1121 ErrorCode::InvalidDuration,
1122 "boundary-generated ingress exceeds the supported time range",
1123 )
1124 })?;
1125 let receipt = self.append_ingress(
1126 due_at,
1127 descriptor.class,
1128 priority,
1129 IngressPayload::Plugin {
1130 plugin: staged.plugin.clone(),
1131 packet_type,
1132 payload,
1133 affected_entities: affected,
1134 archive_retention: Vec::new(),
1135 },
1136 Some(CauseRef::Boundary(boundary_id)),
1137 true,
1138 )?;
1139 generated_ingress.push(BoundaryIngressGeneration {
1140 ingress: receipt.ingress_id,
1141 plugin: staged.plugin,
1142 system: staged.system,
1143 phase: staged.phase,
1144 visibility: staged.visibility,
1145 });
1146 }
1147 BoundaryDirective::SchedulePluginIngress {
1148 target_plugin,
1149 after,
1150 packet_type,
1151 priority,
1152 payload,
1153 mut affected,
1154 } => {
1155 self.ensure_canonical_ingress_can_start()?;
1156 let descriptor = self
1157 .plugins
1158 .ingress
1159 .get(&(target_plugin.clone(), packet_type.clone()))
1160 .ok_or_else(|| {
1161 CanwuError::new(
1162 ErrorCode::InvalidPayload,
1163 format!(
1164 "boundary system {}.{} scheduled undeclared target ingress {}.{packet_type}",
1165 staged.plugin, staged.system, target_plugin
1166 ),
1167 )
1168 })?
1169 .clone();
1170 descriptor.payload_schema.validate(&payload)?;
1171 affected.sort();
1172 affected.dedup();
1173 let due_at = self.state.scheduler.now.checked_add(after).ok_or_else(|| {
1174 CanwuError::new(
1175 ErrorCode::InvalidDuration,
1176 "boundary-generated cross-plugin ingress exceeds the supported time range",
1177 )
1178 })?;
1179 let receipt = self.append_ingress(
1180 due_at,
1181 descriptor.class,
1182 priority,
1183 IngressPayload::Plugin {
1184 plugin: target_plugin,
1185 packet_type,
1186 payload,
1187 affected_entities: affected,
1188 archive_retention: Vec::new(),
1189 },
1190 Some(CauseRef::Boundary(boundary_id)),
1191 true,
1192 )?;
1193 generated_ingress.push(BoundaryIngressGeneration {
1194 ingress: receipt.ingress_id,
1195 plugin: staged.plugin,
1196 system: staged.system,
1197 phase: staged.phase,
1198 visibility: staged.visibility,
1199 });
1200 }
1201 BoundaryDirective::CancelPluginIngress { ingress_id, reason } => {
1202 self.ensure_canonical_ingress_can_start()?;
1203 let cancelled_in_this_boundary = generated_ingress.iter().any(|generation| {
1204 self.state
1205 .evidence
1206 .retained_ingress(generation.ingress)
1207 .is_some_and(|record| {
1208 matches!(
1209 record.payload,
1210 IngressPayload::PluginCancellation { cancelled, .. }
1211 if cancelled == ingress_id
1212 )
1213 })
1214 });
1215 if cancelled_in_this_boundary {
1216 return Err(CanwuError::new(
1217 ErrorCode::InvalidBoundary,
1218 format!(
1219 "multiple boundary proposals cancel ingress {ingress_id}; boundary system {}.{} must not repeat a cancellation",
1220 staged.plugin, staged.system
1221 ),
1222 ));
1223 }
1224 let target = self.plugin_ingress_cancellation_target(
1225 ingress_id,
1226 IngressCancellationAuthority::BoundarySystem,
1227 PluginIngressCancellationProof {
1228 permit: None,
1229 replay: false,
1230 boundary_plugin: Some(&staged.plugin),
1231 current_generations: generated_ingress,
1232 },
1233 &reason,
1234 )?;
1235 let receipt = self.append_plugin_ingress_cancellation(
1236 target,
1237 IngressCancellationAuthority::BoundarySystem,
1238 reason,
1239 Some(CauseRef::Boundary(boundary_id)),
1240 true,
1241 )?;
1242 generated_ingress.push(BoundaryIngressGeneration {
1243 ingress: receipt.ingress_id,
1244 plugin: staged.plugin,
1245 system: staged.system,
1246 phase: staged.phase,
1247 visibility: staged.visibility,
1248 });
1249 }
1250 BoundaryDirective::ResolveDecisionRandomly { resolution } => {
1251 evidence
1252 .random_decisions
1253 .push(PendingRandomDecisionResolution {
1254 plugin: staged.plugin,
1255 system: staged.system,
1256 phase: staged.phase,
1257 visibility: staged.visibility,
1258 resolution,
1259 });
1260 }
1261 BoundaryDirective::SetPersonAvailability {
1262 person,
1263 availability,
1264 summary,
1265 } => {
1266 let change = self.apply_person_availability(
1267 (
1268 &staged.plugin,
1269 &staged.system,
1270 staged.phase,
1271 staged.visibility,
1272 ),
1273 person,
1274 availability,
1275 summary,
1276 )?;
1277 evidence.person_availability_changes.push(change);
1278 }
1279 BoundaryDirective::CreatePerson {
1280 draft,
1281 correlation,
1282 summary,
1283 } => {
1284 let creation = self.apply_person_creation(
1285 &staged.plugin,
1286 &staged.system,
1287 draft,
1288 correlation,
1289 summary,
1290 )?;
1291 evidence.created_persons.push(creation);
1292 }
1293 BoundaryDirective::RecordEvaluationTrace { .. } => {
1294 return Err(CanwuError::new(
1295 ErrorCode::InvalidBoundary,
1296 "evaluation traces are boundary evidence and never reach a commit stage",
1297 ));
1298 }
1299 BoundaryDirective::RegisterTransitionManifest { .. }
1300 | BoundaryDirective::StageTransitionWrite { .. } => {
1301 return Err(transition_directive_not_admitted(
1302 &staged.plugin,
1303 &staged.system,
1304 ));
1305 }
1306 }
1307 }
1308 validate_runtime_domain_dependents(&self.state)?;
1309 Ok(())
1310 }
1311
1312 fn materialize_boundary_random_decisions(
1313 &mut self,
1314 boundary_id: BoundaryId,
1315 pending: &[PendingRandomDecisionResolution],
1316 committed_draws: &[CommittedBoundaryRandomDraw],
1317 generated_ingress: &mut Vec<BoundaryIngressGeneration>,
1318 ) -> Result<(), CanwuError> {
1319 let expected_revision = self.revision().checked_add(1).ok_or_else(|| {
1320 CanwuError::new(
1321 ErrorCode::IdentifierExhausted,
1322 "random decision target revision is exhausted",
1323 )
1324 })?;
1325 let draws = committed_draws
1326 .iter()
1327 .map(|draw| ((draw.stream.clone(), draw.address.clone()), draw.id))
1328 .collect::<BTreeMap<_, _>>();
1329 for pending in pending {
1330 let resolution = &pending.resolution;
1331 let draw_id = draws
1332 .get(&(
1333 resolution.sample.stream.clone(),
1334 resolution.sample.address.clone(),
1335 ))
1336 .copied()
1337 .ok_or_else(|| {
1338 CanwuError::new(
1339 ErrorCode::InvalidRandomDraw,
1340 "random decision draw was not committed by its source boundary",
1341 )
1342 })?;
1343 let ticket = self
1344 .state
1345 .current
1346 .decisions
1347 .ticket(resolution.ticket_id)
1348 .cloned()
1349 .ok_or_else(|| {
1350 CanwuError::new(
1351 ErrorCode::InvalidDecision,
1352 "random decision ticket disappeared before ingress generation",
1353 )
1354 })?;
1355 let controller = self
1356 .state
1357 .current
1358 .decisions
1359 .controller(&resolution.controller_id)
1360 .cloned()
1361 .ok_or_else(|| {
1362 CanwuError::new(
1363 ErrorCode::InvalidDecision,
1364 "random decision controller disappeared before ingress generation",
1365 )
1366 })?;
1367 let option_id = random_resolution_selection(&ticket, resolution)?;
1368 let option = ticket.option(&option_id).ok_or_else(|| {
1369 CanwuError::new(
1370 ErrorCode::InvalidDecision,
1371 "random decision selected an unavailable option",
1372 )
1373 })?;
1374 let due_at = self.state.scheduler.now;
1375 let command = match &option.action {
1376 DecisionAction::Command { command } => {
1377 let request_id = resolution.command_request_id.ok_or_else(|| {
1378 CanwuError::new(
1379 ErrorCode::InvalidDecision,
1380 "random decision command option lacks a command request ID",
1381 )
1382 })?;
1383 let command: Command =
1384 serde_json::from_value(command.clone()).map_err(|error| {
1385 CanwuError::new(
1386 ErrorCode::InvalidDecision,
1387 format!("decision option contains an invalid command: {error}"),
1388 )
1389 })?;
1390 Some(CommandRequest::new(
1391 request_id,
1392 expected_revision,
1393 CommandEnvelope::new(
1394 super::decision::controller_issuer(&controller),
1395 command,
1396 )
1397 .with_authority(super::decision::controller_authority(&controller))
1398 .at_time(due_at),
1399 ))
1400 }
1401 DecisionAction::None => None,
1402 };
1403 let random = Some(DecisionRandomEvidence {
1404 draw_id,
1405 value: resolution.sample.value,
1406 upper_exclusive: resolution.sample.upper_exclusive,
1407 option_weights: resolution.option_weights.clone(),
1408 });
1409 let outcome = DecisionOutcome::Selected {
1410 option_id: option_id.clone(),
1411 };
1412 let decision = match &resolution.tie_break {
1413 None => PolicyDecision {
1414 outcome,
1415 summary: random_policy_summary(&option_id),
1416 evaluations: Vec::new(),
1417 external: None,
1418 random,
1419 stage: None,
1420 fired_guards: Vec::new(),
1421 },
1422 Some(pending) => PolicyDecision {
1423 outcome,
1424 summary: random_tie_break_summary(&option_id),
1425 evaluations: pending.evaluations.clone(),
1426 external: None,
1427 random,
1428 stage: Some(DecisionStage::Random),
1429 fired_guards: pending.fired_guards.clone(),
1430 },
1431 };
1432 let mutation = DecisionMutation::Resolve {
1433 ticket_id: ticket.id,
1434 expected_version: ticket.version,
1435 controller_id: controller.id.clone(),
1436 policy: controller.policy.clone(),
1437 decision,
1438 command_request_id: resolution.command_request_id,
1439 };
1440 let mut request = DecisionIngressRequest::new(
1441 resolution.decision_request_id,
1442 expected_revision,
1443 mutation,
1444 );
1445 if let Some(command) = command {
1446 request = request.with_command(command);
1447 }
1448 let receipt = self.append_boundary_decision_ingress(
1449 boundary_id,
1450 due_at,
1451 resolution.priority,
1452 request,
1453 )?;
1454 generated_ingress.push(BoundaryIngressGeneration {
1455 ingress: receipt.ingress_id,
1456 plugin: pending.plugin.clone(),
1457 system: pending.system.clone(),
1458 phase: pending.phase,
1459 visibility: pending.visibility,
1460 });
1461 }
1462 Ok(())
1463 }
1464
1465 fn stage_knowledge_publications(
1466 &mut self,
1467 phase: BoundaryPhase,
1468 directives: Vec<StagedBoundaryDirective>,
1469 visible_overlay: &mut BTreeMap<
1470 KnowledgeHolderRef,
1471 BTreeMap<KnowledgeRecordId, KnowledgeRecord>,
1472 >,
1473 pending: &mut Vec<BoundaryKnowledgeChange>,
1474 correlations: &mut BTreeSet<(String, String, String)>,
1475 ) -> Result<(), CanwuError> {
1476 let new_record_count = directives
1477 .iter()
1478 .map(|staged| match &staged.directive {
1479 BoundaryDirective::PublishKnowledge { records, .. } => records.len(),
1480 _ => 0,
1481 })
1482 .sum::<usize>();
1483 let total_records = pending
1484 .iter()
1485 .map(|change| change.records.len())
1486 .sum::<usize>()
1487 .checked_add(new_record_count)
1488 .ok_or_else(|| {
1489 CanwuError::new(
1490 ErrorCode::KnowledgeLimitExceeded,
1491 "boundary knowledge record count exceeds platform range",
1492 )
1493 })?;
1494 if total_records > crate::KnowledgeLimitsV1::CURRENT.records_per_boundary {
1495 return Err(CanwuError::new(
1496 ErrorCode::KnowledgeLimitExceeded,
1497 "boundary knowledge record limit exceeded",
1498 ));
1499 }
1500 for staged in directives {
1501 let BoundaryDirective::PublishKnowledge {
1502 holder,
1503 visibility,
1504 producer_correlation,
1505 records: drafts,
1506 summary,
1507 } = staged.directive
1508 else {
1509 return Err(CanwuError::new(
1510 ErrorCode::InvalidBoundary,
1511 "knowledge stage received an ordinary directive",
1512 ));
1513 };
1514 if let Some(value) = &producer_correlation
1515 && !correlations.insert((
1516 staged.plugin.clone(),
1517 staged.system.clone(),
1518 value.clone(),
1519 ))
1520 {
1521 return Err(CanwuError::new(
1522 ErrorCode::InvalidKnowledgeRecord,
1523 "producer correlation is duplicated within one system and boundary",
1524 ));
1525 }
1526 let mut records = Vec::with_capacity(drafts.len());
1527 for draft in drafts {
1528 let (id, next_id) = claim_counter(
1529 self.state.counters.next_knowledge_record_id,
1530 "knowledge record ID",
1531 )?;
1532 self.state.counters.next_knowledge_record_id = next_id;
1533 let record = KnowledgeRecord {
1534 id: KnowledgeRecordId::new(id),
1535 holder: holder.clone(),
1536 schema: draft.schema,
1537 subjects: draft.subjects,
1538 payload: draft.payload,
1539 as_of: draft.as_of,
1540 learned_at: self.state.scheduler.now,
1541 confidence_per_mille: draft.confidence_per_mille,
1542 origin: draft.origin,
1543 supersedes: draft.supersedes,
1544 contradicts: draft.contradicts,
1545 };
1546 if visibility == StateVisibility::SameBoundary {
1547 visible_overlay
1548 .entry(holder.clone())
1549 .or_default()
1550 .insert(record.id, record.clone());
1551 }
1552 records.push(record);
1553 }
1554 pending.push(BoundaryKnowledgeChange {
1555 plugin: staged.plugin,
1556 system: staged.system,
1557 phase,
1558 holder,
1559 producer_correlation,
1560 records,
1561 visibility,
1562 summary,
1563 });
1564 }
1565 Ok(())
1566 }
1567
1568 fn commit_knowledge_publications(
1569 &mut self,
1570 boundary_id: BoundaryId,
1571 correlation_id: u64,
1572 changes: &[BoundaryKnowledgeChange],
1573 emissions: &mut Vec<BoundaryEmission>,
1574 ) -> Result<(), CanwuError> {
1575 if changes.is_empty() {
1576 return Ok(());
1577 }
1578 let mut ledger = self.state.current.knowledge.records.clone();
1579 for change in changes {
1580 let holder = ledger.entry(change.holder.clone()).or_default();
1581 for record in &change.records {
1582 if holder.insert(record.id, record.clone()).is_some() {
1583 return Err(CanwuError::new(
1584 ErrorCode::InvalidKnowledgeRecord,
1585 "knowledge publication attempted to reuse a global record ID",
1586 ));
1587 }
1588 }
1589 }
1590 self.state.current.knowledge.records = ledger;
1591 self.invalidate_commitments(CommitmentDomains::KNOWLEDGE);
1592 for (index, change) in changes.iter().enumerate() {
1593 let record_count = u32::try_from(change.records.len()).map_err(|_| {
1594 CanwuError::new(
1595 ErrorCode::KnowledgeLimitExceeded,
1596 "knowledge publication event count exceeds u32",
1597 )
1598 })?;
1599 let affected = match &change.holder {
1600 KnowledgeHolderRef::Person(person) => vec![EntityRef::Person(*person)],
1601 KnowledgeHolderRef::Entity(entity) => vec![entity.clone()],
1602 };
1603 let event = self.append_event(
1604 KnowledgePublished {
1605 holder: change.holder.clone(),
1606 record_count,
1607 }
1608 .into_kind(),
1609 affected,
1610 change.summary.clone(),
1611 Some(CauseRef::Boundary(boundary_id)),
1612 correlation_id,
1613 )?;
1614 emissions.push(BoundaryEmission {
1615 plugin: change.plugin.clone(),
1616 system: change.system.clone(),
1617 event: event.id,
1618 kind: BoundaryEmissionKind::KnowledgeChange {
1619 change_index: u64::try_from(index).map_err(|_| {
1620 CanwuError::new(
1621 ErrorCode::IdentifierExhausted,
1622 "knowledge change index exceeds identifier space",
1623 )
1624 })?,
1625 },
1626 });
1627 }
1628 Ok(())
1629 }
1630
1631 pub(super) fn apply_directives(
1632 &mut self,
1633 plugin: &str,
1634 directives: Vec<SystemDirective>,
1635 allowed_writes: &[StateKey],
1636 cause: &CauseRef,
1637 correlation_id: u64,
1638 ) -> Result<(), CanwuError> {
1639 for directive in directives {
1640 match directive {
1641 SystemDirective::SetComponent {
1642 state,
1643 entity,
1644 component,
1645 value,
1646 summary,
1647 } => {
1648 let key = component_key(plugin, &state, &entity, &component);
1649 self.invalidate_commitments(CommitmentDomains::PLUGIN_COMPONENTS);
1650 self.state.current.plugin_components.insert(
1651 key,
1652 PluginComponentRecord {
1653 plugin: plugin.to_owned(),
1654 state,
1655 entity: entity.clone(),
1656 component: component.clone(),
1657 value,
1658 },
1659 );
1660 self.emit(
1661 EventKind::plugin(plugin, format!("{component}_changed")),
1662 vec![entity],
1663 summary,
1664 Some(cause.clone()),
1665 correlation_id,
1666 )?;
1667 }
1668 SystemDirective::Emit {
1669 event_type,
1670 summary,
1671 affected,
1672 } => {
1673 self.emit(
1674 EventKind::plugin(plugin, event_type),
1675 affected,
1676 summary,
1677 Some(cause.clone()),
1678 correlation_id,
1679 )?;
1680 }
1681 SystemDirective::Schedule { after, directive } => {
1682 let at = self.state.scheduler.now.checked_add(after).ok_or_else(|| {
1683 CanwuError::new(
1684 ErrorCode::InvalidDuration,
1685 "plugin scheduled time exceeds the supported range",
1686 )
1687 })?;
1688 self.schedule_at(
1689 at,
1690 ScheduledAction::PluginDirective {
1691 plugin: plugin.to_owned(),
1692 directive,
1693 allowed_writes: allowed_writes.to_vec(),
1694 cause: cause.clone(),
1695 correlation_id,
1696 },
1697 )?;
1698 }
1699 SystemDirective::EnqueuePluginIngress {
1700 after,
1701 packet_type,
1702 priority,
1703 payload,
1704 mut affected,
1705 } => {
1706 self.ensure_canonical_ingress_can_start()?;
1707 let descriptor = self
1708 .plugins
1709 .ingress
1710 .get(&(plugin.to_owned(), packet_type.clone()))
1711 .ok_or_else(|| {
1712 CanwuError::new(
1713 ErrorCode::InvalidPayload,
1714 format!(
1715 "plugin command scheduled unregistered ingress type {plugin}.{packet_type}"
1716 ),
1717 )
1718 })?
1719 .clone();
1720 descriptor.payload_schema.validate(&payload)?;
1721 affected.sort();
1722 affected.dedup();
1723 if affected
1724 .iter()
1725 .any(|entity| !runtime_entity_identity_exists(&self.state, entity))
1726 {
1727 return Err(CanwuError::new(
1728 ErrorCode::EntityNotFound,
1729 "plugin command ingress references an unknown entity identity",
1730 ));
1731 }
1732 let due_at = self.state.scheduler.now.checked_add(after).ok_or_else(|| {
1733 CanwuError::new(
1734 ErrorCode::InvalidDuration,
1735 "plugin command ingress exceeds the supported time range",
1736 )
1737 })?;
1738 self.append_ingress(
1739 due_at,
1740 descriptor.class,
1741 priority,
1742 IngressPayload::Plugin {
1743 plugin: plugin.to_owned(),
1744 packet_type,
1745 payload,
1746 affected_entities: affected,
1747 archive_retention: Vec::new(),
1748 },
1749 Some(cause.clone()),
1750 true,
1751 )?;
1752 }
1753 }
1754 }
1755 Ok(())
1756 }
1757}
1758
1759fn index_current_domain_record_version(
1760 state: &mut RuntimeState,
1761 boundary: BoundaryId,
1762 change_index: u64,
1763 change: &DomainRecordChange,
1764) {
1765 let established_at = state.scheduler.now;
1768 state.metadata.current_domain_record_versions.insert(
1769 change.current.reference.clone(),
1770 super::CurrentDomainRecordVersion {
1771 version: super::DomainRecordVersionRef {
1772 record: change.current.reference.clone(),
1773 version: change.current.version,
1774 established_by: super::DomainRecordVersionSource::BoundaryChange {
1775 boundary,
1776 change_index,
1777 },
1778 },
1779 established_at,
1780 },
1781 );
1782}
1783
1784struct PendingReservationOffer {
1785 plugin: String,
1786 system: String,
1787 offer: ReservationOffer,
1788}
1789
1790struct PendingReservationRequest {
1791 reservation: ReservationRef,
1792 request: ReservationRequest,
1793}
1794
1795struct ReservationAllocationResult {
1796 by_reservation: BTreeMap<ReservationRef, ReservationAllocation>,
1797 offers: Vec<ReservationOfferRecord>,
1798 requests: Vec<ReservationRequestRecord>,
1799 records: Vec<ReservationAllocation>,
1800}
1801
1802struct StagedBoundaryDirective {
1803 plugin: String,
1804 system: String,
1805 phase: BoundaryPhase,
1806 visibility: StateVisibility,
1807 directive: BoundaryDirective,
1808}
1809
1810#[derive(Default)]
1811struct PendingBoundaryEvidence {
1812 changes: Vec<BoundaryChange>,
1813 record_changes: Vec<DomainRecordChange>,
1814 emissions: Vec<BoundaryEmission>,
1815 generated_ingress: Vec<BoundaryIngressGeneration>,
1816 random_decisions: Vec<PendingRandomDecisionResolution>,
1817 person_availability_changes: Vec<super::BoundaryPersonAvailabilityChange>,
1818 created_persons: Vec<super::BoundaryPersonCreation>,
1819}
1820
1821struct PendingRandomDecisionResolution {
1822 plugin: String,
1823 system: String,
1824 phase: BoundaryPhase,
1825 visibility: StateVisibility,
1826 resolution: super::RandomDecisionResolution,
1827}
1828
1829fn proposal_evidence_refs(
1830 boundary: BoundaryId,
1831 pending: &PendingBoundaryEvidence,
1832) -> BTreeSet<EvidenceRef> {
1833 let mut values = BTreeSet::new();
1834 for (index, change) in pending.record_changes.iter().enumerate() {
1835 if let Ok(change_index) = u64::try_from(index) {
1836 values.insert(EvidenceRef::DomainRecordVersion(
1837 super::DomainRecordVersionRef {
1838 record: change.current.reference.clone(),
1839 version: change.current.version,
1840 established_by: DomainRecordVersionSource::BoundaryChange {
1841 boundary,
1842 change_index,
1843 },
1844 },
1845 ));
1846 }
1847 }
1848 values.extend(
1849 pending
1850 .emissions
1851 .iter()
1852 .map(|emission| EvidenceRef::Event(emission.event)),
1853 );
1854 values
1855}
1856
1857pub(super) struct PendingBoundaryRandomDraw {
1858 pub(super) plugin: String,
1859 pub(super) system: String,
1860 pub(super) draw: random::PendingRandomDraw,
1861}
1862
1863pub(super) struct CommittedBoundaryRandomDraw {
1864 pub(super) id: super::RandomDrawId,
1865 pub(super) stream: super::RandomStreamKey,
1866 pub(super) address: RandomDrawAddress,
1867}
1868
1869pub(super) fn boundary_system_due(
1870 contract: &BoundarySystemContract,
1871 cadences: &[SystemCadence],
1872 has_admitted_events: bool,
1873) -> bool {
1874 match contract.cadence {
1875 SystemCadence::EventDriven => has_admitted_events,
1876 _ => cadences.contains(&contract.cadence),
1877 }
1878}
1879
1880pub(super) fn boundary_has_event_ingress(record: &BoundaryRecord) -> bool {
1881 !record.admitted_events.is_empty() || !record.admitted_ingress.is_empty()
1882}
1883
1884#[allow(clippy::too_many_arguments)]
1885fn validate_boundary_proposal(
1886 plugin: &str,
1887 contract: &BoundarySystemContract,
1888 current: &RuntimeCurrentState,
1889 committed_availability: &BTreeMap<super::PersonId, super::PersonAvailability>,
1890 now: SimTime,
1891 runtime: &RuntimeState,
1892 boundary_id: BoundaryId,
1893 pending_evidence: &PendingBoundaryEvidence,
1894 plugins: &PluginRegistry,
1895 record_overlay: &BTreeMap<DomainRecordRef, DomainRecord>,
1896 knowledge_overlay: &BTreeMap<KnowledgeHolderRef, BTreeMap<KnowledgeRecordId, KnowledgeRecord>>,
1897 proposal: &BoundaryProposal,
1898 pending_random_draws: &[random::PendingRandomDraw],
1899) -> Result<(), CanwuError> {
1900 if contract.phase != BoundaryPhase::ReservationAndAllocation
1901 && (!proposal.offers.is_empty() || !proposal.requests.is_empty())
1902 {
1903 return Err(CanwuError::new(
1904 ErrorCode::InvalidBoundary,
1905 format!(
1906 "boundary system {plugin}.{} proposed reservations in phase {:?}",
1907 contract.name, contract.phase
1908 ),
1909 ));
1910 }
1911
1912 let entity_exists = |entity: &EntityRef| {
1913 proposal_entity_exists(
1914 current,
1915 &plugins.record_schemas,
1916 record_overlay,
1917 proposal,
1918 entity,
1919 )
1920 };
1921 let mut offered_pools = BTreeSet::new();
1922 for offer in &proposal.offers {
1923 validate_reservation_pool(&offer.pool, &entity_exists)?;
1924 if !contract.reservation_offers.contains(&offer.pool.state)
1925 || plugins
1926 .state_owners
1927 .get(&offer.pool.state)
1928 .is_none_or(|owner| owner != plugin)
1929 {
1930 return Err(CanwuError::new(
1931 ErrorCode::InvalidBoundary,
1932 format!(
1933 "boundary system {plugin}.{} offered undeclared state {}.{}",
1934 contract.name, offer.pool.state.namespace, offer.pool.state.name
1935 ),
1936 ));
1937 }
1938 if !offered_pools.insert(&offer.pool) {
1939 return Err(CanwuError::new(
1940 ErrorCode::InvalidBoundary,
1941 format!(
1942 "boundary system {plugin}.{} offered the same reservation pool twice",
1943 contract.name
1944 ),
1945 ));
1946 }
1947 }
1948
1949 let mut request_names = BTreeSet::new();
1950 for request in &proposal.requests {
1951 validate_reservation_pool(&request.pool, &entity_exists)?;
1952 if request.request.trim().is_empty()
1953 || request.request != request.request.trim()
1954 || request.tie_break.trim().is_empty()
1955 || request.tie_break != request.tie_break.trim()
1956 || request.quantity == 0
1957 || !request_names.insert(&request.request)
1958 || !contract.reservation_requests.contains(&request.pool.state)
1959 {
1960 return Err(CanwuError::new(
1961 ErrorCode::InvalidBoundary,
1962 format!(
1963 "boundary system {plugin}.{} produced an invalid reservation request",
1964 contract.name
1965 ),
1966 ));
1967 }
1968 }
1969
1970 let publication_count = proposal
1971 .directives
1972 .iter()
1973 .filter(|directive| matches!(directive, BoundaryDirective::PublishKnowledge { .. }))
1974 .count();
1975 if publication_count > crate::KnowledgeLimitsV1::CURRENT.batches_per_system_boundary {
1976 return Err(CanwuError::new(
1977 ErrorCode::KnowledgeLimitExceeded,
1978 "system knowledge publication batch limit exceeded",
1979 ));
1980 }
1981 let mut component_keys = BTreeSet::new();
1982 let mut record_targets = BTreeSet::new();
1983 let mut producer_correlations = BTreeSet::new();
1984 let mut canonical_drafts = BTreeSet::new();
1985 let mut random_decision_samples = BTreeSet::new();
1986 let mut cancelled_ingress = BTreeSet::new();
1987 for directive in &proposal.directives {
1988 match directive {
1989 BoundaryDirective::SetComponent {
1990 state: state_key,
1991 entity,
1992 component,
1993 ..
1994 } => {
1995 if component.trim().is_empty()
1996 || component != component.trim()
1997 || !contract.writes.contains(state_key)
1998 || plugins
1999 .state_owners
2000 .get(state_key)
2001 .is_none_or(|owner| owner != plugin)
2002 || is_domain_record_state(&plugins.record_schemas, state_key)
2003 {
2004 return Err(CanwuError::new(
2005 ErrorCode::UndeclaredStateWrite,
2006 format!(
2007 "boundary system {plugin}.{} produced an undeclared component write",
2008 contract.name
2009 ),
2010 ));
2011 }
2012 if !entity_exists(entity) {
2013 return Err(CanwuError::new(
2014 ErrorCode::EntityNotFound,
2015 format!(
2016 "boundary system {plugin}.{} targeted missing entity {entity}",
2017 contract.name
2018 ),
2019 )
2020 .with_entity(entity.clone()));
2021 }
2022 let key = component_key(plugin, state_key, entity, component);
2023 if !component_keys.insert(key) {
2024 return Err(CanwuError::new(
2025 ErrorCode::InvalidBoundary,
2026 format!(
2027 "boundary system {plugin}.{} wrote the same component twice",
2028 contract.name
2029 ),
2030 ));
2031 }
2032 }
2033 BoundaryDirective::MutateRecord { mutation, summary } => {
2034 let target = mutation.target();
2035 let state_key = records::record_state_key(&target.kind);
2036 if !canonical_text(summary)
2037 || !contract.writes.contains(&state_key)
2038 || plugins
2039 .state_owners
2040 .get(&state_key)
2041 .is_none_or(|owner| owner != plugin)
2042 || plugins
2043 .record_schemas
2044 .get(&target.kind)
2045 .is_none_or(|(owner, _)| owner != plugin)
2046 {
2047 return Err(CanwuError::new(
2048 ErrorCode::UndeclaredStateWrite,
2049 format!(
2050 "boundary system {plugin}.{} produced an undeclared record mutation",
2051 contract.name
2052 ),
2053 ));
2054 }
2055 if !record_targets.insert(target.clone()) {
2056 return Err(CanwuError::new(
2057 ErrorCode::InvalidBoundary,
2058 format!(
2059 "boundary system {plugin}.{} mutated the same record twice",
2060 contract.name
2061 ),
2062 ));
2063 }
2064 }
2065 BoundaryDirective::Emit {
2066 event_type,
2067 affected,
2068 ..
2069 } => {
2070 if event_type.trim().is_empty()
2071 || event_type != event_type.trim()
2072 || !contract.emits.contains(event_type)
2073 {
2074 return Err(CanwuError::new(
2075 ErrorCode::InvalidBoundary,
2076 format!(
2077 "boundary system {plugin}.{} emitted an undeclared event type",
2078 contract.name
2079 ),
2080 ));
2081 }
2082 if affected.iter().any(|entity| !entity_exists(entity)) {
2083 return Err(CanwuError::new(
2084 ErrorCode::EntityNotFound,
2085 format!(
2086 "boundary system {plugin}.{} emitted an event for a missing entity",
2087 contract.name
2088 ),
2089 ));
2090 }
2091 }
2092 BoundaryDirective::PublishKnowledge {
2093 holder,
2094 visibility,
2095 producer_correlation,
2096 records,
2097 summary,
2098 } => {
2099 if !matches!(
2100 contract.phase,
2101 BoundaryPhase::PerceptionAndAttentionRefresh
2102 | BoundaryPhase::PerspectiveAndReportMaterialization
2103 ) {
2104 return Err(CanwuError::new(
2105 ErrorCode::UndeclaredKnowledgeWrite,
2106 "knowledge publication is allowed only in phases 4 and 13",
2107 ));
2108 }
2109 if records.is_empty()
2110 || records.len() > crate::KnowledgeLimitsV1::CURRENT.records_per_batch
2111 {
2112 return Err(CanwuError::new(
2113 ErrorCode::KnowledgeLimitExceeded,
2114 "knowledge publication batch is empty or exceeds its record limit",
2115 ));
2116 }
2117 if !canonical_text(summary)
2118 || summary.len() > crate::KnowledgeLimitsV1::CURRENT.text_bytes
2119 {
2120 return Err(CanwuError::new(
2121 ErrorCode::InvalidKnowledgeRecord,
2122 "knowledge publication summary is not canonical or exceeds its limit",
2123 ));
2124 }
2125 if let Some(value) = producer_correlation
2126 && (!canonical_text(value)
2127 || value.len() > 256
2128 || !producer_correlations.insert(value))
2129 {
2130 return Err(CanwuError::new(
2131 ErrorCode::InvalidKnowledgeRecord,
2132 "producer correlation is invalid or duplicated",
2133 ));
2134 }
2135 for draft in records {
2136 let Some(grant) = contract
2137 .knowledge_writes
2138 .iter()
2139 .find(|grant| grant.schema == draft.schema)
2140 else {
2141 return Err(CanwuError::new(
2142 ErrorCode::UndeclaredKnowledgeWrite,
2143 format!(
2144 "boundary system {plugin}.{} did not declare the knowledge schema",
2145 contract.name
2146 ),
2147 ));
2148 };
2149 if !grant.visibilities.contains(visibility) {
2150 return Err(CanwuError::new(
2151 ErrorCode::UndeclaredKnowledgeWrite,
2152 "knowledge publication visibility is not granted",
2153 ));
2154 }
2155 let Some((owner, schema)) = plugins.knowledge_schemas.get(&draft.schema) else {
2156 return Err(CanwuError::new(
2157 ErrorCode::InvalidKnowledgeSchema,
2158 "knowledge publication uses an unregistered schema",
2159 ));
2160 };
2161 if owner != plugin || !schema.writable {
2162 return Err(CanwuError::new(
2163 ErrorCode::UndeclaredKnowledgeWrite,
2164 "knowledge publication uses a foreign or read-only schema",
2165 ));
2166 }
2167 super::knowledge::validate_draft(
2168 draft,
2169 schema,
2170 holder,
2171 current,
2172 &plugins.record_schemas,
2173 )?;
2174 for reference in &draft.origin.evidence {
2175 validate_proposal_evidence_reference(
2176 runtime,
2177 boundary_id,
2178 pending_evidence,
2179 reference,
2180 )?;
2181 }
2182 let existing = current
2183 .knowledge
2184 .records
2185 .get(holder)
2186 .into_iter()
2187 .flat_map(|records| records.iter())
2188 .chain(
2189 knowledge_overlay
2190 .get(holder)
2191 .into_iter()
2192 .flat_map(|records| records.iter()),
2193 )
2194 .collect::<BTreeMap<_, _>>();
2195 for related in draft.supersedes.iter().chain(&draft.contradicts) {
2196 let Some(related_record) = existing.get(related) else {
2197 return Err(CanwuError::new(
2198 ErrorCode::KnowledgeRecordNotFound,
2199 "knowledge relation does not resolve for the same holder at this cut",
2200 ));
2201 };
2202 if related_record.schema.kind != draft.schema.kind {
2203 return Err(CanwuError::new(
2204 ErrorCode::InvalidKnowledgeRecord,
2205 "knowledge supersession and contradiction cannot cross schema kinds",
2206 ));
2207 }
2208 }
2209 let encoded = serde_json::to_vec(&(holder, draft)).map_err(|error| {
2210 CanwuError::new(
2211 ErrorCode::InvalidKnowledgeRecord,
2212 format!("holder-scoped knowledge draft could not be encoded: {error}"),
2213 )
2214 })?;
2215 if !canonical_drafts.insert(encoded) {
2216 return Err(CanwuError::new(
2217 ErrorCode::InvalidKnowledgeRecord,
2218 "one system proposal contains a duplicate canonical knowledge draft",
2219 ));
2220 }
2221 }
2222 }
2223 BoundaryDirective::ScheduleIngress {
2224 after,
2225 packet_type,
2226 payload,
2227 affected,
2228 ..
2229 } => {
2230 let descriptor = plugins
2231 .ingress
2232 .get(&(plugin.to_owned(), packet_type.clone()))
2233 .ok_or_else(|| {
2234 CanwuError::new(
2235 ErrorCode::InvalidPayload,
2236 format!(
2237 "boundary system {plugin}.{} scheduled undeclared ingress type {packet_type}",
2238 contract.name
2239 ),
2240 )
2241 })?;
2242 if after.is_negative() || now.checked_add(*after).is_none() {
2243 return Err(CanwuError::new(
2244 ErrorCode::InvalidDuration,
2245 "boundary-generated ingress requires a nonnegative supported delay",
2246 ));
2247 }
2248 descriptor.payload_schema.validate(payload)?;
2249 if affected.iter().any(|entity| {
2250 !proposal_entity_identity_exists(
2251 current,
2252 &plugins.record_schemas,
2253 proposal,
2254 entity,
2255 )
2256 }) {
2257 return Err(CanwuError::new(
2258 ErrorCode::EntityNotFound,
2259 format!(
2260 "boundary system {plugin}.{} scheduled ingress for an unknown entity identity",
2261 contract.name
2262 ),
2263 ));
2264 }
2265 }
2266 BoundaryDirective::SchedulePluginIngress {
2267 target_plugin,
2268 after,
2269 packet_type,
2270 payload,
2271 affected,
2272 ..
2273 } => {
2274 let grant = super::PluginIngressTarget {
2275 target_plugin: target_plugin.clone(),
2276 packet_type: packet_type.clone(),
2277 };
2278 if !contract.plugin_ingress_targets.contains(&grant) {
2279 return Err(CanwuError::new(
2280 ErrorCode::UndeclaredStateWrite,
2281 format!(
2282 "boundary system {plugin}.{} did not declare target ingress {target_plugin}.{packet_type}",
2283 contract.name
2284 ),
2285 ));
2286 }
2287 let descriptor = plugins
2288 .ingress
2289 .get(&(target_plugin.clone(), packet_type.clone()))
2290 .ok_or_else(|| {
2291 CanwuError::new(
2292 ErrorCode::InvalidPayload,
2293 format!(
2294 "boundary system {plugin}.{} scheduled undeclared target ingress {target_plugin}.{packet_type}",
2295 contract.name
2296 ),
2297 )
2298 })?;
2299 if after.is_negative() || now.checked_add(*after).is_none() {
2300 return Err(CanwuError::new(
2301 ErrorCode::InvalidDuration,
2302 "boundary-generated cross-plugin ingress requires a nonnegative supported delay",
2303 ));
2304 }
2305 descriptor.payload_schema.validate(payload)?;
2306 if affected.iter().any(|entity| {
2307 !proposal_entity_identity_exists(
2308 current,
2309 &plugins.record_schemas,
2310 proposal,
2311 entity,
2312 )
2313 }) {
2314 return Err(CanwuError::new(
2315 ErrorCode::EntityNotFound,
2316 format!(
2317 "boundary system {plugin}.{} scheduled cross-plugin ingress for an unknown entity identity",
2318 contract.name
2319 ),
2320 ));
2321 }
2322 }
2323 BoundaryDirective::CancelPluginIngress { ingress_id, reason } => {
2324 if !valid_ingress_cancellation_reason(reason)
2325 || !cancelled_ingress.insert(*ingress_id)
2326 {
2327 return Err(CanwuError::new(
2328 ErrorCode::InvalidPayload,
2329 format!(
2330 "boundary system {plugin}.{} proposed a duplicate or malformed ingress cancellation",
2331 contract.name
2332 ),
2333 ));
2334 }
2335 }
2336 BoundaryDirective::ResolveDecisionRandomly { resolution } => {
2337 validate_random_decision_resolution(
2338 plugin,
2339 contract,
2340 current,
2341 committed_availability,
2342 pending_random_draws,
2343 &mut random_decision_samples,
2344 resolution,
2345 )?;
2346 }
2347 BoundaryDirective::SetPersonAvailability {
2348 person,
2349 availability,
2350 summary,
2351 } => {
2352 super::persons::validate_availability_directive(
2353 plugin,
2354 contract,
2355 now,
2356 *person,
2357 availability,
2358 summary,
2359 &entity_exists,
2360 )?;
2361 }
2362 BoundaryDirective::RecordEvaluationTrace { trace } => {
2363 super::evaluation::validate_trace_shape(
2364 contract.phase,
2365 trace,
2366 boundary_id,
2367 runtime.metadata.run_configuration.evaluation_limits(),
2368 )?;
2369 if !proposal_entity_identity_exists(
2370 current,
2371 &plugins.record_schemas,
2372 proposal,
2373 &trace.subject,
2374 ) {
2375 return Err(CanwuError::new(
2376 ErrorCode::EntityNotFound,
2377 format!(
2378 "boundary system {plugin}.{} traced an evaluation of unknown subject {}",
2379 contract.name, trace.subject
2380 ),
2381 )
2382 .with_entity(trace.subject.clone()));
2383 }
2384 for reference in trace.terms.iter().flat_map(|term| &term.evidence) {
2385 validate_proposal_evidence_reference(
2386 runtime,
2387 boundary_id,
2388 pending_evidence,
2389 reference,
2390 )?;
2391 }
2392 }
2393 BoundaryDirective::CreatePerson {
2394 draft,
2395 correlation,
2396 summary,
2397 } => {
2398 super::persons::validate_person_draft(
2399 &super::persons::PersonDraftContext {
2400 plugin,
2401 contract,
2402 now,
2403 government_exists: &|id| current.governments.contains_key(&id),
2404 territory_exists: &|id| current.territories.contains_key(&id),
2405 entity_exists: &entity_exists,
2406 },
2407 draft,
2408 correlation,
2409 summary,
2410 )?;
2411 validate_proposal_evidence_reference(
2412 runtime,
2413 boundary_id,
2414 pending_evidence,
2415 &draft.provenance,
2416 )?;
2417 }
2418 BoundaryDirective::RegisterTransitionManifest { .. }
2419 | BoundaryDirective::StageTransitionWrite { .. } => {
2420 return Err(transition_directive_not_admitted(plugin, &contract.name));
2421 }
2422 }
2423 }
2424 Ok(())
2425}
2426
2427fn transition_directive_not_admitted(plugin: &str, system: &str) -> CanwuError {
2430 CanwuError::new(
2431 ErrorCode::InvalidBoundary,
2432 format!(
2433 "boundary system {plugin}.{system} produced a transition directive outside the transition ledger"
2434 ),
2435 )
2436}
2437
2438fn validate_random_decision_resolution(
2452 plugin: &str,
2453 contract: &BoundarySystemContract,
2454 current: &RuntimeCurrentState,
2455 committed_availability: &BTreeMap<super::PersonId, super::PersonAvailability>,
2456 pending_random_draws: &[random::PendingRandomDraw],
2457 used_samples: &mut BTreeSet<(super::RandomStreamKey, RandomDrawAddress)>,
2458 resolution: &super::RandomDecisionResolution,
2459) -> Result<(), CanwuError> {
2460 if resolution.decision_request_id.get() == 0
2461 || resolution
2462 .command_request_id
2463 .is_some_and(|request_id| request_id.get() == 0)
2464 || resolution.expected_version == 0
2465 {
2466 return Err(CanwuError::new(
2467 ErrorCode::InvalidDecision,
2468 "random decision resolution requires nonzero request IDs and ticket version",
2469 ));
2470 }
2471 let ticket = current
2472 .decisions
2473 .ticket(resolution.ticket_id)
2474 .ok_or_else(|| {
2475 CanwuError::new(
2476 ErrorCode::InvalidDecision,
2477 "random decision resolution references an unknown ticket",
2478 )
2479 })?;
2480 if !ticket.is_open()
2481 || ticket.version != resolution.expected_version
2482 || ticket.assigned_controller != resolution.controller_id
2483 {
2484 return Err(CanwuError::new(
2485 ErrorCode::InvalidDecision,
2486 "random decision resolution references a closed, stale, or differently controlled ticket",
2487 ));
2488 }
2489 let controller = current
2490 .decisions
2491 .controller(&resolution.controller_id)
2492 .ok_or_else(|| {
2493 CanwuError::new(
2494 ErrorCode::InvalidDecision,
2495 "random decision resolution references an unknown controller",
2496 )
2497 })?;
2498 super::persons::validate_decision_preparation(committed_availability, ticket, controller)?;
2499 match &resolution.tie_break {
2500 None if controller.policy.kind != DecisionPolicyKind::Random => {
2501 return Err(CanwuError::new(
2502 ErrorCode::InvalidDecision,
2503 "random decision resolution requires a controller with random policy identity",
2504 ));
2505 }
2506 None => {}
2507 Some(pending) => validate_random_tie_break(controller, ticket, resolution, pending)?,
2508 }
2509 if !contract.random_streams.contains(&resolution.sample.stream) {
2510 return Err(CanwuError::new(
2511 ErrorCode::UndeclaredRandomStream,
2512 format!(
2513 "boundary system {plugin}.{} did not declare the random decision stream",
2514 contract.name
2515 ),
2516 ));
2517 }
2518 let RandomDrawAddress::OperationV1(address) = &resolution.sample.address else {
2519 return Err(CanwuError::new(
2520 ErrorCode::InvalidRandomDraw,
2521 "random decisions require an operation-keyed draw",
2522 ));
2523 };
2524 if address.producer_plugin != plugin
2525 || address.target
2526 != (RandomOperationTarget::DecisionTicket {
2527 ticket_id: ticket.id,
2528 ticket_version: ticket.version,
2529 })
2530 {
2531 return Err(CanwuError::new(
2532 ErrorCode::InvalidRandomDraw,
2533 "random decision draw address does not bind the current ticket version",
2534 ));
2535 }
2536 let sample_key = (
2537 resolution.sample.stream.clone(),
2538 resolution.sample.address.clone(),
2539 );
2540 if !used_samples.insert(sample_key.clone()) {
2541 return Err(CanwuError::new(
2542 ErrorCode::InvalidRandomDraw,
2543 "one random draw cannot resolve more than one decision",
2544 ));
2545 }
2546 if !pending_random_draws.iter().any(|draw| {
2547 draw.stream == sample_key.0
2548 && draw.address == sample_key.1
2549 && draw.upper_exclusive == resolution.sample.upper_exclusive
2550 && draw.value == resolution.sample.value
2551 }) {
2552 return Err(CanwuError::new(
2553 ErrorCode::InvalidRandomDraw,
2554 "random decision resolution does not reference a draw produced by this proposal",
2555 ));
2556 }
2557 let total_weight = resolution
2558 .option_weights
2559 .iter()
2560 .try_fold(0_u64, |total, option| total.checked_add(option.weight))
2561 .ok_or_else(|| {
2562 CanwuError::new(
2563 ErrorCode::InvalidDecision,
2564 "random decision option weights overflow the supported range",
2565 )
2566 })?;
2567 if total_weight != resolution.sample.upper_exclusive {
2568 return Err(CanwuError::new(
2569 ErrorCode::InvalidDecision,
2570 "random decision option weights disagree with the draw bound",
2571 ));
2572 }
2573 let selected = random_resolution_selection(ticket, resolution)?;
2574 let action = &ticket
2575 .option(&selected)
2576 .expect("validated random decision selected an existing option")
2577 .action;
2578 if matches!(action, DecisionAction::Command { .. }) != resolution.command_request_id.is_some() {
2579 return Err(CanwuError::new(
2580 ErrorCode::InvalidDecision,
2581 "random decision command options require exactly one command request ID",
2582 ));
2583 }
2584 Ok(())
2585}
2586
2587fn validate_random_tie_break(
2591 controller: &super::DecisionControllerBinding,
2592 ticket: &super::DecisionTicket,
2593 resolution: &super::RandomDecisionResolution,
2594 pending: &PolicyDecision,
2595) -> Result<(), CanwuError> {
2596 if controller.policy.kind != DecisionPolicyKind::Utility || !controller.random_tie_break {
2597 return Err(CanwuError::new(
2598 ErrorCode::InvalidDecision,
2599 "a random tie-break requires a utility-policy controller that permits tie-breaks",
2600 ));
2601 }
2602 let DecisionOutcome::PendingRandom { candidates } = &pending.outcome else {
2603 return Err(CanwuError::new(
2604 ErrorCode::InvalidDecision,
2605 "a random tie-break must carry a pending random policy decision",
2606 ));
2607 };
2608 if candidates != &resolution.option_weights
2609 || !pending.is_random_tie_break()
2610 || pending.external.is_some()
2611 || pending.random.is_some()
2612 {
2613 return Err(CanwuError::new(
2614 ErrorCode::InvalidDecision,
2615 "random tie-break weights must equal the pending candidates of an evidence-free random stage",
2616 ));
2617 }
2618 pending
2619 .validate(ticket)
2620 .map_err(|error| CanwuError::new(ErrorCode::InvalidDecision, error.to_string()))
2621}
2622
2623pub(super) fn random_resolution_selection(
2627 ticket: &super::DecisionTicket,
2628 resolution: &super::RandomDecisionResolution,
2629) -> Result<String, CanwuError> {
2630 if resolution.tie_break.is_some() {
2631 DecisionRandomEvidence::selected_candidate(
2632 ticket,
2633 &resolution.option_weights,
2634 resolution.sample.value,
2635 )
2636 } else {
2637 DecisionRandomEvidence::selected_option(
2638 ticket,
2639 &resolution.option_weights,
2640 resolution.sample.value,
2641 )
2642 }
2643 .map_err(|error| CanwuError::new(ErrorCode::InvalidDecision, error.to_string()))
2644}
2645
2646pub(super) fn random_policy_summary(option_id: &str) -> String {
2647 format!("random policy selected {option_id}")
2648}
2649
2650pub(super) fn random_tie_break_summary(option_id: &str) -> String {
2651 format!("random tie-break selected {option_id}")
2652}
2653
2654fn validate_proposal_evidence_reference(
2655 runtime: &RuntimeState,
2656 boundary_id: BoundaryId,
2657 pending: &PendingBoundaryEvidence,
2658 reference: &EvidenceRef,
2659) -> Result<(), CanwuError> {
2660 if let EvidenceRef::DomainRecordVersion(version) = reference
2661 && let DomainRecordVersionSource::BoundaryChange {
2662 boundary,
2663 change_index,
2664 } = version.established_by
2665 && boundary == boundary_id
2666 {
2667 let resolved = usize::try_from(change_index)
2668 .ok()
2669 .and_then(|index| pending.record_changes.get(index))
2670 .is_some_and(|change| {
2671 change.current.reference == version.record
2672 && change.current.version == version.version
2673 });
2674 return if resolved {
2675 Ok(())
2676 } else {
2677 Err(CanwuError::new(
2678 ErrorCode::EvidenceUnavailable,
2679 "knowledge origin references an unavailable current-boundary record version",
2680 ))
2681 };
2682 }
2683
2684 if let EvidenceRef::Event(id) = reference
2685 && runtime
2686 .evidence
2687 .retained_event(*id)
2688 .is_some_and(|event| event.cause == Some(CauseRef::Boundary(boundary_id)))
2689 {
2690 if pending
2691 .emissions
2692 .iter()
2693 .any(|emission| emission.event == *id)
2694 {
2695 return Ok(());
2696 }
2697 return Err(CanwuError::new(
2698 ErrorCode::EvidenceUnavailable,
2699 "knowledge origin references an event outside the proposal-visible boundary cut",
2700 ));
2701 }
2702
2703 if let EvidenceRef::Ingress(id) = reference
2704 && runtime
2705 .evidence
2706 .retained_ingress(*id)
2707 .is_some_and(|record| record.cause == Some(CauseRef::Boundary(boundary_id)))
2708 {
2709 return Err(CanwuError::new(
2710 ErrorCode::EvidenceUnavailable,
2711 "current-boundary generated ingress is not proposal-visible evidence",
2712 ));
2713 }
2714
2715 match resolve_evidence_reference(&RuntimeValidationContext::new(runtime), reference) {
2716 EvidenceAvailability::Retained | EvidenceAvailability::Archived => Ok(()),
2717 EvidenceAvailability::Missing => Err(CanwuError::new(
2718 ErrorCode::EvidenceUnavailable,
2719 "knowledge origin references missing or wrong-version evidence",
2720 )),
2721 }
2722}
2723
2724fn validate_reservation_pool(
2725 pool: &ReservationPoolKey,
2726 entity_exists: &dyn Fn(&EntityRef) -> bool,
2727) -> Result<(), CanwuError> {
2728 if pool.resource.trim().is_empty()
2729 || pool.resource != pool.resource.trim()
2730 || !entity_exists(&pool.entity)
2731 {
2732 return Err(CanwuError::new(
2733 ErrorCode::InvalidBoundary,
2734 "reservation pools require a canonical resource and an existing entity",
2735 ));
2736 }
2737 Ok(())
2738}
2739
2740fn extend_boundary_overlay(
2741 current: &RuntimeCurrentState,
2742 record_overlay: &BTreeMap<DomainRecordRef, DomainRecord>,
2743 overlay: &mut BTreeMap<PluginComponentKey, PluginComponentRecord>,
2744 directives: &[StagedBoundaryDirective],
2745) -> Result<(), CanwuError> {
2746 extend_boundary_component_overlay(current, record_overlay, overlay, directives, false)
2747}
2748
2749fn extend_boundary_candidate_overlay(
2750 current: &RuntimeCurrentState,
2751 record_overlay: &BTreeMap<DomainRecordRef, DomainRecord>,
2752 overlay: &mut BTreeMap<PluginComponentKey, PluginComponentRecord>,
2753 directives: &[StagedBoundaryDirective],
2754) -> Result<(), CanwuError> {
2755 extend_boundary_component_overlay(current, record_overlay, overlay, directives, true)
2756}
2757
2758fn extend_boundary_component_overlay(
2759 current: &RuntimeCurrentState,
2760 record_overlay: &BTreeMap<DomainRecordRef, DomainRecord>,
2761 overlay: &mut BTreeMap<PluginComponentKey, PluginComponentRecord>,
2762 directives: &[StagedBoundaryDirective],
2763 include_next_boundary: bool,
2764) -> Result<(), CanwuError> {
2765 for staged in directives.iter().filter(|staged| {
2766 include_next_boundary || staged.visibility == StateVisibility::SameBoundary
2767 }) {
2768 if let BoundaryDirective::SetComponent {
2769 state: state_key,
2770 entity,
2771 component,
2772 value,
2773 ..
2774 } = &staged.directive
2775 {
2776 let key = component_key(&staged.plugin, state_key, entity, component);
2777 if overlay.contains_key(&key) {
2778 return Err(CanwuError::new(
2779 ErrorCode::InvalidBoundary,
2780 "multiple boundary proposals target the same component",
2781 ));
2782 }
2783 if !runtime_entity_exists_with_record_overlay(current, record_overlay, entity) {
2784 return Err(CanwuError::new(
2785 ErrorCode::EntityNotFound,
2786 format!("boundary proposal targeted missing entity {entity}"),
2787 ));
2788 }
2789 overlay.insert(
2790 key,
2791 PluginComponentRecord {
2792 plugin: staged.plugin.clone(),
2793 state: state_key.clone(),
2794 entity: entity.clone(),
2795 component: component.clone(),
2796 value: value.clone(),
2797 },
2798 );
2799 }
2800 }
2801 Ok(())
2802}
2803
2804fn extend_boundary_record_overlay(
2805 context: &BoundaryRecordOverlayContext<'_>,
2806 overlay: &mut BTreeMap<DomainRecordRef, DomainRecord>,
2807 directives: &[StagedBoundaryDirective],
2808) -> Result<(), CanwuError> {
2809 extend_boundary_domain_record_overlay(context, overlay, directives, false)
2810}
2811
2812fn extend_boundary_record_candidate_overlay(
2813 context: &BoundaryRecordOverlayContext<'_>,
2814 overlay: &mut BTreeMap<DomainRecordRef, DomainRecord>,
2815 directives: &[StagedBoundaryDirective],
2816) -> Result<(), CanwuError> {
2817 extend_boundary_domain_record_overlay(context, overlay, directives, true)
2818}
2819
2820struct BoundaryRecordOverlayContext<'a> {
2821 current: &'a RuntimeCurrentState,
2822 now: SimTime,
2823 scheduled_actions: &'a BTreeMap<ScheduleKey, ScheduledAction>,
2824 run_configuration: &'a RunConfigurationSnapshot,
2825 schemas: &'a records::DomainRecordSchemas,
2826}
2827
2828fn extend_boundary_domain_record_overlay(
2829 context: &BoundaryRecordOverlayContext<'_>,
2830 overlay: &mut BTreeMap<DomainRecordRef, DomainRecord>,
2831 directives: &[StagedBoundaryDirective],
2832 include_next_boundary: bool,
2833) -> Result<(), CanwuError> {
2834 let requests: Vec<_> = directives
2835 .iter()
2836 .filter(|staged| {
2837 include_next_boundary || staged.visibility == StateVisibility::SameBoundary
2838 })
2839 .filter_map(|staged| match &staged.directive {
2840 BoundaryDirective::MutateRecord { mutation, summary } => {
2841 Some(records::DomainMutationRequest {
2842 plugin: &staged.plugin,
2843 system: &staged.system,
2844 visibility: staged.visibility,
2845 mutation,
2846 summary,
2847 })
2848 }
2849 BoundaryDirective::SetComponent { .. }
2850 | BoundaryDirective::Emit { .. }
2851 | BoundaryDirective::ScheduleIngress { .. }
2852 | BoundaryDirective::SchedulePluginIngress { .. }
2853 | BoundaryDirective::ResolveDecisionRandomly { .. }
2854 | BoundaryDirective::PublishKnowledge { .. }
2855 | BoundaryDirective::SetPersonAvailability { .. }
2856 | BoundaryDirective::CreatePerson { .. }
2857 | BoundaryDirective::CancelPluginIngress { .. }
2858 | BoundaryDirective::RecordEvaluationTrace { .. }
2859 | BoundaryDirective::RegisterTransitionManifest { .. }
2860 | BoundaryDirective::StageTransitionWrite { .. } => None,
2861 })
2862 .collect();
2863 if requests.is_empty() {
2864 return Ok(());
2865 }
2866 let (next, changes) = records::apply_mutation_bundle_cow_with_overlay(
2867 &context.current.domain_records,
2868 overlay,
2869 context.schemas,
2870 context.now,
2871 &|entity| runtime_current_entity_exists(context.current, entity),
2872 requests,
2873 )?;
2874 validate_domain_dependents_with_records(
2875 &context.current.plugin_components,
2876 context.scheduled_actions,
2877 context.run_configuration,
2878 &next,
2879 )?;
2880 for change in changes {
2881 overlay.insert(change.current.reference.clone(), change.current);
2882 }
2883 Ok(())
2884}
2885
2886fn partition_boundary_visibility(
2887 directives: Vec<StagedBoundaryDirective>,
2888) -> (Vec<StagedBoundaryDirective>, Vec<StagedBoundaryDirective>) {
2889 directives
2890 .into_iter()
2891 .partition(|staged| staged.visibility == StateVisibility::SameBoundary)
2892}
2893
2894fn partition_knowledge_directives(
2895 directives: Vec<StagedBoundaryDirective>,
2896) -> (Vec<StagedBoundaryDirective>, Vec<StagedBoundaryDirective>) {
2897 directives
2898 .into_iter()
2899 .partition(|staged| matches!(staged.directive, BoundaryDirective::PublishKnowledge { .. }))
2900}
2901
2902fn allocate_reservations(
2903 mut offers: Vec<PendingReservationOffer>,
2904 mut requests: Vec<PendingReservationRequest>,
2905) -> Result<ReservationAllocationResult, CanwuError> {
2906 offers.sort_by(|left, right| {
2907 left.offer
2908 .pool
2909 .cmp(&right.offer.pool)
2910 .then_with(|| left.plugin.cmp(&right.plugin))
2911 .then_with(|| left.system.cmp(&right.system))
2912 });
2913 let mut remaining = BTreeMap::new();
2914 let mut offer_records = Vec::new();
2915 for pending in offers {
2916 if remaining
2917 .insert(pending.offer.pool.clone(), pending.offer.capacity)
2918 .is_some()
2919 {
2920 return Err(CanwuError::new(
2921 ErrorCode::InvalidBoundary,
2922 format!(
2923 "reservation pool was offered more than once, including by {}.{}",
2924 pending.plugin, pending.system
2925 ),
2926 ));
2927 }
2928 offer_records.push(ReservationOfferRecord {
2929 plugin: pending.plugin,
2930 system: pending.system,
2931 offer: pending.offer,
2932 });
2933 }
2934 requests.sort_by(|left, right| {
2935 left.request
2936 .pool
2937 .cmp(&right.request.pool)
2938 .then_with(|| right.request.priority.cmp(&left.request.priority))
2939 .then_with(|| left.request.tie_break.cmp(&right.request.tie_break))
2940 .then_with(|| left.reservation.cmp(&right.reservation))
2941 });
2942 let mut seen = BTreeSet::new();
2943 let mut by_reservation = BTreeMap::new();
2944 let mut request_records = Vec::new();
2945 let mut records = Vec::new();
2946 for pending in requests {
2947 if !seen.insert(pending.reservation.clone()) {
2948 return Err(CanwuError::new(
2949 ErrorCode::InvalidBoundary,
2950 "reservation request identity is duplicated",
2951 ));
2952 }
2953 request_records.push(ReservationRequestRecord {
2954 reservation: pending.reservation.clone(),
2955 request: pending.request.clone(),
2956 });
2957 let available = remaining.entry(pending.request.pool.clone()).or_default();
2958 let granted = pending.request.quantity.min(*available);
2959 *available -= granted;
2960 let disposition = if granted == pending.request.quantity {
2961 ReservationDisposition::Fulfilled
2962 } else if granted == 0 {
2963 ReservationDisposition::Rejected
2964 } else {
2965 ReservationDisposition::Partial
2966 };
2967 let allocation = ReservationAllocation {
2968 reservation: pending.reservation.clone(),
2969 pool: pending.request.pool,
2970 requested: pending.request.quantity,
2971 granted,
2972 remaining_after: *available,
2973 disposition,
2974 };
2975 by_reservation.insert(pending.reservation, allocation.clone());
2976 records.push(allocation);
2977 }
2978 Ok(ReservationAllocationResult {
2979 by_reservation,
2980 offers: offer_records,
2981 requests: request_records,
2982 records,
2983 })
2984}