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