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