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