Skip to main content

canwu_sim/runtime/
settlement.rs

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