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 state.metadata.current_domain_record_versions.insert(
1766 change.current.reference.clone(),
1767 super::DomainRecordVersionRef {
1768 record: change.current.reference.clone(),
1769 version: change.current.version,
1770 established_by: super::DomainRecordVersionSource::BoundaryChange {
1771 boundary,
1772 change_index,
1773 },
1774 },
1775 );
1776}
1777
1778struct PendingReservationOffer {
1779 plugin: String,
1780 system: String,
1781 offer: ReservationOffer,
1782}
1783
1784struct PendingReservationRequest {
1785 reservation: ReservationRef,
1786 request: ReservationRequest,
1787}
1788
1789struct ReservationAllocationResult {
1790 by_reservation: BTreeMap<ReservationRef, ReservationAllocation>,
1791 offers: Vec<ReservationOfferRecord>,
1792 requests: Vec<ReservationRequestRecord>,
1793 records: Vec<ReservationAllocation>,
1794}
1795
1796struct StagedBoundaryDirective {
1797 plugin: String,
1798 system: String,
1799 phase: BoundaryPhase,
1800 visibility: StateVisibility,
1801 directive: BoundaryDirective,
1802}
1803
1804#[derive(Default)]
1805struct PendingBoundaryEvidence {
1806 changes: Vec<BoundaryChange>,
1807 record_changes: Vec<DomainRecordChange>,
1808 emissions: Vec<BoundaryEmission>,
1809 generated_ingress: Vec<BoundaryIngressGeneration>,
1810 random_decisions: Vec<PendingRandomDecisionResolution>,
1811 person_availability_changes: Vec<super::BoundaryPersonAvailabilityChange>,
1812 created_persons: Vec<super::BoundaryPersonCreation>,
1813}
1814
1815struct PendingRandomDecisionResolution {
1816 plugin: String,
1817 system: String,
1818 phase: BoundaryPhase,
1819 visibility: StateVisibility,
1820 resolution: super::RandomDecisionResolution,
1821}
1822
1823fn proposal_evidence_refs(
1824 boundary: BoundaryId,
1825 pending: &PendingBoundaryEvidence,
1826) -> BTreeSet<EvidenceRef> {
1827 let mut values = BTreeSet::new();
1828 for (index, change) in pending.record_changes.iter().enumerate() {
1829 if let Ok(change_index) = u64::try_from(index) {
1830 values.insert(EvidenceRef::DomainRecordVersion(
1831 super::DomainRecordVersionRef {
1832 record: change.current.reference.clone(),
1833 version: change.current.version,
1834 established_by: DomainRecordVersionSource::BoundaryChange {
1835 boundary,
1836 change_index,
1837 },
1838 },
1839 ));
1840 }
1841 }
1842 values.extend(
1843 pending
1844 .emissions
1845 .iter()
1846 .map(|emission| EvidenceRef::Event(emission.event)),
1847 );
1848 values
1849}
1850
1851pub(super) struct PendingBoundaryRandomDraw {
1852 pub(super) plugin: String,
1853 pub(super) system: String,
1854 pub(super) draw: random::PendingRandomDraw,
1855}
1856
1857pub(super) struct CommittedBoundaryRandomDraw {
1858 pub(super) id: super::RandomDrawId,
1859 pub(super) stream: super::RandomStreamKey,
1860 pub(super) address: RandomDrawAddress,
1861}
1862
1863pub(super) fn boundary_system_due(
1864 contract: &BoundarySystemContract,
1865 cadences: &[SystemCadence],
1866 has_admitted_events: bool,
1867) -> bool {
1868 match contract.cadence {
1869 SystemCadence::EventDriven => has_admitted_events,
1870 _ => cadences.contains(&contract.cadence),
1871 }
1872}
1873
1874pub(super) fn boundary_has_event_ingress(record: &BoundaryRecord) -> bool {
1875 !record.admitted_events.is_empty() || !record.admitted_ingress.is_empty()
1876}
1877
1878#[allow(clippy::too_many_arguments)]
1879fn validate_boundary_proposal(
1880 plugin: &str,
1881 contract: &BoundarySystemContract,
1882 current: &RuntimeCurrentState,
1883 committed_availability: &BTreeMap<super::PersonId, super::PersonAvailability>,
1884 now: SimTime,
1885 runtime: &RuntimeState,
1886 boundary_id: BoundaryId,
1887 pending_evidence: &PendingBoundaryEvidence,
1888 plugins: &PluginRegistry,
1889 record_overlay: &BTreeMap<DomainRecordRef, DomainRecord>,
1890 knowledge_overlay: &BTreeMap<KnowledgeHolderRef, BTreeMap<KnowledgeRecordId, KnowledgeRecord>>,
1891 proposal: &BoundaryProposal,
1892 pending_random_draws: &[random::PendingRandomDraw],
1893) -> Result<(), CanwuError> {
1894 if contract.phase != BoundaryPhase::ReservationAndAllocation
1895 && (!proposal.offers.is_empty() || !proposal.requests.is_empty())
1896 {
1897 return Err(CanwuError::new(
1898 ErrorCode::InvalidBoundary,
1899 format!(
1900 "boundary system {plugin}.{} proposed reservations in phase {:?}",
1901 contract.name, contract.phase
1902 ),
1903 ));
1904 }
1905
1906 let entity_exists = |entity: &EntityRef| {
1907 proposal_entity_exists(
1908 current,
1909 &plugins.record_schemas,
1910 record_overlay,
1911 proposal,
1912 entity,
1913 )
1914 };
1915 let mut offered_pools = BTreeSet::new();
1916 for offer in &proposal.offers {
1917 validate_reservation_pool(&offer.pool, &entity_exists)?;
1918 if !contract.reservation_offers.contains(&offer.pool.state)
1919 || plugins
1920 .state_owners
1921 .get(&offer.pool.state)
1922 .is_none_or(|owner| owner != plugin)
1923 {
1924 return Err(CanwuError::new(
1925 ErrorCode::InvalidBoundary,
1926 format!(
1927 "boundary system {plugin}.{} offered undeclared state {}.{}",
1928 contract.name, offer.pool.state.namespace, offer.pool.state.name
1929 ),
1930 ));
1931 }
1932 if !offered_pools.insert(&offer.pool) {
1933 return Err(CanwuError::new(
1934 ErrorCode::InvalidBoundary,
1935 format!(
1936 "boundary system {plugin}.{} offered the same reservation pool twice",
1937 contract.name
1938 ),
1939 ));
1940 }
1941 }
1942
1943 let mut request_names = BTreeSet::new();
1944 for request in &proposal.requests {
1945 validate_reservation_pool(&request.pool, &entity_exists)?;
1946 if request.request.trim().is_empty()
1947 || request.request != request.request.trim()
1948 || request.tie_break.trim().is_empty()
1949 || request.tie_break != request.tie_break.trim()
1950 || request.quantity == 0
1951 || !request_names.insert(&request.request)
1952 || !contract.reservation_requests.contains(&request.pool.state)
1953 {
1954 return Err(CanwuError::new(
1955 ErrorCode::InvalidBoundary,
1956 format!(
1957 "boundary system {plugin}.{} produced an invalid reservation request",
1958 contract.name
1959 ),
1960 ));
1961 }
1962 }
1963
1964 let publication_count = proposal
1965 .directives
1966 .iter()
1967 .filter(|directive| matches!(directive, BoundaryDirective::PublishKnowledge { .. }))
1968 .count();
1969 if publication_count > crate::KnowledgeLimitsV1::CURRENT.batches_per_system_boundary {
1970 return Err(CanwuError::new(
1971 ErrorCode::KnowledgeLimitExceeded,
1972 "system knowledge publication batch limit exceeded",
1973 ));
1974 }
1975 let mut component_keys = BTreeSet::new();
1976 let mut record_targets = BTreeSet::new();
1977 let mut producer_correlations = BTreeSet::new();
1978 let mut canonical_drafts = BTreeSet::new();
1979 let mut random_decision_samples = BTreeSet::new();
1980 let mut cancelled_ingress = BTreeSet::new();
1981 for directive in &proposal.directives {
1982 match directive {
1983 BoundaryDirective::SetComponent {
1984 state: state_key,
1985 entity,
1986 component,
1987 ..
1988 } => {
1989 if component.trim().is_empty()
1990 || component != component.trim()
1991 || !contract.writes.contains(state_key)
1992 || plugins
1993 .state_owners
1994 .get(state_key)
1995 .is_none_or(|owner| owner != plugin)
1996 || is_domain_record_state(&plugins.record_schemas, state_key)
1997 {
1998 return Err(CanwuError::new(
1999 ErrorCode::UndeclaredStateWrite,
2000 format!(
2001 "boundary system {plugin}.{} produced an undeclared component write",
2002 contract.name
2003 ),
2004 ));
2005 }
2006 if !entity_exists(entity) {
2007 return Err(CanwuError::new(
2008 ErrorCode::EntityNotFound,
2009 format!(
2010 "boundary system {plugin}.{} targeted missing entity {entity}",
2011 contract.name
2012 ),
2013 )
2014 .with_entity(entity.clone()));
2015 }
2016 let key = component_key(plugin, state_key, entity, component);
2017 if !component_keys.insert(key) {
2018 return Err(CanwuError::new(
2019 ErrorCode::InvalidBoundary,
2020 format!(
2021 "boundary system {plugin}.{} wrote the same component twice",
2022 contract.name
2023 ),
2024 ));
2025 }
2026 }
2027 BoundaryDirective::MutateRecord { mutation, summary } => {
2028 let target = mutation.target();
2029 let state_key = records::record_state_key(&target.kind);
2030 if !canonical_text(summary)
2031 || !contract.writes.contains(&state_key)
2032 || plugins
2033 .state_owners
2034 .get(&state_key)
2035 .is_none_or(|owner| owner != plugin)
2036 || plugins
2037 .record_schemas
2038 .get(&target.kind)
2039 .is_none_or(|(owner, _)| owner != plugin)
2040 {
2041 return Err(CanwuError::new(
2042 ErrorCode::UndeclaredStateWrite,
2043 format!(
2044 "boundary system {plugin}.{} produced an undeclared record mutation",
2045 contract.name
2046 ),
2047 ));
2048 }
2049 if !record_targets.insert(target.clone()) {
2050 return Err(CanwuError::new(
2051 ErrorCode::InvalidBoundary,
2052 format!(
2053 "boundary system {plugin}.{} mutated the same record twice",
2054 contract.name
2055 ),
2056 ));
2057 }
2058 }
2059 BoundaryDirective::Emit {
2060 event_type,
2061 affected,
2062 ..
2063 } => {
2064 if event_type.trim().is_empty()
2065 || event_type != event_type.trim()
2066 || !contract.emits.contains(event_type)
2067 {
2068 return Err(CanwuError::new(
2069 ErrorCode::InvalidBoundary,
2070 format!(
2071 "boundary system {plugin}.{} emitted an undeclared event type",
2072 contract.name
2073 ),
2074 ));
2075 }
2076 if affected.iter().any(|entity| !entity_exists(entity)) {
2077 return Err(CanwuError::new(
2078 ErrorCode::EntityNotFound,
2079 format!(
2080 "boundary system {plugin}.{} emitted an event for a missing entity",
2081 contract.name
2082 ),
2083 ));
2084 }
2085 }
2086 BoundaryDirective::PublishKnowledge {
2087 holder,
2088 visibility,
2089 producer_correlation,
2090 records,
2091 summary,
2092 } => {
2093 if !matches!(
2094 contract.phase,
2095 BoundaryPhase::PerceptionAndAttentionRefresh
2096 | BoundaryPhase::PerspectiveAndReportMaterialization
2097 ) {
2098 return Err(CanwuError::new(
2099 ErrorCode::UndeclaredKnowledgeWrite,
2100 "knowledge publication is allowed only in phases 4 and 13",
2101 ));
2102 }
2103 if records.is_empty()
2104 || records.len() > crate::KnowledgeLimitsV1::CURRENT.records_per_batch
2105 {
2106 return Err(CanwuError::new(
2107 ErrorCode::KnowledgeLimitExceeded,
2108 "knowledge publication batch is empty or exceeds its record limit",
2109 ));
2110 }
2111 if !canonical_text(summary)
2112 || summary.len() > crate::KnowledgeLimitsV1::CURRENT.text_bytes
2113 {
2114 return Err(CanwuError::new(
2115 ErrorCode::InvalidKnowledgeRecord,
2116 "knowledge publication summary is not canonical or exceeds its limit",
2117 ));
2118 }
2119 if let Some(value) = producer_correlation
2120 && (!canonical_text(value)
2121 || value.len() > 256
2122 || !producer_correlations.insert(value))
2123 {
2124 return Err(CanwuError::new(
2125 ErrorCode::InvalidKnowledgeRecord,
2126 "producer correlation is invalid or duplicated",
2127 ));
2128 }
2129 for draft in records {
2130 let Some(grant) = contract
2131 .knowledge_writes
2132 .iter()
2133 .find(|grant| grant.schema == draft.schema)
2134 else {
2135 return Err(CanwuError::new(
2136 ErrorCode::UndeclaredKnowledgeWrite,
2137 format!(
2138 "boundary system {plugin}.{} did not declare the knowledge schema",
2139 contract.name
2140 ),
2141 ));
2142 };
2143 if !grant.visibilities.contains(visibility) {
2144 return Err(CanwuError::new(
2145 ErrorCode::UndeclaredKnowledgeWrite,
2146 "knowledge publication visibility is not granted",
2147 ));
2148 }
2149 let Some((owner, schema)) = plugins.knowledge_schemas.get(&draft.schema) else {
2150 return Err(CanwuError::new(
2151 ErrorCode::InvalidKnowledgeSchema,
2152 "knowledge publication uses an unregistered schema",
2153 ));
2154 };
2155 if owner != plugin || !schema.writable {
2156 return Err(CanwuError::new(
2157 ErrorCode::UndeclaredKnowledgeWrite,
2158 "knowledge publication uses a foreign or read-only schema",
2159 ));
2160 }
2161 super::knowledge::validate_draft(
2162 draft,
2163 schema,
2164 holder,
2165 current,
2166 &plugins.record_schemas,
2167 )?;
2168 for reference in &draft.origin.evidence {
2169 validate_proposal_evidence_reference(
2170 runtime,
2171 boundary_id,
2172 pending_evidence,
2173 reference,
2174 )?;
2175 }
2176 let existing = current
2177 .knowledge
2178 .records
2179 .get(holder)
2180 .into_iter()
2181 .flat_map(|records| records.iter())
2182 .chain(
2183 knowledge_overlay
2184 .get(holder)
2185 .into_iter()
2186 .flat_map(|records| records.iter()),
2187 )
2188 .collect::<BTreeMap<_, _>>();
2189 for related in draft.supersedes.iter().chain(&draft.contradicts) {
2190 let Some(related_record) = existing.get(related) else {
2191 return Err(CanwuError::new(
2192 ErrorCode::KnowledgeRecordNotFound,
2193 "knowledge relation does not resolve for the same holder at this cut",
2194 ));
2195 };
2196 if related_record.schema.kind != draft.schema.kind {
2197 return Err(CanwuError::new(
2198 ErrorCode::InvalidKnowledgeRecord,
2199 "knowledge supersession and contradiction cannot cross schema kinds",
2200 ));
2201 }
2202 }
2203 let encoded = serde_json::to_vec(&(holder, draft)).map_err(|error| {
2204 CanwuError::new(
2205 ErrorCode::InvalidKnowledgeRecord,
2206 format!("holder-scoped knowledge draft could not be encoded: {error}"),
2207 )
2208 })?;
2209 if !canonical_drafts.insert(encoded) {
2210 return Err(CanwuError::new(
2211 ErrorCode::InvalidKnowledgeRecord,
2212 "one system proposal contains a duplicate canonical knowledge draft",
2213 ));
2214 }
2215 }
2216 }
2217 BoundaryDirective::ScheduleIngress {
2218 after,
2219 packet_type,
2220 payload,
2221 affected,
2222 ..
2223 } => {
2224 let descriptor = plugins
2225 .ingress
2226 .get(&(plugin.to_owned(), packet_type.clone()))
2227 .ok_or_else(|| {
2228 CanwuError::new(
2229 ErrorCode::InvalidPayload,
2230 format!(
2231 "boundary system {plugin}.{} scheduled undeclared ingress type {packet_type}",
2232 contract.name
2233 ),
2234 )
2235 })?;
2236 if after.is_negative() || now.checked_add(*after).is_none() {
2237 return Err(CanwuError::new(
2238 ErrorCode::InvalidDuration,
2239 "boundary-generated ingress requires a nonnegative supported delay",
2240 ));
2241 }
2242 descriptor.payload_schema.validate(payload)?;
2243 if affected.iter().any(|entity| {
2244 !proposal_entity_identity_exists(
2245 current,
2246 &plugins.record_schemas,
2247 proposal,
2248 entity,
2249 )
2250 }) {
2251 return Err(CanwuError::new(
2252 ErrorCode::EntityNotFound,
2253 format!(
2254 "boundary system {plugin}.{} scheduled ingress for an unknown entity identity",
2255 contract.name
2256 ),
2257 ));
2258 }
2259 }
2260 BoundaryDirective::SchedulePluginIngress {
2261 target_plugin,
2262 after,
2263 packet_type,
2264 payload,
2265 affected,
2266 ..
2267 } => {
2268 let grant = super::PluginIngressTarget {
2269 target_plugin: target_plugin.clone(),
2270 packet_type: packet_type.clone(),
2271 };
2272 if !contract.plugin_ingress_targets.contains(&grant) {
2273 return Err(CanwuError::new(
2274 ErrorCode::UndeclaredStateWrite,
2275 format!(
2276 "boundary system {plugin}.{} did not declare target ingress {target_plugin}.{packet_type}",
2277 contract.name
2278 ),
2279 ));
2280 }
2281 let descriptor = plugins
2282 .ingress
2283 .get(&(target_plugin.clone(), packet_type.clone()))
2284 .ok_or_else(|| {
2285 CanwuError::new(
2286 ErrorCode::InvalidPayload,
2287 format!(
2288 "boundary system {plugin}.{} scheduled undeclared target ingress {target_plugin}.{packet_type}",
2289 contract.name
2290 ),
2291 )
2292 })?;
2293 if after.is_negative() || now.checked_add(*after).is_none() {
2294 return Err(CanwuError::new(
2295 ErrorCode::InvalidDuration,
2296 "boundary-generated cross-plugin ingress requires a nonnegative supported delay",
2297 ));
2298 }
2299 descriptor.payload_schema.validate(payload)?;
2300 if affected.iter().any(|entity| {
2301 !proposal_entity_identity_exists(
2302 current,
2303 &plugins.record_schemas,
2304 proposal,
2305 entity,
2306 )
2307 }) {
2308 return Err(CanwuError::new(
2309 ErrorCode::EntityNotFound,
2310 format!(
2311 "boundary system {plugin}.{} scheduled cross-plugin ingress for an unknown entity identity",
2312 contract.name
2313 ),
2314 ));
2315 }
2316 }
2317 BoundaryDirective::CancelPluginIngress { ingress_id, reason } => {
2318 if !valid_ingress_cancellation_reason(reason)
2319 || !cancelled_ingress.insert(*ingress_id)
2320 {
2321 return Err(CanwuError::new(
2322 ErrorCode::InvalidPayload,
2323 format!(
2324 "boundary system {plugin}.{} proposed a duplicate or malformed ingress cancellation",
2325 contract.name
2326 ),
2327 ));
2328 }
2329 }
2330 BoundaryDirective::ResolveDecisionRandomly { resolution } => {
2331 validate_random_decision_resolution(
2332 plugin,
2333 contract,
2334 current,
2335 committed_availability,
2336 pending_random_draws,
2337 &mut random_decision_samples,
2338 resolution,
2339 )?;
2340 }
2341 BoundaryDirective::SetPersonAvailability {
2342 person,
2343 availability,
2344 summary,
2345 } => {
2346 super::persons::validate_availability_directive(
2347 plugin,
2348 contract,
2349 now,
2350 *person,
2351 availability,
2352 summary,
2353 &entity_exists,
2354 )?;
2355 }
2356 BoundaryDirective::RecordEvaluationTrace { trace } => {
2357 super::evaluation::validate_trace_shape(
2358 contract.phase,
2359 trace,
2360 boundary_id,
2361 runtime.metadata.run_configuration.evaluation_limits(),
2362 )?;
2363 if !proposal_entity_identity_exists(
2364 current,
2365 &plugins.record_schemas,
2366 proposal,
2367 &trace.subject,
2368 ) {
2369 return Err(CanwuError::new(
2370 ErrorCode::EntityNotFound,
2371 format!(
2372 "boundary system {plugin}.{} traced an evaluation of unknown subject {}",
2373 contract.name, trace.subject
2374 ),
2375 )
2376 .with_entity(trace.subject.clone()));
2377 }
2378 for reference in trace.terms.iter().flat_map(|term| &term.evidence) {
2379 validate_proposal_evidence_reference(
2380 runtime,
2381 boundary_id,
2382 pending_evidence,
2383 reference,
2384 )?;
2385 }
2386 }
2387 BoundaryDirective::CreatePerson {
2388 draft,
2389 correlation,
2390 summary,
2391 } => {
2392 super::persons::validate_person_draft(
2393 &super::persons::PersonDraftContext {
2394 plugin,
2395 contract,
2396 now,
2397 government_exists: &|id| current.governments.contains_key(&id),
2398 territory_exists: &|id| current.territories.contains_key(&id),
2399 entity_exists: &entity_exists,
2400 },
2401 draft,
2402 correlation,
2403 summary,
2404 )?;
2405 validate_proposal_evidence_reference(
2406 runtime,
2407 boundary_id,
2408 pending_evidence,
2409 &draft.provenance,
2410 )?;
2411 }
2412 BoundaryDirective::RegisterTransitionManifest { .. }
2413 | BoundaryDirective::StageTransitionWrite { .. } => {
2414 return Err(transition_directive_not_admitted(plugin, &contract.name));
2415 }
2416 }
2417 }
2418 Ok(())
2419}
2420
2421fn transition_directive_not_admitted(plugin: &str, system: &str) -> CanwuError {
2424 CanwuError::new(
2425 ErrorCode::InvalidBoundary,
2426 format!(
2427 "boundary system {plugin}.{system} produced a transition directive outside the transition ledger"
2428 ),
2429 )
2430}
2431
2432fn validate_random_decision_resolution(
2446 plugin: &str,
2447 contract: &BoundarySystemContract,
2448 current: &RuntimeCurrentState,
2449 committed_availability: &BTreeMap<super::PersonId, super::PersonAvailability>,
2450 pending_random_draws: &[random::PendingRandomDraw],
2451 used_samples: &mut BTreeSet<(super::RandomStreamKey, RandomDrawAddress)>,
2452 resolution: &super::RandomDecisionResolution,
2453) -> Result<(), CanwuError> {
2454 if resolution.decision_request_id.get() == 0
2455 || resolution
2456 .command_request_id
2457 .is_some_and(|request_id| request_id.get() == 0)
2458 || resolution.expected_version == 0
2459 {
2460 return Err(CanwuError::new(
2461 ErrorCode::InvalidDecision,
2462 "random decision resolution requires nonzero request IDs and ticket version",
2463 ));
2464 }
2465 let ticket = current
2466 .decisions
2467 .ticket(resolution.ticket_id)
2468 .ok_or_else(|| {
2469 CanwuError::new(
2470 ErrorCode::InvalidDecision,
2471 "random decision resolution references an unknown ticket",
2472 )
2473 })?;
2474 if !ticket.is_open()
2475 || ticket.version != resolution.expected_version
2476 || ticket.assigned_controller != resolution.controller_id
2477 {
2478 return Err(CanwuError::new(
2479 ErrorCode::InvalidDecision,
2480 "random decision resolution references a closed, stale, or differently controlled ticket",
2481 ));
2482 }
2483 let controller = current
2484 .decisions
2485 .controller(&resolution.controller_id)
2486 .ok_or_else(|| {
2487 CanwuError::new(
2488 ErrorCode::InvalidDecision,
2489 "random decision resolution references an unknown controller",
2490 )
2491 })?;
2492 super::persons::validate_decision_preparation(committed_availability, ticket, controller)?;
2493 match &resolution.tie_break {
2494 None if controller.policy.kind != DecisionPolicyKind::Random => {
2495 return Err(CanwuError::new(
2496 ErrorCode::InvalidDecision,
2497 "random decision resolution requires a controller with random policy identity",
2498 ));
2499 }
2500 None => {}
2501 Some(pending) => validate_random_tie_break(controller, ticket, resolution, pending)?,
2502 }
2503 if !contract.random_streams.contains(&resolution.sample.stream) {
2504 return Err(CanwuError::new(
2505 ErrorCode::UndeclaredRandomStream,
2506 format!(
2507 "boundary system {plugin}.{} did not declare the random decision stream",
2508 contract.name
2509 ),
2510 ));
2511 }
2512 let RandomDrawAddress::OperationV1(address) = &resolution.sample.address else {
2513 return Err(CanwuError::new(
2514 ErrorCode::InvalidRandomDraw,
2515 "random decisions require an operation-keyed draw",
2516 ));
2517 };
2518 if address.producer_plugin != plugin
2519 || address.target
2520 != (RandomOperationTarget::DecisionTicket {
2521 ticket_id: ticket.id,
2522 ticket_version: ticket.version,
2523 })
2524 {
2525 return Err(CanwuError::new(
2526 ErrorCode::InvalidRandomDraw,
2527 "random decision draw address does not bind the current ticket version",
2528 ));
2529 }
2530 let sample_key = (
2531 resolution.sample.stream.clone(),
2532 resolution.sample.address.clone(),
2533 );
2534 if !used_samples.insert(sample_key.clone()) {
2535 return Err(CanwuError::new(
2536 ErrorCode::InvalidRandomDraw,
2537 "one random draw cannot resolve more than one decision",
2538 ));
2539 }
2540 if !pending_random_draws.iter().any(|draw| {
2541 draw.stream == sample_key.0
2542 && draw.address == sample_key.1
2543 && draw.upper_exclusive == resolution.sample.upper_exclusive
2544 && draw.value == resolution.sample.value
2545 }) {
2546 return Err(CanwuError::new(
2547 ErrorCode::InvalidRandomDraw,
2548 "random decision resolution does not reference a draw produced by this proposal",
2549 ));
2550 }
2551 let total_weight = resolution
2552 .option_weights
2553 .iter()
2554 .try_fold(0_u64, |total, option| total.checked_add(option.weight))
2555 .ok_or_else(|| {
2556 CanwuError::new(
2557 ErrorCode::InvalidDecision,
2558 "random decision option weights overflow the supported range",
2559 )
2560 })?;
2561 if total_weight != resolution.sample.upper_exclusive {
2562 return Err(CanwuError::new(
2563 ErrorCode::InvalidDecision,
2564 "random decision option weights disagree with the draw bound",
2565 ));
2566 }
2567 let selected = random_resolution_selection(ticket, resolution)?;
2568 let action = &ticket
2569 .option(&selected)
2570 .expect("validated random decision selected an existing option")
2571 .action;
2572 if matches!(action, DecisionAction::Command { .. }) != resolution.command_request_id.is_some() {
2573 return Err(CanwuError::new(
2574 ErrorCode::InvalidDecision,
2575 "random decision command options require exactly one command request ID",
2576 ));
2577 }
2578 Ok(())
2579}
2580
2581fn validate_random_tie_break(
2585 controller: &super::DecisionControllerBinding,
2586 ticket: &super::DecisionTicket,
2587 resolution: &super::RandomDecisionResolution,
2588 pending: &PolicyDecision,
2589) -> Result<(), CanwuError> {
2590 if controller.policy.kind != DecisionPolicyKind::Utility || !controller.random_tie_break {
2591 return Err(CanwuError::new(
2592 ErrorCode::InvalidDecision,
2593 "a random tie-break requires a utility-policy controller that permits tie-breaks",
2594 ));
2595 }
2596 let DecisionOutcome::PendingRandom { candidates } = &pending.outcome else {
2597 return Err(CanwuError::new(
2598 ErrorCode::InvalidDecision,
2599 "a random tie-break must carry a pending random policy decision",
2600 ));
2601 };
2602 if candidates != &resolution.option_weights
2603 || !pending.is_random_tie_break()
2604 || pending.external.is_some()
2605 || pending.random.is_some()
2606 {
2607 return Err(CanwuError::new(
2608 ErrorCode::InvalidDecision,
2609 "random tie-break weights must equal the pending candidates of an evidence-free random stage",
2610 ));
2611 }
2612 pending
2613 .validate(ticket)
2614 .map_err(|error| CanwuError::new(ErrorCode::InvalidDecision, error.to_string()))
2615}
2616
2617pub(super) fn random_resolution_selection(
2621 ticket: &super::DecisionTicket,
2622 resolution: &super::RandomDecisionResolution,
2623) -> Result<String, CanwuError> {
2624 if resolution.tie_break.is_some() {
2625 DecisionRandomEvidence::selected_candidate(
2626 ticket,
2627 &resolution.option_weights,
2628 resolution.sample.value,
2629 )
2630 } else {
2631 DecisionRandomEvidence::selected_option(
2632 ticket,
2633 &resolution.option_weights,
2634 resolution.sample.value,
2635 )
2636 }
2637 .map_err(|error| CanwuError::new(ErrorCode::InvalidDecision, error.to_string()))
2638}
2639
2640pub(super) fn random_policy_summary(option_id: &str) -> String {
2641 format!("random policy selected {option_id}")
2642}
2643
2644pub(super) fn random_tie_break_summary(option_id: &str) -> String {
2645 format!("random tie-break selected {option_id}")
2646}
2647
2648fn validate_proposal_evidence_reference(
2649 runtime: &RuntimeState,
2650 boundary_id: BoundaryId,
2651 pending: &PendingBoundaryEvidence,
2652 reference: &EvidenceRef,
2653) -> Result<(), CanwuError> {
2654 if let EvidenceRef::DomainRecordVersion(version) = reference
2655 && let DomainRecordVersionSource::BoundaryChange {
2656 boundary,
2657 change_index,
2658 } = version.established_by
2659 && boundary == boundary_id
2660 {
2661 let resolved = usize::try_from(change_index)
2662 .ok()
2663 .and_then(|index| pending.record_changes.get(index))
2664 .is_some_and(|change| {
2665 change.current.reference == version.record
2666 && change.current.version == version.version
2667 });
2668 return if resolved {
2669 Ok(())
2670 } else {
2671 Err(CanwuError::new(
2672 ErrorCode::EvidenceUnavailable,
2673 "knowledge origin references an unavailable current-boundary record version",
2674 ))
2675 };
2676 }
2677
2678 if let EvidenceRef::Event(id) = reference
2679 && runtime
2680 .evidence
2681 .retained_event(*id)
2682 .is_some_and(|event| event.cause == Some(CauseRef::Boundary(boundary_id)))
2683 {
2684 if pending
2685 .emissions
2686 .iter()
2687 .any(|emission| emission.event == *id)
2688 {
2689 return Ok(());
2690 }
2691 return Err(CanwuError::new(
2692 ErrorCode::EvidenceUnavailable,
2693 "knowledge origin references an event outside the proposal-visible boundary cut",
2694 ));
2695 }
2696
2697 if let EvidenceRef::Ingress(id) = reference
2698 && runtime
2699 .evidence
2700 .retained_ingress(*id)
2701 .is_some_and(|record| record.cause == Some(CauseRef::Boundary(boundary_id)))
2702 {
2703 return Err(CanwuError::new(
2704 ErrorCode::EvidenceUnavailable,
2705 "current-boundary generated ingress is not proposal-visible evidence",
2706 ));
2707 }
2708
2709 match resolve_evidence_reference(&RuntimeValidationContext::new(runtime), reference) {
2710 EvidenceAvailability::Retained | EvidenceAvailability::Archived => Ok(()),
2711 EvidenceAvailability::Missing => Err(CanwuError::new(
2712 ErrorCode::EvidenceUnavailable,
2713 "knowledge origin references missing or wrong-version evidence",
2714 )),
2715 }
2716}
2717
2718fn validate_reservation_pool(
2719 pool: &ReservationPoolKey,
2720 entity_exists: &dyn Fn(&EntityRef) -> bool,
2721) -> Result<(), CanwuError> {
2722 if pool.resource.trim().is_empty()
2723 || pool.resource != pool.resource.trim()
2724 || !entity_exists(&pool.entity)
2725 {
2726 return Err(CanwuError::new(
2727 ErrorCode::InvalidBoundary,
2728 "reservation pools require a canonical resource and an existing entity",
2729 ));
2730 }
2731 Ok(())
2732}
2733
2734fn extend_boundary_overlay(
2735 current: &RuntimeCurrentState,
2736 record_overlay: &BTreeMap<DomainRecordRef, DomainRecord>,
2737 overlay: &mut BTreeMap<PluginComponentKey, PluginComponentRecord>,
2738 directives: &[StagedBoundaryDirective],
2739) -> Result<(), CanwuError> {
2740 extend_boundary_component_overlay(current, record_overlay, overlay, directives, false)
2741}
2742
2743fn extend_boundary_candidate_overlay(
2744 current: &RuntimeCurrentState,
2745 record_overlay: &BTreeMap<DomainRecordRef, DomainRecord>,
2746 overlay: &mut BTreeMap<PluginComponentKey, PluginComponentRecord>,
2747 directives: &[StagedBoundaryDirective],
2748) -> Result<(), CanwuError> {
2749 extend_boundary_component_overlay(current, record_overlay, overlay, directives, true)
2750}
2751
2752fn extend_boundary_component_overlay(
2753 current: &RuntimeCurrentState,
2754 record_overlay: &BTreeMap<DomainRecordRef, DomainRecord>,
2755 overlay: &mut BTreeMap<PluginComponentKey, PluginComponentRecord>,
2756 directives: &[StagedBoundaryDirective],
2757 include_next_boundary: bool,
2758) -> Result<(), CanwuError> {
2759 for staged in directives.iter().filter(|staged| {
2760 include_next_boundary || staged.visibility == StateVisibility::SameBoundary
2761 }) {
2762 if let BoundaryDirective::SetComponent {
2763 state: state_key,
2764 entity,
2765 component,
2766 value,
2767 ..
2768 } = &staged.directive
2769 {
2770 let key = component_key(&staged.plugin, state_key, entity, component);
2771 if overlay.contains_key(&key) {
2772 return Err(CanwuError::new(
2773 ErrorCode::InvalidBoundary,
2774 "multiple boundary proposals target the same component",
2775 ));
2776 }
2777 if !runtime_entity_exists_with_record_overlay(current, record_overlay, entity) {
2778 return Err(CanwuError::new(
2779 ErrorCode::EntityNotFound,
2780 format!("boundary proposal targeted missing entity {entity}"),
2781 ));
2782 }
2783 overlay.insert(
2784 key,
2785 PluginComponentRecord {
2786 plugin: staged.plugin.clone(),
2787 state: state_key.clone(),
2788 entity: entity.clone(),
2789 component: component.clone(),
2790 value: value.clone(),
2791 },
2792 );
2793 }
2794 }
2795 Ok(())
2796}
2797
2798fn extend_boundary_record_overlay(
2799 context: &BoundaryRecordOverlayContext<'_>,
2800 overlay: &mut BTreeMap<DomainRecordRef, DomainRecord>,
2801 directives: &[StagedBoundaryDirective],
2802) -> Result<(), CanwuError> {
2803 extend_boundary_domain_record_overlay(context, overlay, directives, false)
2804}
2805
2806fn extend_boundary_record_candidate_overlay(
2807 context: &BoundaryRecordOverlayContext<'_>,
2808 overlay: &mut BTreeMap<DomainRecordRef, DomainRecord>,
2809 directives: &[StagedBoundaryDirective],
2810) -> Result<(), CanwuError> {
2811 extend_boundary_domain_record_overlay(context, overlay, directives, true)
2812}
2813
2814struct BoundaryRecordOverlayContext<'a> {
2815 current: &'a RuntimeCurrentState,
2816 now: SimTime,
2817 scheduled_actions: &'a BTreeMap<ScheduleKey, ScheduledAction>,
2818 run_configuration: &'a RunConfigurationSnapshot,
2819 schemas: &'a records::DomainRecordSchemas,
2820}
2821
2822fn extend_boundary_domain_record_overlay(
2823 context: &BoundaryRecordOverlayContext<'_>,
2824 overlay: &mut BTreeMap<DomainRecordRef, DomainRecord>,
2825 directives: &[StagedBoundaryDirective],
2826 include_next_boundary: bool,
2827) -> Result<(), CanwuError> {
2828 let requests: Vec<_> = directives
2829 .iter()
2830 .filter(|staged| {
2831 include_next_boundary || staged.visibility == StateVisibility::SameBoundary
2832 })
2833 .filter_map(|staged| match &staged.directive {
2834 BoundaryDirective::MutateRecord { mutation, summary } => {
2835 Some(records::DomainMutationRequest {
2836 plugin: &staged.plugin,
2837 system: &staged.system,
2838 visibility: staged.visibility,
2839 mutation,
2840 summary,
2841 })
2842 }
2843 BoundaryDirective::SetComponent { .. }
2844 | BoundaryDirective::Emit { .. }
2845 | BoundaryDirective::ScheduleIngress { .. }
2846 | BoundaryDirective::SchedulePluginIngress { .. }
2847 | BoundaryDirective::ResolveDecisionRandomly { .. }
2848 | BoundaryDirective::PublishKnowledge { .. }
2849 | BoundaryDirective::SetPersonAvailability { .. }
2850 | BoundaryDirective::CreatePerson { .. }
2851 | BoundaryDirective::CancelPluginIngress { .. }
2852 | BoundaryDirective::RecordEvaluationTrace { .. }
2853 | BoundaryDirective::RegisterTransitionManifest { .. }
2854 | BoundaryDirective::StageTransitionWrite { .. } => None,
2855 })
2856 .collect();
2857 if requests.is_empty() {
2858 return Ok(());
2859 }
2860 let (next, changes) = records::apply_mutation_bundle_cow_with_overlay(
2861 &context.current.domain_records,
2862 overlay,
2863 context.schemas,
2864 context.now,
2865 &|entity| runtime_current_entity_exists(context.current, entity),
2866 requests,
2867 )?;
2868 validate_domain_dependents_with_records(
2869 &context.current.plugin_components,
2870 context.scheduled_actions,
2871 context.run_configuration,
2872 &next,
2873 )?;
2874 for change in changes {
2875 overlay.insert(change.current.reference.clone(), change.current);
2876 }
2877 Ok(())
2878}
2879
2880fn partition_boundary_visibility(
2881 directives: Vec<StagedBoundaryDirective>,
2882) -> (Vec<StagedBoundaryDirective>, Vec<StagedBoundaryDirective>) {
2883 directives
2884 .into_iter()
2885 .partition(|staged| staged.visibility == StateVisibility::SameBoundary)
2886}
2887
2888fn partition_knowledge_directives(
2889 directives: Vec<StagedBoundaryDirective>,
2890) -> (Vec<StagedBoundaryDirective>, Vec<StagedBoundaryDirective>) {
2891 directives
2892 .into_iter()
2893 .partition(|staged| matches!(staged.directive, BoundaryDirective::PublishKnowledge { .. }))
2894}
2895
2896fn allocate_reservations(
2897 mut offers: Vec<PendingReservationOffer>,
2898 mut requests: Vec<PendingReservationRequest>,
2899) -> Result<ReservationAllocationResult, CanwuError> {
2900 offers.sort_by(|left, right| {
2901 left.offer
2902 .pool
2903 .cmp(&right.offer.pool)
2904 .then_with(|| left.plugin.cmp(&right.plugin))
2905 .then_with(|| left.system.cmp(&right.system))
2906 });
2907 let mut remaining = BTreeMap::new();
2908 let mut offer_records = Vec::new();
2909 for pending in offers {
2910 if remaining
2911 .insert(pending.offer.pool.clone(), pending.offer.capacity)
2912 .is_some()
2913 {
2914 return Err(CanwuError::new(
2915 ErrorCode::InvalidBoundary,
2916 format!(
2917 "reservation pool was offered more than once, including by {}.{}",
2918 pending.plugin, pending.system
2919 ),
2920 ));
2921 }
2922 offer_records.push(ReservationOfferRecord {
2923 plugin: pending.plugin,
2924 system: pending.system,
2925 offer: pending.offer,
2926 });
2927 }
2928 requests.sort_by(|left, right| {
2929 left.request
2930 .pool
2931 .cmp(&right.request.pool)
2932 .then_with(|| right.request.priority.cmp(&left.request.priority))
2933 .then_with(|| left.request.tie_break.cmp(&right.request.tie_break))
2934 .then_with(|| left.reservation.cmp(&right.reservation))
2935 });
2936 let mut seen = BTreeSet::new();
2937 let mut by_reservation = BTreeMap::new();
2938 let mut request_records = Vec::new();
2939 let mut records = Vec::new();
2940 for pending in requests {
2941 if !seen.insert(pending.reservation.clone()) {
2942 return Err(CanwuError::new(
2943 ErrorCode::InvalidBoundary,
2944 "reservation request identity is duplicated",
2945 ));
2946 }
2947 request_records.push(ReservationRequestRecord {
2948 reservation: pending.reservation.clone(),
2949 request: pending.request.clone(),
2950 });
2951 let available = remaining.entry(pending.request.pool.clone()).or_default();
2952 let granted = pending.request.quantity.min(*available);
2953 *available -= granted;
2954 let disposition = if granted == pending.request.quantity {
2955 ReservationDisposition::Fulfilled
2956 } else if granted == 0 {
2957 ReservationDisposition::Rejected
2958 } else {
2959 ReservationDisposition::Partial
2960 };
2961 let allocation = ReservationAllocation {
2962 reservation: pending.reservation.clone(),
2963 pool: pending.request.pool,
2964 requested: pending.request.quantity,
2965 granted,
2966 remaining_after: *available,
2967 disposition,
2968 };
2969 by_reservation.insert(pending.reservation, allocation.clone());
2970 records.push(allocation);
2971 }
2972 Ok(ReservationAllocationResult {
2973 by_reservation,
2974 offers: offer_records,
2975 requests: request_records,
2976 records,
2977 })
2978}