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