1use super::validation::{
2 EvidenceAvailability, RuntimeValidationContext, resolve_evidence_reference,
3};
4use super::{
5 AssertUnwindSafe, BTreeMap, BTreeSet, BoundaryChange, BoundaryContext, BoundaryDirective,
6 BoundaryEmission, BoundaryEmissionKind, BoundaryId, BoundaryIngressGeneration,
7 BoundaryKnowledgeChange, BoundaryPhase, BoundaryProposal, BoundaryReceipt, BoundaryRecord,
8 BoundaryRequest, BoundaryStateHashFormat, BoundarySystemContract,
9 BoundaryTransactionCheckpoint, CanwuError, CauseRef, CommandIngress, CommandRequest,
10 CommitmentDomains, DomainRecord, DomainRecordChange, DomainRecordRef,
11 DomainRecordVersionSource, EntityRef, ErrorCode, EventKind, EvidenceRef, GENESIS_BOUNDARY_HASH,
12 HashSet, IngressPayload, KnowledgeHolderRef, KnowledgeRecord, KnowledgeRecordId,
13 PluginComponentKey, PluginComponentRecord, PluginRegistry, RefCell, ReservationAllocation,
14 ReservationDisposition, ReservationOffer, ReservationOfferRecord, ReservationPoolKey,
15 ReservationRef, ReservationRequest, ReservationRequestRecord, RunConfigurationSnapshot,
16 RuntimeCurrentState, RuntimeState, ScheduleKey, ScheduledAction, SimTime, Simulation,
17 SimulationView, SimulationViewState, StateKey, StateVisibility, SystemCadence, SystemDirective,
18 canonical_text, catch_unwind, claim_counter, component_key, compute_boundary_hash,
19 invalid_snapshot_error, is_domain_record_state, proposal_entity_exists,
20 proposal_entity_identity_exists, random, record_change_affected_entities, records,
21 runtime_current_entity_exists, runtime_entity_exists,
22 runtime_entity_exists_with_record_overlay, runtime_entity_identity_exists,
23 validate_domain_dependents_with_records, validate_runtime_domain_dependents,
24};
25
26impl Simulation {
27 pub fn settle_boundary(
28 &mut self,
29 request: BoundaryRequest,
30 ) -> Result<BoundaryReceipt, CanwuError> {
31 self.settle_boundary_with_state_hash_format(request, BoundaryStateHashFormat::CommitmentsV1)
32 }
33
34 pub(super) fn settle_boundary_with_state_hash_format(
35 &mut self,
36 mut request: BoundaryRequest,
37 state_hash_format: BoundaryStateHashFormat,
38 ) -> Result<BoundaryReceipt, CanwuError> {
39 self.ensure_runtime_ready()?;
40 if request.at < self.state.scheduler.now {
41 return Err(CanwuError::new(
42 ErrorCode::InvalidBoundary,
43 "a settlement boundary cannot precede committed simulation time",
44 ));
45 }
46 if self
47 .state
48 .scheduler
49 .pending_ingress
50 .first()
51 .is_some_and(|key| key.due_at < request.at)
52 {
53 return Err(CanwuError::new(
54 ErrorCode::InvalidBoundary,
55 "a settlement boundary cannot step past earlier canonical ingress",
56 ));
57 }
58 if request.cadences.contains(&SystemCadence::EventDriven) {
59 return Err(CanwuError::new(
60 ErrorCode::InvalidBoundary,
61 "event-driven cadence is derived from admitted events, not caller supplied",
62 ));
63 }
64 request.cadences.sort();
65 request.cadences.dedup();
66
67 let transaction = BoundaryTransactionCheckpoint::capture(&self.state);
68 match self.settle_boundary_inner(request, state_hash_format) {
69 Ok(receipt) => Ok(receipt),
70 Err(error) => {
71 transaction.restore(&mut self.state);
72 Err(error)
73 }
74 }
75 }
76
77 fn settle_boundary_inner(
78 &mut self,
79 mut request: BoundaryRequest,
80 state_hash_format: BoundaryStateHashFormat,
81 ) -> Result<BoundaryReceipt, CanwuError> {
82 self.advance_to_before_boundary(request.at)?;
83
84 let admitted_ingress = self.take_due_ingress(request.at);
85 let admitted_ingress_index: HashSet<_> = admitted_ingress.iter().copied().collect();
86 for ingress_id in &admitted_ingress {
87 let record = self
88 .state
89 .evidence
90 .retained_ingress(*ingress_id)
91 .cloned()
92 .ok_or_else(|| {
93 CanwuError::new(
94 ErrorCode::InvalidSnapshot,
95 "pending ingress references an unknown record",
96 )
97 })?;
98 match record.payload {
99 IngressPayload::Command { request: command } => {
100 let CommandRequest {
101 request_id,
102 expected_revision,
103 envelope,
104 } = *command;
105 self.admit_command(
106 Some(request_id),
107 Some(expected_revision),
108 envelope,
109 CommandIngress::LiveRequest,
110 None,
111 true,
112 )?;
113 }
114 IngressPayload::Calendar { cadences } => request.cadences.extend(cadences),
115 IngressPayload::Plugin { .. } => {}
116 IngressPayload::Decision { request } => {
117 self.apply_decision_request(*request)?;
118 }
119 }
120 }
121 self.state
122 .current
123 .decisions
124 .advance_time(request.at)
125 .map_err(super::decision::decision_error)?;
126 self.invalidate_commitments(CommitmentDomains::DECISIONS);
127 self.execute_scheduled_at(request.at)?;
128 request.cadences.sort();
129 request.cadences.dedup();
130
131 let admitted_attempt_count = self
132 .state
133 .evidence
134 .archived
135 .command_attempt_count
136 .checked_add(
137 u64::try_from(self.state.evidence.command_attempts.len()).map_err(|_| {
138 invalid_snapshot_error("attempt journal exceeds admission cursor range")
139 })?,
140 )
141 .ok_or_else(|| invalid_snapshot_error("attempt journal cursor is exhausted"))?;
142 let admitted_command_count = self
143 .state
144 .evidence
145 .archived
146 .command_count
147 .checked_add(
148 u64::try_from(self.state.evidence.commands.len()).map_err(|_| {
149 invalid_snapshot_error("command journal exceeds admission cursor range")
150 })?,
151 )
152 .ok_or_else(|| invalid_snapshot_error("command journal cursor is exhausted"))?;
153 let admitted_event_count = self
154 .state
155 .evidence
156 .archived
157 .event_count
158 .checked_add(
159 u64::try_from(self.state.evidence.events.len()).map_err(|_| {
160 invalid_snapshot_error("event journal exceeds admission cursor range")
161 })?,
162 )
163 .ok_or_else(|| invalid_snapshot_error("event journal cursor is exhausted"))?;
164 let admitted_attempt_start = self
165 .state
166 .counters
167 .admitted_attempt_count
168 .checked_sub(self.state.evidence.archived.command_attempt_count)
169 .ok_or_else(|| {
170 invalid_snapshot_error("runtime attempt admission cursor precedes live evidence")
171 })?;
172 let admitted_attempt_start = usize::try_from(admitted_attempt_start).map_err(|_| {
173 invalid_snapshot_error("runtime attempt admission cursor exceeds platform range")
174 })?;
175 let admitted_command_start = self
176 .state
177 .counters
178 .admitted_command_count
179 .checked_sub(self.state.evidence.archived.command_count)
180 .ok_or_else(|| {
181 invalid_snapshot_error("runtime command admission cursor precedes live evidence")
182 })?;
183 let admitted_command_start = usize::try_from(admitted_command_start).map_err(|_| {
184 invalid_snapshot_error("runtime command admission cursor exceeds platform range")
185 })?;
186 let admitted_event_start = self
187 .state
188 .counters
189 .admitted_event_count
190 .checked_sub(self.state.evidence.archived.event_count)
191 .ok_or_else(|| {
192 invalid_snapshot_error("runtime event admission cursor precedes live evidence")
193 })?;
194 let admitted_event_start = usize::try_from(admitted_event_start).map_err(|_| {
195 invalid_snapshot_error("runtime event admission cursor exceeds platform range")
196 })?;
197 let admitted_attempts: Vec<_> = self
198 .state
199 .evidence
200 .command_attempts
201 .get(admitted_attempt_start..)
202 .ok_or_else(|| {
203 invalid_snapshot_error("runtime attempt admission cursor exceeds its journal")
204 })?
205 .iter()
206 .map(|record| record.id)
207 .collect();
208 let admitted_commands: Vec<_> = self
209 .state
210 .evidence
211 .commands
212 .get(admitted_command_start..)
213 .ok_or_else(|| {
214 invalid_snapshot_error("runtime command admission cursor exceeds its journal")
215 })?
216 .iter()
217 .map(|record| record.id)
218 .collect();
219 let admitted_events: Vec<_> = self
220 .state
221 .evidence
222 .events
223 .get(admitted_event_start..)
224 .ok_or_else(|| {
225 invalid_snapshot_error("runtime event admission cursor exceeds its journal")
226 })?
227 .iter()
228 .map(|event| event.id)
229 .collect();
230
231 let (boundary_id_value, next_boundary_id) =
232 claim_counter(self.state.counters.next_boundary_id, "boundary ID")?;
233 let (correlation_id, next_correlation_id) = claim_counter(
234 self.state.counters.next_correlation_id,
235 "boundary correlation ID",
236 )?;
237 self.state.counters.next_boundary_id = next_boundary_id;
238 self.state.counters.next_correlation_id = next_correlation_id;
239 let boundary_id = BoundaryId::new(boundary_id_value);
240
241 let boundary_snapshot = self.state.current.clone();
242 let boundary_time = self.state.scheduler.now;
243 let systems = self.plugins.boundary_systems.clone();
244 let state_owners = self.plugins.state_owners.clone();
245 let record_schemas = self.plugins.record_schemas.clone();
246 let mut allocations = BTreeMap::new();
247 let mut allocation_records = Vec::new();
248 let mut reservation_offer_records = Vec::new();
249 let mut reservation_request_records = Vec::new();
250 let mut offers = Vec::new();
251 let mut requests = Vec::new();
252 let mut random_overlay = boundary_snapshot.random_streams.clone();
253 let mut pending_random_draws = Vec::new();
254 let mut keyed_random_draws = random::keyed_draws_with_reservations(
255 &self.state.evidence.random_draws,
256 &self.state.evidence.keyed_draw_reservations,
257 )?;
258 let mut visible_overlay = BTreeMap::new();
259 let mut candidate_overlay = BTreeMap::new();
260 let mut visible_record_overlay = BTreeMap::new();
261 let mut candidate_record_overlay = BTreeMap::new();
262 let mut visible_knowledge_overlay = BTreeMap::new();
263 let mut pending_knowledge_changes = Vec::new();
264 let mut knowledge_correlations = BTreeSet::new();
265 let mut ordinary = Vec::new();
266 let mut transitions = Vec::new();
267 let mut deferred = Vec::new();
268 let mut evidence = PendingBoundaryEvidence::default();
269
270 for phase in BoundaryPhase::ALL {
271 match phase {
272 BoundaryPhase::AtomicDomainCommit => {
273 let (same_boundary, next_boundary) =
274 partition_boundary_visibility(std::mem::take(&mut ordinary));
275 self.apply_boundary_stage(
276 boundary_id,
277 correlation_id,
278 same_boundary,
279 &mut evidence,
280 )?;
281 deferred.extend(next_boundary);
282 visible_overlay.clear();
283 candidate_overlay.clear();
284 visible_record_overlay.clear();
285 candidate_record_overlay.clear();
286 }
287 BoundaryPhase::ConditionalTransitionCommit => {
288 let (same_boundary, next_boundary) =
289 partition_boundary_visibility(std::mem::take(&mut transitions));
290 self.apply_boundary_stage(
291 boundary_id,
292 correlation_id,
293 same_boundary,
294 &mut evidence,
295 )?;
296 deferred.extend(next_boundary);
297 visible_overlay.clear();
298 visible_record_overlay.clear();
299 }
300 _ => {}
301 }
302
303 let mut phase_directives = Vec::new();
304 for registered in systems.iter().filter(|registered| {
305 registered.contract.phase == phase
306 && boundary_system_due(
307 ®istered.contract,
308 &request.cadences,
309 !admitted_events.is_empty() || !admitted_ingress.is_empty(),
310 )
311 }) {
312 let reader = format!("{}.{}", registered.plugin, registered.contract.name);
313 let (view_current, view_now) = if phase <= BoundaryPhase::InvariantValidation {
314 (&boundary_snapshot, boundary_time)
315 } else {
316 (&self.state.current, self.state.scheduler.now)
317 };
318 let random_session = random::RandomSession::new(
319 &random_overlay,
320 ®istered.contract.random_streams,
321 boundary_snapshot.root_seed,
322 ®istered.plugin,
323 &keyed_random_draws,
324 )?;
325 let proposal_evidence = proposal_evidence_refs(boundary_id, &evidence);
326 let view = SimulationView {
327 state: SimulationViewState::Boundary {
328 current: view_current,
329 now: view_now,
330 runtime: &self.state,
331 },
332 state_owners: &state_owners,
333 reader: Some(&reader),
334 allowed_reads: Some(®istered.contract.reads),
335 allowed_ingress: Some(&admitted_ingress_index),
336 ingress_plugin: Some(®istered.plugin),
337 component_overlay: Some(&visible_overlay),
338 proposed_components: (phase == BoundaryPhase::InvariantValidation)
339 .then_some(&candidate_overlay),
340 record_overlay: Some(&visible_record_overlay),
341 proposed_records: (phase == BoundaryPhase::InvariantValidation)
342 .then_some(&candidate_record_overlay),
343 boundary_id: Some(boundary_id),
344 proposal_evidence: Some(&proposal_evidence),
345 knowledge_overlay: Some(&visible_knowledge_overlay),
346 allocations: Some(&allocations),
347 allowed_reservations: Some(®istered.contract.reservation_reads),
348 random_session: Some(RefCell::new(random_session)),
349 };
350 let context = BoundaryContext {
351 boundary_id,
352 at: request.at,
353 phase,
354 plugin: registered.plugin.clone(),
355 system: registered.contract.name.clone(),
356 admitted_attempts: admitted_attempts.clone(),
357 admitted_commands: admitted_commands.clone(),
358 admitted_ingress: admitted_ingress.clone(),
359 admitted_events: admitted_events.clone(),
360 emitted_events: evidence
361 .emissions
362 .iter()
363 .map(|emission| emission.event)
364 .collect(),
365 };
366 let proposal =
367 catch_unwind(AssertUnwindSafe(|| (registered.handler)(&view, &context)))
368 .map_err(|_| {
369 CanwuError::new(
370 ErrorCode::PluginPanicked,
371 format!(
372 "boundary system {}.{} panicked",
373 registered.plugin, registered.contract.name
374 ),
375 )
376 })??;
377 validate_boundary_proposal(
378 ®istered.plugin,
379 ®istered.contract,
380 view_current,
381 view_now,
382 &self.state,
383 boundary_id,
384 &evidence,
385 &self.plugins,
386 &visible_record_overlay,
387 &visible_knowledge_overlay,
388 &proposal,
389 )?;
390 let random_execution = view
391 .finish_random_session()
392 .expect("boundary views always have a random session");
393 random::extend_keyed_draws(&mut keyed_random_draws, &random_execution.draws)?;
394 random_overlay.extend(random_execution.states);
395 pending_random_draws.extend(random_execution.draws.into_iter().map(|draw| {
396 PendingBoundaryRandomDraw {
397 plugin: registered.plugin.clone(),
398 system: registered.contract.name.clone(),
399 draw,
400 }
401 }));
402 offers.extend(
403 proposal
404 .offers
405 .into_iter()
406 .map(|offer| PendingReservationOffer {
407 plugin: registered.plugin.clone(),
408 system: registered.contract.name.clone(),
409 offer,
410 }),
411 );
412 requests.extend(proposal.requests.into_iter().map(|request| {
413 PendingReservationRequest {
414 reservation: ReservationRef::new(
415 ®istered.plugin,
416 ®istered.contract.name,
417 &request.request,
418 ),
419 request,
420 }
421 }));
422 phase_directives.extend(proposal.directives.into_iter().map(|directive| {
423 StagedBoundaryDirective {
424 plugin: registered.plugin.clone(),
425 system: registered.contract.name.clone(),
426 phase,
427 visibility: registered.contract.visibility,
428 directive,
429 }
430 }));
431 }
432
433 let (knowledge_directives, phase_directives) =
434 partition_knowledge_directives(phase_directives);
435 if !knowledge_directives.is_empty()
436 && !matches!(
437 phase,
438 BoundaryPhase::PerceptionAndAttentionRefresh
439 | BoundaryPhase::PerspectiveAndReportMaterialization
440 )
441 {
442 return Err(CanwuError::new(
443 ErrorCode::UndeclaredKnowledgeWrite,
444 "knowledge publication is allowed only in phases 4 and 13",
445 ));
446 }
447
448 match phase {
449 BoundaryPhase::PerceptionAndAttentionRefresh => {
450 if !phase_directives.is_empty() {
451 return Err(CanwuError::new(
452 ErrorCode::InvalidBoundary,
453 "phase 4 accepts knowledge publications but no ordinary directives",
454 ));
455 }
456 self.stage_knowledge_publications(
457 phase,
458 knowledge_directives,
459 &mut visible_knowledge_overlay,
460 &mut pending_knowledge_changes,
461 &mut knowledge_correlations,
462 )?;
463 }
464 BoundaryPhase::ReservationAndAllocation => {
465 let result = allocate_reservations(
466 std::mem::take(&mut offers),
467 std::mem::take(&mut requests),
468 )?;
469 allocations = result.by_reservation;
470 allocation_records = result.records;
471 reservation_offer_records = result.offers;
472 reservation_request_records = result.requests;
473 }
474 BoundaryPhase::DomainDeltaProposal => {
475 let record_context = BoundaryRecordOverlayContext {
476 current: &boundary_snapshot,
477 now: boundary_time,
478 scheduled_actions: &self.state.scheduler.actions,
479 run_configuration: &self.state.metadata.run_configuration,
480 schemas: &record_schemas,
481 };
482 extend_boundary_record_candidate_overlay(
483 &record_context,
484 &mut candidate_record_overlay,
485 &phase_directives,
486 )?;
487 extend_boundary_candidate_overlay(
488 &boundary_snapshot,
489 &candidate_record_overlay,
490 &mut candidate_overlay,
491 &phase_directives,
492 )?;
493 extend_boundary_record_overlay(
494 &record_context,
495 &mut visible_record_overlay,
496 &phase_directives,
497 )?;
498 extend_boundary_overlay(
499 &boundary_snapshot,
500 &visible_record_overlay,
501 &mut visible_overlay,
502 &phase_directives,
503 )?;
504 ordinary.extend(phase_directives);
505 }
506 BoundaryPhase::HistoricalCandidateEvaluation => {
507 let record_context = BoundaryRecordOverlayContext {
508 current: &self.state.current,
509 now: self.state.scheduler.now,
510 scheduled_actions: &self.state.scheduler.actions,
511 run_configuration: &self.state.metadata.run_configuration,
512 schemas: &record_schemas,
513 };
514 extend_boundary_record_overlay(
515 &record_context,
516 &mut visible_record_overlay,
517 &phase_directives,
518 )?;
519 extend_boundary_overlay(
520 &self.state.current,
521 &visible_record_overlay,
522 &mut visible_overlay,
523 &phase_directives,
524 )?;
525 transitions.extend(phase_directives);
526 }
527 BoundaryPhase::StrategicAggregation
528 | BoundaryPhase::PerspectiveAndReportMaterialization => {
529 if phase == BoundaryPhase::PerspectiveAndReportMaterialization {
530 self.stage_knowledge_publications(
531 phase,
532 knowledge_directives,
533 &mut visible_knowledge_overlay,
534 &mut pending_knowledge_changes,
535 &mut knowledge_correlations,
536 )?;
537 }
538 let (same_boundary, next_boundary) =
539 partition_boundary_visibility(phase_directives);
540 self.apply_boundary_stage(
541 boundary_id,
542 correlation_id,
543 same_boundary,
544 &mut evidence,
545 )?;
546 deferred.extend(next_boundary);
547 }
548 _ if !phase_directives.is_empty() => {
549 return Err(CanwuError::new(
550 ErrorCode::InvalidBoundary,
551 format!("boundary phase {phase:?} cannot produce state directives"),
552 ));
553 }
554 _ => {}
555 }
556 }
557
558 self.apply_boundary_stage(boundary_id, correlation_id, deferred, &mut evidence)?;
559 self.commit_knowledge_publications(
560 boundary_id,
561 correlation_id,
562 &pending_knowledge_changes,
563 &mut evidence.emissions,
564 )?;
565 let PendingBoundaryEvidence {
566 changes,
567 record_changes,
568 emissions,
569 generated_ingress,
570 } = evidence;
571 self.invalidate_commitments(CommitmentDomains::RANDOM_STREAMS);
572 self.state.current.random_streams = random_overlay;
573 let random_draws =
574 self.append_boundary_random_draws(boundary_id, correlation_id, pending_random_draws)?;
575 self.state.metadata.plugin_registration_closed = true;
576 let state_hash = self.compute_boundary_state_hash_for(state_hash_format)?;
577 let previous_hash = self
578 .state
579 .evidence
580 .boundary_head_hash()
581 .map_or_else(|| GENESIS_BOUNDARY_HASH.to_owned(), str::to_owned);
582 let mut record = BoundaryRecord {
583 id: boundary_id,
584 at: request.at,
585 correlation_id,
586 cadences: request.cadences,
587 admitted_attempts,
588 admitted_commands,
589 admitted_ingress,
590 generated_ingress: generated_ingress.clone(),
591 admitted_events,
592 reservation_offers: reservation_offer_records,
593 reservation_requests: reservation_request_records,
594 allocations: allocation_records.clone(),
595 random_draws: random_draws.clone(),
596 changes: changes.clone(),
597 record_changes: record_changes.clone(),
598 knowledge_changes: pending_knowledge_changes.clone(),
599 emissions: emissions.clone(),
600 state_hash: Some(state_hash),
601 previous_hash,
602 hash: String::new(),
603 };
604 record.hash = compute_boundary_hash(&record)?;
605 let boundary_hash = record.hash.clone();
606 self.state.evidence.boundaries.push(record);
607 self.state.counters.admitted_attempt_count = admitted_attempt_count;
608 self.state.counters.admitted_command_count = admitted_command_count;
609 self.state.counters.admitted_event_count = admitted_event_count;
610 self.advance_state_revision()?;
611 self.refresh_checkpoint_hash()?;
612 Ok(BoundaryReceipt {
613 boundary_id,
614 settled_at: request.at,
615 emitted_events: emissions
616 .into_iter()
617 .map(|emission| emission.event)
618 .collect(),
619 generated_ingress: generated_ingress
620 .into_iter()
621 .map(|generation| generation.ingress)
622 .collect(),
623 random_draws,
624 boundary_hash,
625 change_count: changes.len(),
626 record_change_count: record_changes.len(),
627 knowledge_batch_count: pending_knowledge_changes.len(),
628 knowledge_record_count: pending_knowledge_changes
629 .iter()
630 .map(|change| change.records.len())
631 .sum(),
632 allocations: allocation_records,
633 })
634 }
635
636 fn apply_boundary_stage(
637 &mut self,
638 boundary_id: BoundaryId,
639 correlation_id: u64,
640 directives: Vec<StagedBoundaryDirective>,
641 evidence: &mut PendingBoundaryEvidence,
642 ) -> Result<(), CanwuError> {
643 let changes = &mut evidence.changes;
644 let record_changes = &mut evidence.record_changes;
645 let emissions = &mut evidence.emissions;
646 let generated_ingress = &mut evidence.generated_ingress;
647 let mutation_requests: Vec<_> = directives
648 .iter()
649 .filter_map(|staged| match &staged.directive {
650 BoundaryDirective::MutateRecord { mutation, summary } => {
651 Some(records::DomainMutationRequest {
652 plugin: &staged.plugin,
653 system: &staged.system,
654 visibility: staged.visibility,
655 mutation,
656 summary,
657 })
658 }
659 BoundaryDirective::SetComponent { .. }
660 | BoundaryDirective::Emit { .. }
661 | BoundaryDirective::ScheduleIngress { .. }
662 | BoundaryDirective::SchedulePluginIngress { .. }
663 | BoundaryDirective::PublishKnowledge { .. } => None,
664 })
665 .collect();
666 let mut stage_record_changes = BTreeMap::new();
667 if !mutation_requests.is_empty() {
668 let (next_records, applied) = records::apply_mutation_bundle(
669 &self.state.current.domain_records,
670 &self.plugins.record_schemas,
671 self.state.scheduler.now,
672 &|entity| runtime_entity_exists(&self.state, entity),
673 mutation_requests,
674 )?;
675 let first_index = record_changes.len();
676 for (offset, change) in applied.iter().enumerate() {
677 let index = first_index.checked_add(offset).ok_or_else(|| {
678 CanwuError::new(
679 ErrorCode::IdentifierExhausted,
680 "boundary record-change index exceeds the persistent identifier space",
681 )
682 })?;
683 let index = u64::try_from(index).map_err(|_| {
684 CanwuError::new(
685 ErrorCode::IdentifierExhausted,
686 "boundary record-change index exceeds the persistent identifier space",
687 )
688 })?;
689 stage_record_changes
690 .insert(change.current.reference.clone(), (index, change.clone()));
691 }
692 self.invalidate_commitments(CommitmentDomains::DOMAIN_RECORDS);
693 self.state.current.domain_records = next_records;
694 record_changes.extend(applied);
695 }
696
697 for staged in &directives {
698 let unavailable = match &staged.directive {
699 BoundaryDirective::SetComponent { entity, .. } => {
700 (!runtime_entity_exists(&self.state, entity)).then_some(entity)
701 }
702 BoundaryDirective::Emit { affected, .. } => affected
703 .iter()
704 .find(|entity| !runtime_entity_exists(&self.state, entity)),
705 BoundaryDirective::ScheduleIngress { affected, .. } => affected
706 .iter()
707 .find(|entity| !runtime_entity_identity_exists(&self.state, entity)),
708 BoundaryDirective::SchedulePluginIngress { affected, .. } => affected
709 .iter()
710 .find(|entity| !runtime_entity_identity_exists(&self.state, entity)),
711 BoundaryDirective::MutateRecord { .. }
712 | BoundaryDirective::PublishKnowledge { .. } => None,
713 };
714 if let Some(entity) = unavailable {
715 return Err(CanwuError::new(
716 ErrorCode::EntityNotFound,
717 format!(
718 "boundary stage {}.{} references unavailable entity {entity}",
719 staged.plugin, staged.system
720 ),
721 )
722 .with_entity(entity.clone()));
723 }
724 }
725
726 for staged in directives {
727 match staged.directive {
728 BoundaryDirective::SetComponent {
729 state,
730 entity,
731 component,
732 value,
733 summary,
734 } => {
735 let key = component_key(&staged.plugin, &state, &entity, &component);
736 self.invalidate_commitments(CommitmentDomains::PLUGIN_COMPONENTS);
737 let previous = self
738 .state
739 .current
740 .plugin_components
741 .get(&key)
742 .map(|record| record.value.clone());
743 self.state.current.plugin_components.insert(
744 key,
745 PluginComponentRecord {
746 plugin: staged.plugin.clone(),
747 state: state.clone(),
748 entity: entity.clone(),
749 component: component.clone(),
750 value: value.clone(),
751 },
752 );
753 let change_index = u64::try_from(changes.len()).map_err(|_| {
754 CanwuError::new(
755 ErrorCode::IdentifierExhausted,
756 "boundary change index exceeds the persistent identifier space",
757 )
758 })?;
759 changes.push(BoundaryChange {
760 plugin: staged.plugin.clone(),
761 system: staged.system.clone(),
762 state,
763 entity: entity.clone(),
764 component: component.clone(),
765 previous,
766 value,
767 visibility: staged.visibility,
768 summary: summary.clone(),
769 });
770 let event = self.append_event(
771 EventKind::Plugin {
772 plugin: staged.plugin.clone(),
773 event_type: format!("{component}_changed"),
774 },
775 vec![entity],
776 summary,
777 Some(CauseRef::Boundary(boundary_id)),
778 correlation_id,
779 )?;
780 emissions.push(BoundaryEmission {
781 plugin: staged.plugin,
782 system: staged.system,
783 event: event.id,
784 kind: BoundaryEmissionKind::Change { change_index },
785 });
786 }
787 BoundaryDirective::MutateRecord { mutation, .. } => {
788 let Some((change_index, change)) = stage_record_changes.get(mutation.target())
789 else {
790 return Err(CanwuError::new(
791 ErrorCode::InvalidBoundary,
792 "record mutation is missing its committed change evidence",
793 ));
794 };
795 let event = self.append_event(
796 EventKind::Plugin {
797 plugin: staged.plugin.clone(),
798 event_type: change.operation.event_type().to_owned(),
799 },
800 record_change_affected_entities(change),
801 change.summary.clone(),
802 Some(CauseRef::Boundary(boundary_id)),
803 correlation_id,
804 )?;
805 emissions.push(BoundaryEmission {
806 plugin: staged.plugin,
807 system: staged.system,
808 event: event.id,
809 kind: BoundaryEmissionKind::RecordChange {
810 change_index: *change_index,
811 },
812 });
813 }
814 BoundaryDirective::Emit {
815 event_type,
816 summary,
817 affected,
818 } => {
819 let event = self.append_event(
820 EventKind::Plugin {
821 plugin: staged.plugin.clone(),
822 event_type,
823 },
824 affected,
825 summary,
826 Some(CauseRef::Boundary(boundary_id)),
827 correlation_id,
828 )?;
829 emissions.push(BoundaryEmission {
830 plugin: staged.plugin,
831 system: staged.system,
832 event: event.id,
833 kind: BoundaryEmissionKind::Explicit,
834 });
835 }
836 BoundaryDirective::PublishKnowledge { .. } => {
837 return Err(CanwuError::new(
838 ErrorCode::InvalidBoundary,
839 "knowledge publication execution is not enabled in this runtime slice",
840 ));
841 }
842 BoundaryDirective::ScheduleIngress {
843 after,
844 packet_type,
845 priority,
846 payload,
847 mut affected,
848 } => {
849 self.ensure_canonical_ingress_can_start()?;
850 let descriptor = self
851 .plugins
852 .ingress
853 .get(&(staged.plugin.clone(), packet_type.clone()))
854 .ok_or_else(|| {
855 CanwuError::new(
856 ErrorCode::InvalidPayload,
857 format!(
858 "boundary system {}.{} scheduled undeclared ingress type {packet_type}",
859 staged.plugin, staged.system
860 ),
861 )
862 })?
863 .clone();
864 descriptor.payload_schema.validate(&payload)?;
865 affected.sort();
866 affected.dedup();
867 let due_at = self.state.scheduler.now.checked_add(after).ok_or_else(|| {
868 CanwuError::new(
869 ErrorCode::InvalidDuration,
870 "boundary-generated ingress exceeds the supported time range",
871 )
872 })?;
873 let receipt = self.append_ingress(
874 due_at,
875 descriptor.class,
876 priority,
877 IngressPayload::Plugin {
878 plugin: staged.plugin.clone(),
879 packet_type,
880 payload,
881 affected_entities: affected,
882 },
883 Some(CauseRef::Boundary(boundary_id)),
884 true,
885 )?;
886 generated_ingress.push(BoundaryIngressGeneration {
887 ingress: receipt.ingress_id,
888 plugin: staged.plugin,
889 system: staged.system,
890 phase: staged.phase,
891 visibility: staged.visibility,
892 });
893 }
894 BoundaryDirective::SchedulePluginIngress {
895 target_plugin,
896 after,
897 packet_type,
898 priority,
899 payload,
900 mut affected,
901 } => {
902 self.ensure_canonical_ingress_can_start()?;
903 let descriptor = self
904 .plugins
905 .ingress
906 .get(&(target_plugin.clone(), packet_type.clone()))
907 .ok_or_else(|| {
908 CanwuError::new(
909 ErrorCode::InvalidPayload,
910 format!(
911 "boundary system {}.{} scheduled undeclared target ingress {}.{packet_type}",
912 staged.plugin, staged.system, target_plugin
913 ),
914 )
915 })?
916 .clone();
917 descriptor.payload_schema.validate(&payload)?;
918 affected.sort();
919 affected.dedup();
920 let due_at = self.state.scheduler.now.checked_add(after).ok_or_else(|| {
921 CanwuError::new(
922 ErrorCode::InvalidDuration,
923 "boundary-generated cross-plugin ingress exceeds the supported time range",
924 )
925 })?;
926 let receipt = self.append_ingress(
927 due_at,
928 descriptor.class,
929 priority,
930 IngressPayload::Plugin {
931 plugin: target_plugin,
932 packet_type,
933 payload,
934 affected_entities: affected,
935 },
936 Some(CauseRef::Boundary(boundary_id)),
937 true,
938 )?;
939 generated_ingress.push(BoundaryIngressGeneration {
940 ingress: receipt.ingress_id,
941 plugin: staged.plugin,
942 system: staged.system,
943 phase: staged.phase,
944 visibility: staged.visibility,
945 });
946 }
947 }
948 }
949 validate_runtime_domain_dependents(&self.state)?;
950 Ok(())
951 }
952
953 fn stage_knowledge_publications(
954 &mut self,
955 phase: BoundaryPhase,
956 directives: Vec<StagedBoundaryDirective>,
957 visible_overlay: &mut BTreeMap<
958 KnowledgeHolderRef,
959 BTreeMap<KnowledgeRecordId, KnowledgeRecord>,
960 >,
961 pending: &mut Vec<BoundaryKnowledgeChange>,
962 correlations: &mut BTreeSet<(String, String, String)>,
963 ) -> Result<(), CanwuError> {
964 let new_record_count = directives
965 .iter()
966 .map(|staged| match &staged.directive {
967 BoundaryDirective::PublishKnowledge { records, .. } => records.len(),
968 _ => 0,
969 })
970 .sum::<usize>();
971 let total_records = pending
972 .iter()
973 .map(|change| change.records.len())
974 .sum::<usize>()
975 .checked_add(new_record_count)
976 .ok_or_else(|| {
977 CanwuError::new(
978 ErrorCode::KnowledgeLimitExceeded,
979 "boundary knowledge record count exceeds platform range",
980 )
981 })?;
982 if total_records > crate::KnowledgeLimitsV1::CURRENT.records_per_boundary {
983 return Err(CanwuError::new(
984 ErrorCode::KnowledgeLimitExceeded,
985 "boundary knowledge record limit exceeded",
986 ));
987 }
988 for staged in directives {
989 let BoundaryDirective::PublishKnowledge {
990 holder,
991 visibility,
992 producer_correlation,
993 records: drafts,
994 summary,
995 } = staged.directive
996 else {
997 return Err(CanwuError::new(
998 ErrorCode::InvalidBoundary,
999 "knowledge stage received an ordinary directive",
1000 ));
1001 };
1002 if let Some(value) = &producer_correlation
1003 && !correlations.insert((
1004 staged.plugin.clone(),
1005 staged.system.clone(),
1006 value.clone(),
1007 ))
1008 {
1009 return Err(CanwuError::new(
1010 ErrorCode::InvalidKnowledgeRecord,
1011 "producer correlation is duplicated within one system and boundary",
1012 ));
1013 }
1014 let mut records = Vec::with_capacity(drafts.len());
1015 for draft in drafts {
1016 let (id, next_id) = claim_counter(
1017 self.state.counters.next_knowledge_record_id,
1018 "knowledge record ID",
1019 )?;
1020 self.state.counters.next_knowledge_record_id = next_id;
1021 let record = KnowledgeRecord {
1022 id: KnowledgeRecordId::new(id),
1023 holder: holder.clone(),
1024 schema: draft.schema,
1025 subjects: draft.subjects,
1026 payload: draft.payload,
1027 as_of: draft.as_of,
1028 learned_at: self.state.scheduler.now,
1029 confidence_per_mille: draft.confidence_per_mille,
1030 origin: draft.origin,
1031 supersedes: draft.supersedes,
1032 contradicts: draft.contradicts,
1033 };
1034 if visibility == StateVisibility::SameBoundary {
1035 visible_overlay
1036 .entry(holder.clone())
1037 .or_default()
1038 .insert(record.id, record.clone());
1039 }
1040 records.push(record);
1041 }
1042 pending.push(BoundaryKnowledgeChange {
1043 plugin: staged.plugin,
1044 system: staged.system,
1045 phase,
1046 holder,
1047 producer_correlation,
1048 records,
1049 visibility,
1050 summary,
1051 });
1052 }
1053 Ok(())
1054 }
1055
1056 fn commit_knowledge_publications(
1057 &mut self,
1058 boundary_id: BoundaryId,
1059 correlation_id: u64,
1060 changes: &[BoundaryKnowledgeChange],
1061 emissions: &mut Vec<BoundaryEmission>,
1062 ) -> Result<(), CanwuError> {
1063 if changes.is_empty() {
1064 return Ok(());
1065 }
1066 let mut ledger = self.state.current.knowledge.records.clone();
1067 for change in changes {
1068 let holder = ledger.entry(change.holder.clone()).or_default();
1069 for record in &change.records {
1070 if holder.insert(record.id, record.clone()).is_some() {
1071 return Err(CanwuError::new(
1072 ErrorCode::InvalidKnowledgeRecord,
1073 "knowledge publication attempted to reuse a global record ID",
1074 ));
1075 }
1076 }
1077 }
1078 self.state.current.knowledge.records = ledger;
1079 self.invalidate_commitments(CommitmentDomains::KNOWLEDGE);
1080 for (index, change) in changes.iter().enumerate() {
1081 let record_count = u32::try_from(change.records.len()).map_err(|_| {
1082 CanwuError::new(
1083 ErrorCode::KnowledgeLimitExceeded,
1084 "knowledge publication event count exceeds u32",
1085 )
1086 })?;
1087 let affected = match &change.holder {
1088 KnowledgeHolderRef::Person(person) => vec![EntityRef::Person(*person)],
1089 KnowledgeHolderRef::Entity(entity) => vec![entity.clone()],
1090 };
1091 let event = self.append_event(
1092 EventKind::KnowledgePublished {
1093 holder: change.holder.clone(),
1094 record_count,
1095 },
1096 affected,
1097 change.summary.clone(),
1098 Some(CauseRef::Boundary(boundary_id)),
1099 correlation_id,
1100 )?;
1101 emissions.push(BoundaryEmission {
1102 plugin: change.plugin.clone(),
1103 system: change.system.clone(),
1104 event: event.id,
1105 kind: BoundaryEmissionKind::KnowledgeChange {
1106 change_index: u64::try_from(index).map_err(|_| {
1107 CanwuError::new(
1108 ErrorCode::IdentifierExhausted,
1109 "knowledge change index exceeds identifier space",
1110 )
1111 })?,
1112 },
1113 });
1114 }
1115 Ok(())
1116 }
1117
1118 pub(super) fn apply_directives(
1119 &mut self,
1120 plugin: &str,
1121 directives: Vec<SystemDirective>,
1122 allowed_writes: &[StateKey],
1123 cause: &CauseRef,
1124 correlation_id: u64,
1125 ) -> Result<(), CanwuError> {
1126 for directive in directives {
1127 match directive {
1128 SystemDirective::SetComponent {
1129 state,
1130 entity,
1131 component,
1132 value,
1133 summary,
1134 } => {
1135 let key = component_key(plugin, &state, &entity, &component);
1136 self.invalidate_commitments(CommitmentDomains::PLUGIN_COMPONENTS);
1137 self.state.current.plugin_components.insert(
1138 key,
1139 PluginComponentRecord {
1140 plugin: plugin.to_owned(),
1141 state,
1142 entity: entity.clone(),
1143 component: component.clone(),
1144 value,
1145 },
1146 );
1147 self.emit(
1148 EventKind::Plugin {
1149 plugin: plugin.to_owned(),
1150 event_type: format!("{component}_changed"),
1151 },
1152 vec![entity],
1153 summary,
1154 Some(cause.clone()),
1155 correlation_id,
1156 )?;
1157 }
1158 SystemDirective::Emit {
1159 event_type,
1160 summary,
1161 affected,
1162 } => {
1163 self.emit(
1164 EventKind::Plugin {
1165 plugin: plugin.to_owned(),
1166 event_type,
1167 },
1168 affected,
1169 summary,
1170 Some(cause.clone()),
1171 correlation_id,
1172 )?;
1173 }
1174 SystemDirective::Schedule { after, directive } => {
1175 let at = self.state.scheduler.now.checked_add(after).ok_or_else(|| {
1176 CanwuError::new(
1177 ErrorCode::InvalidDuration,
1178 "plugin scheduled time exceeds the supported range",
1179 )
1180 })?;
1181 self.schedule_at(
1182 at,
1183 ScheduledAction::PluginDirective {
1184 plugin: plugin.to_owned(),
1185 directive,
1186 allowed_writes: allowed_writes.to_vec(),
1187 cause: cause.clone(),
1188 correlation_id,
1189 },
1190 )?;
1191 }
1192 SystemDirective::EnqueuePluginIngress {
1193 after,
1194 packet_type,
1195 priority,
1196 payload,
1197 mut affected,
1198 } => {
1199 self.ensure_canonical_ingress_can_start()?;
1200 let descriptor = self
1201 .plugins
1202 .ingress
1203 .get(&(plugin.to_owned(), packet_type.clone()))
1204 .ok_or_else(|| {
1205 CanwuError::new(
1206 ErrorCode::InvalidPayload,
1207 format!(
1208 "plugin command scheduled unregistered ingress type {plugin}.{packet_type}"
1209 ),
1210 )
1211 })?
1212 .clone();
1213 descriptor.payload_schema.validate(&payload)?;
1214 affected.sort();
1215 affected.dedup();
1216 if affected
1217 .iter()
1218 .any(|entity| !runtime_entity_identity_exists(&self.state, entity))
1219 {
1220 return Err(CanwuError::new(
1221 ErrorCode::EntityNotFound,
1222 "plugin command ingress references an unknown entity identity",
1223 ));
1224 }
1225 let due_at = self.state.scheduler.now.checked_add(after).ok_or_else(|| {
1226 CanwuError::new(
1227 ErrorCode::InvalidDuration,
1228 "plugin command ingress exceeds the supported time range",
1229 )
1230 })?;
1231 self.append_ingress(
1232 due_at,
1233 descriptor.class,
1234 priority,
1235 IngressPayload::Plugin {
1236 plugin: plugin.to_owned(),
1237 packet_type,
1238 payload,
1239 affected_entities: affected,
1240 },
1241 Some(cause.clone()),
1242 true,
1243 )?;
1244 }
1245 }
1246 }
1247 Ok(())
1248 }
1249}
1250
1251struct PendingReservationOffer {
1252 plugin: String,
1253 system: String,
1254 offer: ReservationOffer,
1255}
1256
1257struct PendingReservationRequest {
1258 reservation: ReservationRef,
1259 request: ReservationRequest,
1260}
1261
1262struct ReservationAllocationResult {
1263 by_reservation: BTreeMap<ReservationRef, ReservationAllocation>,
1264 offers: Vec<ReservationOfferRecord>,
1265 requests: Vec<ReservationRequestRecord>,
1266 records: Vec<ReservationAllocation>,
1267}
1268
1269struct StagedBoundaryDirective {
1270 plugin: String,
1271 system: String,
1272 phase: BoundaryPhase,
1273 visibility: StateVisibility,
1274 directive: BoundaryDirective,
1275}
1276
1277#[derive(Default)]
1278struct PendingBoundaryEvidence {
1279 changes: Vec<BoundaryChange>,
1280 record_changes: Vec<DomainRecordChange>,
1281 emissions: Vec<BoundaryEmission>,
1282 generated_ingress: Vec<BoundaryIngressGeneration>,
1283}
1284
1285fn proposal_evidence_refs(
1286 boundary: BoundaryId,
1287 pending: &PendingBoundaryEvidence,
1288) -> BTreeSet<EvidenceRef> {
1289 let mut values = BTreeSet::new();
1290 for (index, change) in pending.record_changes.iter().enumerate() {
1291 if let Ok(change_index) = u64::try_from(index) {
1292 values.insert(EvidenceRef::DomainRecordVersion(
1293 super::DomainRecordVersionRef {
1294 record: change.current.reference.clone(),
1295 version: change.current.version,
1296 established_by: DomainRecordVersionSource::BoundaryChange {
1297 boundary,
1298 change_index,
1299 },
1300 },
1301 ));
1302 }
1303 }
1304 values.extend(
1305 pending
1306 .emissions
1307 .iter()
1308 .map(|emission| EvidenceRef::Event(emission.event)),
1309 );
1310 values
1311}
1312
1313pub(super) struct PendingBoundaryRandomDraw {
1314 pub(super) plugin: String,
1315 pub(super) system: String,
1316 pub(super) draw: random::PendingRandomDraw,
1317}
1318
1319pub(super) fn boundary_system_due(
1320 contract: &BoundarySystemContract,
1321 cadences: &[SystemCadence],
1322 has_admitted_events: bool,
1323) -> bool {
1324 match contract.cadence {
1325 SystemCadence::EventDriven => has_admitted_events,
1326 _ => cadences.contains(&contract.cadence),
1327 }
1328}
1329
1330pub(super) fn boundary_has_event_ingress(record: &BoundaryRecord) -> bool {
1331 !record.admitted_events.is_empty() || !record.admitted_ingress.is_empty()
1332}
1333
1334#[allow(clippy::too_many_arguments)]
1335fn validate_boundary_proposal(
1336 plugin: &str,
1337 contract: &BoundarySystemContract,
1338 current: &RuntimeCurrentState,
1339 now: SimTime,
1340 runtime: &RuntimeState,
1341 boundary_id: BoundaryId,
1342 pending_evidence: &PendingBoundaryEvidence,
1343 plugins: &PluginRegistry,
1344 record_overlay: &BTreeMap<DomainRecordRef, DomainRecord>,
1345 knowledge_overlay: &BTreeMap<KnowledgeHolderRef, BTreeMap<KnowledgeRecordId, KnowledgeRecord>>,
1346 proposal: &BoundaryProposal,
1347) -> Result<(), CanwuError> {
1348 if contract.phase != BoundaryPhase::ReservationAndAllocation
1349 && (!proposal.offers.is_empty() || !proposal.requests.is_empty())
1350 {
1351 return Err(CanwuError::new(
1352 ErrorCode::InvalidBoundary,
1353 format!(
1354 "boundary system {plugin}.{} proposed reservations in phase {:?}",
1355 contract.name, contract.phase
1356 ),
1357 ));
1358 }
1359
1360 let entity_exists = |entity: &EntityRef| {
1361 proposal_entity_exists(
1362 current,
1363 &plugins.record_schemas,
1364 record_overlay,
1365 proposal,
1366 entity,
1367 )
1368 };
1369 let mut offered_pools = BTreeSet::new();
1370 for offer in &proposal.offers {
1371 validate_reservation_pool(&offer.pool, &entity_exists)?;
1372 if !contract.reservation_offers.contains(&offer.pool.state)
1373 || plugins
1374 .state_owners
1375 .get(&offer.pool.state)
1376 .is_none_or(|owner| owner != plugin)
1377 {
1378 return Err(CanwuError::new(
1379 ErrorCode::InvalidBoundary,
1380 format!(
1381 "boundary system {plugin}.{} offered undeclared state {}.{}",
1382 contract.name, offer.pool.state.namespace, offer.pool.state.name
1383 ),
1384 ));
1385 }
1386 if !offered_pools.insert(&offer.pool) {
1387 return Err(CanwuError::new(
1388 ErrorCode::InvalidBoundary,
1389 format!(
1390 "boundary system {plugin}.{} offered the same reservation pool twice",
1391 contract.name
1392 ),
1393 ));
1394 }
1395 }
1396
1397 let mut request_names = BTreeSet::new();
1398 for request in &proposal.requests {
1399 validate_reservation_pool(&request.pool, &entity_exists)?;
1400 if request.request.trim().is_empty()
1401 || request.request != request.request.trim()
1402 || request.tie_break.trim().is_empty()
1403 || request.tie_break != request.tie_break.trim()
1404 || request.quantity == 0
1405 || !request_names.insert(&request.request)
1406 || !contract.reservation_requests.contains(&request.pool.state)
1407 {
1408 return Err(CanwuError::new(
1409 ErrorCode::InvalidBoundary,
1410 format!(
1411 "boundary system {plugin}.{} produced an invalid reservation request",
1412 contract.name
1413 ),
1414 ));
1415 }
1416 }
1417
1418 let publication_count = proposal
1419 .directives
1420 .iter()
1421 .filter(|directive| matches!(directive, BoundaryDirective::PublishKnowledge { .. }))
1422 .count();
1423 if publication_count > crate::KnowledgeLimitsV1::CURRENT.batches_per_system_boundary {
1424 return Err(CanwuError::new(
1425 ErrorCode::KnowledgeLimitExceeded,
1426 "system knowledge publication batch limit exceeded",
1427 ));
1428 }
1429 let mut component_keys = BTreeSet::new();
1430 let mut record_targets = BTreeSet::new();
1431 let mut producer_correlations = BTreeSet::new();
1432 let mut canonical_drafts = BTreeSet::new();
1433 for directive in &proposal.directives {
1434 match directive {
1435 BoundaryDirective::SetComponent {
1436 state: state_key,
1437 entity,
1438 component,
1439 ..
1440 } => {
1441 if component.trim().is_empty()
1442 || component != component.trim()
1443 || !contract.writes.contains(state_key)
1444 || plugins
1445 .state_owners
1446 .get(state_key)
1447 .is_none_or(|owner| owner != plugin)
1448 || is_domain_record_state(&plugins.record_schemas, state_key)
1449 {
1450 return Err(CanwuError::new(
1451 ErrorCode::UndeclaredStateWrite,
1452 format!(
1453 "boundary system {plugin}.{} produced an undeclared component write",
1454 contract.name
1455 ),
1456 ));
1457 }
1458 if !entity_exists(entity) {
1459 return Err(CanwuError::new(
1460 ErrorCode::EntityNotFound,
1461 format!(
1462 "boundary system {plugin}.{} targeted missing entity {entity}",
1463 contract.name
1464 ),
1465 )
1466 .with_entity(entity.clone()));
1467 }
1468 let key = component_key(plugin, state_key, entity, component);
1469 if !component_keys.insert(key) {
1470 return Err(CanwuError::new(
1471 ErrorCode::InvalidBoundary,
1472 format!(
1473 "boundary system {plugin}.{} wrote the same component twice",
1474 contract.name
1475 ),
1476 ));
1477 }
1478 }
1479 BoundaryDirective::MutateRecord { mutation, summary } => {
1480 let target = mutation.target();
1481 let state_key = records::record_state_key(&target.kind);
1482 if !canonical_text(summary)
1483 || !contract.writes.contains(&state_key)
1484 || plugins
1485 .state_owners
1486 .get(&state_key)
1487 .is_none_or(|owner| owner != plugin)
1488 || plugins
1489 .record_schemas
1490 .get(&target.kind)
1491 .is_none_or(|(owner, _)| owner != plugin)
1492 {
1493 return Err(CanwuError::new(
1494 ErrorCode::UndeclaredStateWrite,
1495 format!(
1496 "boundary system {plugin}.{} produced an undeclared record mutation",
1497 contract.name
1498 ),
1499 ));
1500 }
1501 if !record_targets.insert(target.clone()) {
1502 return Err(CanwuError::new(
1503 ErrorCode::InvalidBoundary,
1504 format!(
1505 "boundary system {plugin}.{} mutated the same record twice",
1506 contract.name
1507 ),
1508 ));
1509 }
1510 }
1511 BoundaryDirective::Emit {
1512 event_type,
1513 affected,
1514 ..
1515 } => {
1516 if event_type.trim().is_empty()
1517 || event_type != event_type.trim()
1518 || !contract.emits.contains(event_type)
1519 {
1520 return Err(CanwuError::new(
1521 ErrorCode::InvalidBoundary,
1522 format!(
1523 "boundary system {plugin}.{} emitted an undeclared event type",
1524 contract.name
1525 ),
1526 ));
1527 }
1528 if affected.iter().any(|entity| !entity_exists(entity)) {
1529 return Err(CanwuError::new(
1530 ErrorCode::EntityNotFound,
1531 format!(
1532 "boundary system {plugin}.{} emitted an event for a missing entity",
1533 contract.name
1534 ),
1535 ));
1536 }
1537 }
1538 BoundaryDirective::PublishKnowledge {
1539 holder,
1540 visibility,
1541 producer_correlation,
1542 records,
1543 summary,
1544 } => {
1545 if !matches!(
1546 contract.phase,
1547 BoundaryPhase::PerceptionAndAttentionRefresh
1548 | BoundaryPhase::PerspectiveAndReportMaterialization
1549 ) {
1550 return Err(CanwuError::new(
1551 ErrorCode::UndeclaredKnowledgeWrite,
1552 "knowledge publication is allowed only in phases 4 and 13",
1553 ));
1554 }
1555 if records.is_empty()
1556 || records.len() > crate::KnowledgeLimitsV1::CURRENT.records_per_batch
1557 {
1558 return Err(CanwuError::new(
1559 ErrorCode::KnowledgeLimitExceeded,
1560 "knowledge publication batch is empty or exceeds its record limit",
1561 ));
1562 }
1563 if !canonical_text(summary)
1564 || summary.len() > crate::KnowledgeLimitsV1::CURRENT.text_bytes
1565 {
1566 return Err(CanwuError::new(
1567 ErrorCode::InvalidKnowledgeRecord,
1568 "knowledge publication summary is not canonical or exceeds its limit",
1569 ));
1570 }
1571 if let Some(value) = producer_correlation
1572 && (!canonical_text(value)
1573 || value.len() > 256
1574 || !producer_correlations.insert(value))
1575 {
1576 return Err(CanwuError::new(
1577 ErrorCode::InvalidKnowledgeRecord,
1578 "producer correlation is invalid or duplicated",
1579 ));
1580 }
1581 for draft in records {
1582 let Some(grant) = contract
1583 .knowledge_writes
1584 .iter()
1585 .find(|grant| grant.schema == draft.schema)
1586 else {
1587 return Err(CanwuError::new(
1588 ErrorCode::UndeclaredKnowledgeWrite,
1589 format!(
1590 "boundary system {plugin}.{} did not declare the knowledge schema",
1591 contract.name
1592 ),
1593 ));
1594 };
1595 if !grant.visibilities.contains(visibility) {
1596 return Err(CanwuError::new(
1597 ErrorCode::UndeclaredKnowledgeWrite,
1598 "knowledge publication visibility is not granted",
1599 ));
1600 }
1601 let Some((owner, schema)) = plugins.knowledge_schemas.get(&draft.schema) else {
1602 return Err(CanwuError::new(
1603 ErrorCode::InvalidKnowledgeSchema,
1604 "knowledge publication uses an unregistered schema",
1605 ));
1606 };
1607 if owner != plugin || !schema.writable {
1608 return Err(CanwuError::new(
1609 ErrorCode::UndeclaredKnowledgeWrite,
1610 "knowledge publication uses a foreign or read-only schema",
1611 ));
1612 }
1613 super::knowledge::validate_draft(
1614 draft,
1615 schema,
1616 holder,
1617 current,
1618 &plugins.record_schemas,
1619 )?;
1620 for reference in &draft.origin.evidence {
1621 validate_proposal_evidence_reference(
1622 runtime,
1623 boundary_id,
1624 pending_evidence,
1625 reference,
1626 )?;
1627 }
1628 let existing = current
1629 .knowledge
1630 .records
1631 .get(holder)
1632 .into_iter()
1633 .flat_map(|records| records.iter())
1634 .chain(
1635 knowledge_overlay
1636 .get(holder)
1637 .into_iter()
1638 .flat_map(|records| records.iter()),
1639 )
1640 .collect::<BTreeMap<_, _>>();
1641 for related in draft.supersedes.iter().chain(&draft.contradicts) {
1642 let Some(related_record) = existing.get(related) else {
1643 return Err(CanwuError::new(
1644 ErrorCode::KnowledgeRecordNotFound,
1645 "knowledge relation does not resolve for the same holder at this cut",
1646 ));
1647 };
1648 if related_record.schema.kind != draft.schema.kind {
1649 return Err(CanwuError::new(
1650 ErrorCode::InvalidKnowledgeRecord,
1651 "knowledge supersession and contradiction cannot cross schema kinds",
1652 ));
1653 }
1654 }
1655 let encoded = serde_json::to_vec(&(holder, draft)).map_err(|error| {
1656 CanwuError::new(
1657 ErrorCode::InvalidKnowledgeRecord,
1658 format!("holder-scoped knowledge draft could not be encoded: {error}"),
1659 )
1660 })?;
1661 if !canonical_drafts.insert(encoded) {
1662 return Err(CanwuError::new(
1663 ErrorCode::InvalidKnowledgeRecord,
1664 "one system proposal contains a duplicate canonical knowledge draft",
1665 ));
1666 }
1667 }
1668 }
1669 BoundaryDirective::ScheduleIngress {
1670 after,
1671 packet_type,
1672 payload,
1673 affected,
1674 ..
1675 } => {
1676 let descriptor = plugins
1677 .ingress
1678 .get(&(plugin.to_owned(), packet_type.clone()))
1679 .ok_or_else(|| {
1680 CanwuError::new(
1681 ErrorCode::InvalidPayload,
1682 format!(
1683 "boundary system {plugin}.{} scheduled undeclared ingress type {packet_type}",
1684 contract.name
1685 ),
1686 )
1687 })?;
1688 if after.is_negative() || now.checked_add(*after).is_none() {
1689 return Err(CanwuError::new(
1690 ErrorCode::InvalidDuration,
1691 "boundary-generated ingress requires a nonnegative supported delay",
1692 ));
1693 }
1694 descriptor.payload_schema.validate(payload)?;
1695 if affected.iter().any(|entity| {
1696 !proposal_entity_identity_exists(
1697 current,
1698 &plugins.record_schemas,
1699 proposal,
1700 entity,
1701 )
1702 }) {
1703 return Err(CanwuError::new(
1704 ErrorCode::EntityNotFound,
1705 format!(
1706 "boundary system {plugin}.{} scheduled ingress for an unknown entity identity",
1707 contract.name
1708 ),
1709 ));
1710 }
1711 }
1712 BoundaryDirective::SchedulePluginIngress {
1713 target_plugin,
1714 after,
1715 packet_type,
1716 payload,
1717 affected,
1718 ..
1719 } => {
1720 let grant = super::PluginIngressTarget {
1721 target_plugin: target_plugin.clone(),
1722 packet_type: packet_type.clone(),
1723 };
1724 if !contract.plugin_ingress_targets.contains(&grant) {
1725 return Err(CanwuError::new(
1726 ErrorCode::UndeclaredStateWrite,
1727 format!(
1728 "boundary system {plugin}.{} did not declare target ingress {target_plugin}.{packet_type}",
1729 contract.name
1730 ),
1731 ));
1732 }
1733 let descriptor = plugins
1734 .ingress
1735 .get(&(target_plugin.clone(), packet_type.clone()))
1736 .ok_or_else(|| {
1737 CanwuError::new(
1738 ErrorCode::InvalidPayload,
1739 format!(
1740 "boundary system {plugin}.{} scheduled undeclared target ingress {target_plugin}.{packet_type}",
1741 contract.name
1742 ),
1743 )
1744 })?;
1745 if after.is_negative() || now.checked_add(*after).is_none() {
1746 return Err(CanwuError::new(
1747 ErrorCode::InvalidDuration,
1748 "boundary-generated cross-plugin ingress requires a nonnegative supported delay",
1749 ));
1750 }
1751 descriptor.payload_schema.validate(payload)?;
1752 if affected.iter().any(|entity| {
1753 !proposal_entity_identity_exists(
1754 current,
1755 &plugins.record_schemas,
1756 proposal,
1757 entity,
1758 )
1759 }) {
1760 return Err(CanwuError::new(
1761 ErrorCode::EntityNotFound,
1762 format!(
1763 "boundary system {plugin}.{} scheduled cross-plugin ingress for an unknown entity identity",
1764 contract.name
1765 ),
1766 ));
1767 }
1768 }
1769 }
1770 }
1771 Ok(())
1772}
1773
1774fn validate_proposal_evidence_reference(
1775 runtime: &RuntimeState,
1776 boundary_id: BoundaryId,
1777 pending: &PendingBoundaryEvidence,
1778 reference: &EvidenceRef,
1779) -> Result<(), CanwuError> {
1780 if let EvidenceRef::DomainRecordVersion(version) = reference
1781 && let DomainRecordVersionSource::BoundaryChange {
1782 boundary,
1783 change_index,
1784 } = version.established_by
1785 && boundary == boundary_id
1786 {
1787 let resolved = usize::try_from(change_index)
1788 .ok()
1789 .and_then(|index| pending.record_changes.get(index))
1790 .is_some_and(|change| {
1791 change.current.reference == version.record
1792 && change.current.version == version.version
1793 });
1794 return if resolved {
1795 Ok(())
1796 } else {
1797 Err(CanwuError::new(
1798 ErrorCode::EvidenceUnavailable,
1799 "knowledge origin references an unavailable current-boundary record version",
1800 ))
1801 };
1802 }
1803
1804 if let EvidenceRef::Event(id) = reference
1805 && runtime
1806 .evidence
1807 .retained_event(*id)
1808 .is_some_and(|event| event.cause == Some(CauseRef::Boundary(boundary_id)))
1809 {
1810 if pending
1811 .emissions
1812 .iter()
1813 .any(|emission| emission.event == *id)
1814 {
1815 return Ok(());
1816 }
1817 return Err(CanwuError::new(
1818 ErrorCode::EvidenceUnavailable,
1819 "knowledge origin references an event outside the proposal-visible boundary cut",
1820 ));
1821 }
1822
1823 if let EvidenceRef::Ingress(id) = reference
1824 && runtime
1825 .evidence
1826 .retained_ingress(*id)
1827 .is_some_and(|record| record.cause == Some(CauseRef::Boundary(boundary_id)))
1828 {
1829 return Err(CanwuError::new(
1830 ErrorCode::EvidenceUnavailable,
1831 "current-boundary generated ingress is not proposal-visible evidence",
1832 ));
1833 }
1834
1835 match resolve_evidence_reference(&RuntimeValidationContext::new(runtime), reference) {
1836 EvidenceAvailability::Retained | EvidenceAvailability::Archived => Ok(()),
1837 EvidenceAvailability::Missing => Err(CanwuError::new(
1838 ErrorCode::EvidenceUnavailable,
1839 "knowledge origin references missing or wrong-version evidence",
1840 )),
1841 }
1842}
1843
1844fn validate_reservation_pool(
1845 pool: &ReservationPoolKey,
1846 entity_exists: &dyn Fn(&EntityRef) -> bool,
1847) -> Result<(), CanwuError> {
1848 if pool.resource.trim().is_empty()
1849 || pool.resource != pool.resource.trim()
1850 || !entity_exists(&pool.entity)
1851 {
1852 return Err(CanwuError::new(
1853 ErrorCode::InvalidBoundary,
1854 "reservation pools require a canonical resource and an existing entity",
1855 ));
1856 }
1857 Ok(())
1858}
1859
1860fn extend_boundary_overlay(
1861 current: &RuntimeCurrentState,
1862 record_overlay: &BTreeMap<DomainRecordRef, DomainRecord>,
1863 overlay: &mut BTreeMap<PluginComponentKey, PluginComponentRecord>,
1864 directives: &[StagedBoundaryDirective],
1865) -> Result<(), CanwuError> {
1866 extend_boundary_component_overlay(current, record_overlay, overlay, directives, false)
1867}
1868
1869fn extend_boundary_candidate_overlay(
1870 current: &RuntimeCurrentState,
1871 record_overlay: &BTreeMap<DomainRecordRef, DomainRecord>,
1872 overlay: &mut BTreeMap<PluginComponentKey, PluginComponentRecord>,
1873 directives: &[StagedBoundaryDirective],
1874) -> Result<(), CanwuError> {
1875 extend_boundary_component_overlay(current, record_overlay, overlay, directives, true)
1876}
1877
1878fn extend_boundary_component_overlay(
1879 current: &RuntimeCurrentState,
1880 record_overlay: &BTreeMap<DomainRecordRef, DomainRecord>,
1881 overlay: &mut BTreeMap<PluginComponentKey, PluginComponentRecord>,
1882 directives: &[StagedBoundaryDirective],
1883 include_next_boundary: bool,
1884) -> Result<(), CanwuError> {
1885 for staged in directives.iter().filter(|staged| {
1886 include_next_boundary || staged.visibility == StateVisibility::SameBoundary
1887 }) {
1888 if let BoundaryDirective::SetComponent {
1889 state: state_key,
1890 entity,
1891 component,
1892 value,
1893 ..
1894 } = &staged.directive
1895 {
1896 let key = component_key(&staged.plugin, state_key, entity, component);
1897 if overlay.contains_key(&key) {
1898 return Err(CanwuError::new(
1899 ErrorCode::InvalidBoundary,
1900 "multiple boundary proposals target the same component",
1901 ));
1902 }
1903 if !runtime_entity_exists_with_record_overlay(current, record_overlay, entity) {
1904 return Err(CanwuError::new(
1905 ErrorCode::EntityNotFound,
1906 format!("boundary proposal targeted missing entity {entity}"),
1907 ));
1908 }
1909 overlay.insert(
1910 key,
1911 PluginComponentRecord {
1912 plugin: staged.plugin.clone(),
1913 state: state_key.clone(),
1914 entity: entity.clone(),
1915 component: component.clone(),
1916 value: value.clone(),
1917 },
1918 );
1919 }
1920 }
1921 Ok(())
1922}
1923
1924fn extend_boundary_record_overlay(
1925 context: &BoundaryRecordOverlayContext<'_>,
1926 overlay: &mut BTreeMap<DomainRecordRef, DomainRecord>,
1927 directives: &[StagedBoundaryDirective],
1928) -> Result<(), CanwuError> {
1929 extend_boundary_domain_record_overlay(context, overlay, directives, false)
1930}
1931
1932fn extend_boundary_record_candidate_overlay(
1933 context: &BoundaryRecordOverlayContext<'_>,
1934 overlay: &mut BTreeMap<DomainRecordRef, DomainRecord>,
1935 directives: &[StagedBoundaryDirective],
1936) -> Result<(), CanwuError> {
1937 extend_boundary_domain_record_overlay(context, overlay, directives, true)
1938}
1939
1940struct BoundaryRecordOverlayContext<'a> {
1941 current: &'a RuntimeCurrentState,
1942 now: SimTime,
1943 scheduled_actions: &'a BTreeMap<ScheduleKey, ScheduledAction>,
1944 run_configuration: &'a RunConfigurationSnapshot,
1945 schemas: &'a records::DomainRecordSchemas,
1946}
1947
1948fn extend_boundary_domain_record_overlay(
1949 context: &BoundaryRecordOverlayContext<'_>,
1950 overlay: &mut BTreeMap<DomainRecordRef, DomainRecord>,
1951 directives: &[StagedBoundaryDirective],
1952 include_next_boundary: bool,
1953) -> Result<(), CanwuError> {
1954 let mut base = context.current.domain_records.clone();
1955 base.extend(
1956 overlay
1957 .iter()
1958 .map(|(reference, record)| (reference.clone(), record.clone())),
1959 );
1960 let requests: Vec<_> = directives
1961 .iter()
1962 .filter(|staged| {
1963 include_next_boundary || staged.visibility == StateVisibility::SameBoundary
1964 })
1965 .filter_map(|staged| match &staged.directive {
1966 BoundaryDirective::MutateRecord { mutation, summary } => {
1967 Some(records::DomainMutationRequest {
1968 plugin: &staged.plugin,
1969 system: &staged.system,
1970 visibility: staged.visibility,
1971 mutation,
1972 summary,
1973 })
1974 }
1975 BoundaryDirective::SetComponent { .. }
1976 | BoundaryDirective::Emit { .. }
1977 | BoundaryDirective::ScheduleIngress { .. }
1978 | BoundaryDirective::SchedulePluginIngress { .. }
1979 | BoundaryDirective::PublishKnowledge { .. } => None,
1980 })
1981 .collect();
1982 if requests.is_empty() {
1983 return Ok(());
1984 }
1985 let (next, changes) = records::apply_mutation_bundle(
1986 &base,
1987 context.schemas,
1988 context.now,
1989 &|entity| runtime_current_entity_exists(context.current, entity),
1990 requests,
1991 )?;
1992 validate_domain_dependents_with_records(
1993 &context.current.plugin_components,
1994 context.scheduled_actions,
1995 context.run_configuration,
1996 &next,
1997 )?;
1998 for change in changes {
1999 overlay.insert(change.current.reference.clone(), change.current);
2000 }
2001 Ok(())
2002}
2003
2004fn partition_boundary_visibility(
2005 directives: Vec<StagedBoundaryDirective>,
2006) -> (Vec<StagedBoundaryDirective>, Vec<StagedBoundaryDirective>) {
2007 directives
2008 .into_iter()
2009 .partition(|staged| staged.visibility == StateVisibility::SameBoundary)
2010}
2011
2012fn partition_knowledge_directives(
2013 directives: Vec<StagedBoundaryDirective>,
2014) -> (Vec<StagedBoundaryDirective>, Vec<StagedBoundaryDirective>) {
2015 directives
2016 .into_iter()
2017 .partition(|staged| matches!(staged.directive, BoundaryDirective::PublishKnowledge { .. }))
2018}
2019
2020fn allocate_reservations(
2021 mut offers: Vec<PendingReservationOffer>,
2022 mut requests: Vec<PendingReservationRequest>,
2023) -> Result<ReservationAllocationResult, CanwuError> {
2024 offers.sort_by(|left, right| {
2025 left.offer
2026 .pool
2027 .cmp(&right.offer.pool)
2028 .then_with(|| left.plugin.cmp(&right.plugin))
2029 .then_with(|| left.system.cmp(&right.system))
2030 });
2031 let mut remaining = BTreeMap::new();
2032 let mut offer_records = Vec::new();
2033 for pending in offers {
2034 if remaining
2035 .insert(pending.offer.pool.clone(), pending.offer.capacity)
2036 .is_some()
2037 {
2038 return Err(CanwuError::new(
2039 ErrorCode::InvalidBoundary,
2040 format!(
2041 "reservation pool was offered more than once, including by {}.{}",
2042 pending.plugin, pending.system
2043 ),
2044 ));
2045 }
2046 offer_records.push(ReservationOfferRecord {
2047 plugin: pending.plugin,
2048 system: pending.system,
2049 offer: pending.offer,
2050 });
2051 }
2052 requests.sort_by(|left, right| {
2053 left.request
2054 .pool
2055 .cmp(&right.request.pool)
2056 .then_with(|| right.request.priority.cmp(&left.request.priority))
2057 .then_with(|| left.request.tie_break.cmp(&right.request.tie_break))
2058 .then_with(|| left.reservation.cmp(&right.reservation))
2059 });
2060 let mut seen = BTreeSet::new();
2061 let mut by_reservation = BTreeMap::new();
2062 let mut request_records = Vec::new();
2063 let mut records = Vec::new();
2064 for pending in requests {
2065 if !seen.insert(pending.reservation.clone()) {
2066 return Err(CanwuError::new(
2067 ErrorCode::InvalidBoundary,
2068 "reservation request identity is duplicated",
2069 ));
2070 }
2071 request_records.push(ReservationRequestRecord {
2072 reservation: pending.reservation.clone(),
2073 request: pending.request.clone(),
2074 });
2075 let available = remaining.entry(pending.request.pool.clone()).or_default();
2076 let granted = pending.request.quantity.min(*available);
2077 *available -= granted;
2078 let disposition = if granted == pending.request.quantity {
2079 ReservationDisposition::Fulfilled
2080 } else if granted == 0 {
2081 ReservationDisposition::Rejected
2082 } else {
2083 ReservationDisposition::Partial
2084 };
2085 let allocation = ReservationAllocation {
2086 reservation: pending.reservation.clone(),
2087 pool: pending.request.pool,
2088 requested: pending.request.quantity,
2089 granted,
2090 remaining_after: *available,
2091 disposition,
2092 };
2093 by_reservation.insert(pending.reservation, allocation.clone());
2094 records.push(allocation);
2095 }
2096 Ok(ReservationAllocationResult {
2097 by_reservation,
2098 offers: offer_records,
2099 requests: request_records,
2100 records,
2101 })
2102}