1use super::{
39 BoundaryDirective, BoundaryId, BoundaryPhase, BoundaryProposal, BoundarySystemContract,
40 CanwuError, DomainRecordRef, DomainRecordVersionRef, DomainRecordVersionSource, EntityRef,
41 ErrorCode, PluginRegistry, RuntimeState, Simulation, SimulationSnapshot, StateKey,
42 canonical_text, invalid_snapshot,
43};
44use serde::{Deserialize, Serialize};
45use std::collections::{BTreeMap, BTreeSet};
46use std::fmt::{Display, Formatter};
47
48pub const MAX_PENDING_TRANSITION_MANIFESTS: usize = 128;
52pub const MAX_PENDING_TRANSITION_MANIFESTS_PER_COORDINATOR: usize = 32;
55pub const MAX_TRANSITION_READY_HORIZON: u64 = 1_024;
59pub const MAX_TRANSITION_PARTICIPANTS: usize = 16;
61pub const MAX_TRANSITION_EXPECTED_VERSIONS: usize = 64;
66pub const MAX_TRANSITION_LINEAGE_ID_BYTES: usize = 256;
68
69const REGISTRATION_PHASES: [BoundaryPhase; 3] = [
70 BoundaryPhase::DomainDeltaProposal,
71 BoundaryPhase::HistoricalCandidateEvaluation,
72 BoundaryPhase::StrategicAggregation,
73];
74
75#[derive(Clone, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
78pub struct TransitionManifestId {
79 pub coordinator: String,
80 pub lineage_id: String,
81 pub attempt: u32,
82}
83
84impl Display for TransitionManifestId {
85 fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result {
86 write!(
87 formatter,
88 "{}/{}#{}",
89 self.coordinator, self.lineage_id, self.attempt
90 )
91 }
92}
93
94#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
99pub struct TransitionRecordVersion {
100 pub record: DomainRecordRef,
101 pub version: u64,
102}
103
104#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
106pub struct TransitionParticipant {
107 pub plugin: String,
108 #[serde(default, skip_serializing_if = "Vec::is_empty")]
111 pub expected_pre: Vec<DomainRecordVersionRef>,
112 #[serde(default, skip_serializing_if = "Vec::is_empty")]
117 pub expected_post: Vec<TransitionRecordVersion>,
118}
119
120impl TransitionParticipant {
121 #[must_use]
122 pub fn new(plugin: impl Into<String>) -> Self {
123 Self {
124 plugin: plugin.into(),
125 expected_pre: Vec::new(),
126 expected_post: Vec::new(),
127 }
128 }
129}
130
131#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
134pub struct TransitionManifest {
135 pub lineage_id: String,
136 pub attempt: u32,
137 pub participants: Vec<TransitionParticipant>,
138 pub ready_at: BoundaryId,
139}
140
141#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
144pub struct PendingTransitionManifest {
145 pub coordinator: String,
147 pub system: String,
149 pub registered_at: BoundaryId,
150 pub manifest: TransitionManifest,
151}
152
153impl PendingTransitionManifest {
154 #[must_use]
155 pub fn id(&self) -> TransitionManifestId {
156 TransitionManifestId {
157 coordinator: self.coordinator.clone(),
158 lineage_id: self.manifest.lineage_id.clone(),
159 attempt: self.manifest.attempt,
160 }
161 }
162
163 #[must_use]
165 pub fn lists(&self, plugin: &str) -> bool {
166 self.manifest
167 .participants
168 .iter()
169 .any(|participant| participant.plugin == plugin)
170 }
171
172 pub(super) fn involves(&self, plugin: &str) -> bool {
174 self.coordinator == plugin || self.lists(plugin)
175 }
176}
177
178#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
181#[serde(rename_all = "snake_case")]
182pub enum TransitionAuditOutcome {
183 Committed,
185 Expired,
190}
191
192#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
194pub struct TransitionParticipantAudit {
195 pub plugin: String,
196 pub staged_writes: u64,
197}
198
199#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
201pub struct TransitionAuditRecord {
202 pub manifest_id: TransitionManifestId,
203 pub ready_at: BoundaryId,
208 pub outcome: TransitionAuditOutcome,
209 pub participants: Vec<TransitionParticipantAudit>,
211}
212
213impl TransitionAuditRecord {
214 pub(super) fn involves(&self, plugin: &str) -> bool {
216 self.manifest_id.coordinator == plugin
217 || self
218 .participants
219 .iter()
220 .any(|participant| participant.plugin == plugin)
221 }
222}
223
224pub(super) struct BoundaryTransitionEvidence {
227 pub(super) registered: Vec<PendingTransitionManifest>,
228 pub(super) audits: Vec<TransitionAuditRecord>,
229 pub(super) pending: BTreeMap<TransitionManifestId, PendingTransitionManifest>,
230}
231
232pub(super) struct BoundaryTransitionLedger {
235 boundary: BoundaryId,
236 pending: BTreeMap<TransitionManifestId, PendingTransitionManifest>,
237 proposed: Vec<PendingTransitionManifest>,
238 registered: Vec<PendingTransitionManifest>,
239 staged: BTreeMap<TransitionManifestId, BTreeMap<String, u64>>,
240 audits: Vec<TransitionAuditRecord>,
241 post_checks: Vec<(TransitionManifestId, TransitionRecordVersion)>,
242}
243
244const fn is_transition_directive(directive: &BoundaryDirective) -> bool {
245 matches!(
246 directive,
247 BoundaryDirective::RegisterTransitionManifest { .. }
248 | BoundaryDirective::StageTransitionWrite { .. }
249 )
250}
251
252fn require_transition_write(
253 plugin: &str,
254 contract: &BoundarySystemContract,
255 phases: &[BoundaryPhase],
256 directive: &str,
257) -> Result<(), CanwuError> {
258 let key = StateKey::core_transitions();
259 if !phases.contains(&contract.phase) || !contract.writes.contains(&key) {
260 return Err(CanwuError::new(
261 ErrorCode::UndeclaredStateWrite,
262 format!(
263 "boundary system {plugin}.{} cannot propose {directive}: it must run in phase {} and declare core write {}.{}",
264 contract.name,
265 phases
266 .iter()
267 .map(|phase| (*phase as u8).to_string())
268 .collect::<Vec<_>>()
269 .join(", "),
270 key.namespace,
271 key.name
272 ),
273 ));
274 }
275 Ok(())
276}
277
278fn plugin_can_stage(plugins: &PluginRegistry, plugin: &str) -> bool {
281 let key = StateKey::core_transitions();
282 plugins.descriptors.get(plugin).is_some_and(|descriptor| {
283 descriptor.boundary_systems.iter().any(|contract| {
284 contract.phase == BoundaryPhase::HistoricalCandidateEvaluation
285 && contract.writes.contains(&key)
286 })
287 })
288}
289
290fn invalid_manifest(message: impl Into<String>) -> CanwuError {
291 CanwuError::new(ErrorCode::InvalidBoundary, message)
292}
293
294fn validate_manifest(
296 manifest: &TransitionManifest,
297 plugins: &PluginRegistry,
298) -> Result<(), CanwuError> {
299 if !canonical_text(&manifest.lineage_id)
300 || manifest.lineage_id.len() > MAX_TRANSITION_LINEAGE_ID_BYTES
301 {
302 return Err(invalid_manifest(format!(
303 "transition manifest lineage IDs must be canonical text of at most {MAX_TRANSITION_LINEAGE_ID_BYTES} bytes"
304 )));
305 }
306 if manifest.participants.is_empty() || manifest.participants.len() > MAX_TRANSITION_PARTICIPANTS
307 {
308 return Err(CanwuError::new(
309 ErrorCode::ValueOutOfRange,
310 format!(
311 "transition manifest {} must list between 1 and {MAX_TRANSITION_PARTICIPANTS} participants",
312 manifest.lineage_id
313 ),
314 ));
315 }
316 let expected_versions = manifest
317 .participants
318 .iter()
319 .map(|participant| participant.expected_pre.len() + participant.expected_post.len())
320 .sum::<usize>();
321 if expected_versions > MAX_TRANSITION_EXPECTED_VERSIONS {
322 return Err(CanwuError::new(
323 ErrorCode::ValueOutOfRange,
324 format!(
325 "transition manifest {} declares more than {MAX_TRANSITION_EXPECTED_VERSIONS} expected versions",
326 manifest.lineage_id
327 ),
328 ));
329 }
330 let mut plugin_names = BTreeSet::new();
331 for participant in &manifest.participants {
332 if !canonical_text(&participant.plugin) || !plugin_names.insert(&participant.plugin) {
333 return Err(invalid_manifest(format!(
334 "transition manifest {} must list unique canonical participant plugins",
335 manifest.lineage_id
336 )));
337 }
338 if !plugin_can_stage(plugins, &participant.plugin) {
339 return Err(invalid_manifest(format!(
340 "transition participant {} has no phase-10 system that declares core write {}.{}",
341 participant.plugin,
342 StateKey::core_transitions().namespace,
343 StateKey::core_transitions().name
344 )));
345 }
346 let valid_record = |record: &DomainRecordRef| {
347 plugins.record_schemas.contains_key(&record.kind) && canonical_text(&record.id)
348 };
349 let mut pre = BTreeSet::new();
350 let mut post = BTreeSet::new();
351 if participant.expected_pre.iter().any(|expected| {
352 expected.version == 0
353 || !valid_record(&expected.record)
354 || !pre.insert(&expected.record)
355 }) || participant.expected_post.iter().any(|expected| {
356 expected.version == 0
357 || !valid_record(&expected.record)
358 || !post.insert(&expected.record)
359 }) {
360 return Err(invalid_manifest(format!(
361 "transition participant {} must name each expected record of a registered kind once, with a canonical ID and a nonzero version",
362 participant.plugin
363 )));
364 }
365 }
366 Ok(())
367}
368
369fn ready_at_is_valid(
373 phase: BoundaryPhase,
374 registered_at: BoundaryId,
375 ready_at: BoundaryId,
376) -> bool {
377 let earliest = if phase == BoundaryPhase::DomainDeltaProposal {
378 registered_at.get()
379 } else {
380 registered_at.get().saturating_add(1)
381 };
382 ready_at.get() >= earliest
383 && ready_at.get() - registered_at.get() <= MAX_TRANSITION_READY_HORIZON
384}
385
386fn pending_admission_error<'a>(
389 pending: impl Iterator<Item = &'a PendingTransitionManifest>,
390 coordinator: &str,
391 lineage_id: &str,
392) -> Option<CanwuError> {
393 let mut total = 0;
394 let mut own = 0;
395 for existing in pending {
396 total += 1;
397 if existing.coordinator == coordinator {
398 own += 1;
399 if existing.manifest.lineage_id == lineage_id {
400 return Some(invalid_manifest(format!(
401 "coordinator {coordinator} already has a pending transition manifest for lineage {lineage_id}"
402 )));
403 }
404 }
405 }
406 if total >= MAX_PENDING_TRANSITION_MANIFESTS
407 || own >= MAX_PENDING_TRANSITION_MANIFESTS_PER_COORDINATOR
408 {
409 return Some(CanwuError::new(
410 ErrorCode::ValueOutOfRange,
411 format!(
412 "at most {MAX_PENDING_TRANSITION_MANIFESTS} transition manifests, and {MAX_PENDING_TRANSITION_MANIFESTS_PER_COORDINATOR} per coordinator, may be pending"
413 ),
414 ));
415 }
416 None
417}
418
419fn version_list(values: &[String]) -> String {
420 values.join(", ")
421}
422
423fn version_source(source: &DomainRecordVersionSource) -> String {
424 match source {
425 DomainRecordVersionSource::InitialScenario => "the initial scenario".to_owned(),
426 DomainRecordVersionSource::BoundaryChange {
427 boundary,
428 change_index,
429 } => format!("boundary {boundary} change {change_index}"),
430 }
431}
432
433impl BoundaryTransitionLedger {
434 pub(super) fn new(
435 boundary: BoundaryId,
436 pending: &BTreeMap<TransitionManifestId, PendingTransitionManifest>,
437 ) -> Self {
438 Self {
439 boundary,
440 pending: pending.clone(),
441 proposed: Vec::new(),
442 registered: Vec::new(),
443 staged: BTreeMap::new(),
444 audits: Vec::new(),
445 post_checks: Vec::new(),
446 }
447 }
448
449 pub(super) fn pending(&self) -> impl Iterator<Item = &PendingTransitionManifest> {
451 self.pending.values()
452 }
453
454 pub(super) fn audits(&self) -> &[TransitionAuditRecord] {
456 &self.audits
457 }
458
459 pub(super) fn admit(
463 &mut self,
464 plugin: &str,
465 contract: &BoundarySystemContract,
466 plugins: &PluginRegistry,
467 mut proposal: BoundaryProposal,
468 ) -> Result<BoundaryProposal, CanwuError> {
469 if !proposal.directives.iter().any(is_transition_directive) {
470 return Ok(proposal);
471 }
472 let directives = std::mem::take(&mut proposal.directives);
473 let mut flattened = Vec::with_capacity(directives.len());
474 for directive in directives {
475 match directive {
476 BoundaryDirective::RegisterTransitionManifest { manifest } => {
477 self.admit_registration(plugin, contract, plugins, manifest)?;
478 }
479 BoundaryDirective::StageTransitionWrite {
480 manifest_id,
481 writes,
482 } => {
483 self.admit_stage(plugin, contract, &manifest_id, &writes)?;
484 flattened.extend(writes);
485 }
486 directive => flattened.push(directive),
487 }
488 }
489 proposal.directives = flattened;
490 Ok(proposal)
491 }
492
493 fn admit_registration(
494 &mut self,
495 plugin: &str,
496 contract: &BoundarySystemContract,
497 plugins: &PluginRegistry,
498 manifest: TransitionManifest,
499 ) -> Result<(), CanwuError> {
500 require_transition_write(
501 plugin,
502 contract,
503 ®ISTRATION_PHASES,
504 "RegisterTransitionManifest",
505 )?;
506 validate_manifest(&manifest, plugins)?;
507 if !ready_at_is_valid(contract.phase, self.boundary, manifest.ready_at) {
508 return Err(invalid_manifest(format!(
509 "transition manifest {} registered in phase {} of boundary {} cannot be ready at boundary {}; it must be ready within {MAX_TRANSITION_READY_HORIZON} boundaries",
510 manifest.lineage_id, contract.phase as u8, self.boundary, manifest.ready_at
511 )));
512 }
513 if let Some(error) = pending_admission_error(
514 self.pending.values().chain(&self.proposed),
515 plugin,
516 &manifest.lineage_id,
517 ) {
518 return Err(error);
519 }
520 self.proposed.push(PendingTransitionManifest {
521 coordinator: plugin.to_owned(),
522 system: contract.name.clone(),
523 registered_at: self.boundary,
524 manifest,
525 });
526 Ok(())
527 }
528
529 fn admit_stage(
530 &mut self,
531 plugin: &str,
532 contract: &BoundarySystemContract,
533 manifest_id: &TransitionManifestId,
534 writes: &[BoundaryDirective],
535 ) -> Result<(), CanwuError> {
536 require_transition_write(
537 plugin,
538 contract,
539 &[BoundaryPhase::HistoricalCandidateEvaluation],
540 "StageTransitionWrite",
541 )?;
542 let Some(pending) = self
543 .pending
544 .get(manifest_id)
545 .filter(|pending| pending.manifest.ready_at == self.boundary)
546 else {
547 return Err(invalid_manifest(format!(
548 "transition manifest {manifest_id} is not pending and ready at boundary {}",
549 self.boundary
550 )));
551 };
552 if !pending.lists(plugin) {
553 return Err(CanwuError::new(
554 ErrorCode::InvalidAuthority,
555 format!(
556 "plugin {plugin} is not a participant of transition manifest {manifest_id}"
557 ),
558 ));
559 }
560 if writes.iter().any(is_transition_directive) {
561 return Err(invalid_manifest(
562 "staged transition writes cannot register or stage transition manifests",
563 ));
564 }
565 let count = self
566 .staged
567 .entry(manifest_id.clone())
568 .or_default()
569 .entry(plugin.to_owned())
570 .or_default();
571 *count = count
572 .checked_add(u64::try_from(writes.len()).unwrap_or(u64::MAX))
573 .ok_or_else(|| {
574 CanwuError::new(
575 ErrorCode::IdentifierExhausted,
576 "staged transition write count is exhausted",
577 )
578 })?;
579 Ok(())
580 }
581
582 pub(super) fn close_phase(&mut self) {
584 for manifest in self.proposed.drain(..) {
585 self.pending.insert(manifest.id(), manifest.clone());
586 self.registered.push(manifest);
587 }
588 }
589
590 pub(super) fn settle_ready(&mut self, state: &RuntimeState) -> Result<(), CanwuError> {
593 if let Some(overdue) = self
594 .pending
595 .values()
596 .find(|pending| pending.manifest.ready_at < self.boundary)
597 {
598 return invalid_snapshot(format!(
599 "transition manifest {} passed its ready boundary without an audit",
600 overdue.id()
601 ));
602 }
603 let ready: Vec<_> = self
604 .pending
605 .iter()
606 .filter(|(_, pending)| pending.manifest.ready_at == self.boundary)
607 .map(|(id, _)| id.clone())
608 .collect();
609 for id in ready {
610 let Some(pending) = self.pending.remove(&id) else {
611 continue;
612 };
613 let staged = self.staged.remove(&id).unwrap_or_default();
614 let participants: Vec<_> = pending
615 .manifest
616 .participants
617 .iter()
618 .map(|participant| TransitionParticipantAudit {
619 plugin: participant.plugin.clone(),
620 staged_writes: staged.get(&participant.plugin).copied().unwrap_or(0),
621 })
622 .collect();
623 let missing: Vec<_> = pending
624 .manifest
625 .participants
626 .iter()
627 .filter(|participant| !staged.contains_key(&participant.plugin))
628 .map(|participant| participant.plugin.clone())
629 .collect();
630 let outcome = if missing.len() == pending.manifest.participants.len() {
631 TransitionAuditOutcome::Expired
632 } else if !missing.is_empty() {
633 return Err(CanwuError::new(
634 ErrorCode::TransitionParticipantMissing,
635 format!(
636 "transition manifest {id} at boundary {} is missing participants: {}",
637 self.boundary,
638 version_list(&missing)
639 ),
640 ));
641 } else {
642 let mut mismatched = Vec::new();
643 let mut related = Vec::new();
644 for participant in &pending.manifest.participants {
645 for expected in &participant.expected_pre {
646 let current =
647 super::current_domain_record_version(state, &expected.record)?;
648 if current.as_ref() != Some(expected) {
649 mismatched.push(format!(
650 "{} expected version {} from {} but found {}",
651 expected.record,
652 expected.version,
653 version_source(&expected.established_by),
654 current.map_or_else(
655 || "no record".to_owned(),
656 |current| format!(
657 "version {} from {}",
658 current.version,
659 version_source(¤t.established_by)
660 )
661 )
662 ));
663 related.push(EntityRef::Domain(expected.record.clone()));
664 }
665 }
666 }
667 if !mismatched.is_empty() {
668 return Err(version_mismatch(&id, "pre", &mismatched, related));
669 }
670 self.post_checks
671 .extend(
672 pending
673 .manifest
674 .participants
675 .iter()
676 .flat_map(|participant| {
677 participant
678 .expected_post
679 .iter()
680 .map(|expected| (id.clone(), expected.clone()))
681 }),
682 );
683 TransitionAuditOutcome::Committed
684 };
685 self.audits.push(TransitionAuditRecord {
686 manifest_id: id,
687 ready_at: self.boundary,
688 outcome,
689 participants,
690 });
691 }
692 if let Some((id, _)) = self.staged.iter().next() {
693 return invalid_snapshot(format!(
694 "transition writes were staged for manifest {id}, which is not ready"
695 ));
696 }
697 Ok(())
698 }
699
700 pub(super) fn requires_post_check(&self) -> bool {
703 !self.post_checks.is_empty()
704 }
705
706 pub(super) fn check_expected_post(
709 &mut self,
710 version_of: &dyn Fn(&DomainRecordRef) -> Option<u64>,
711 ) -> Result<(), CanwuError> {
712 let mut failures: BTreeMap<TransitionManifestId, (Vec<String>, Vec<EntityRef>)> =
713 BTreeMap::new();
714 for (id, expected) in std::mem::take(&mut self.post_checks) {
715 let actual = version_of(&expected.record);
716 if actual != Some(expected.version) {
717 let entry = failures.entry(id).or_default();
718 entry.0.push(format!(
719 "{} expected version {} but found {}",
720 expected.record,
721 expected.version,
722 actual.map_or_else(
723 || "no record".to_owned(),
724 |version| format!("version {version}")
725 )
726 ));
727 entry.1.push(EntityRef::Domain(expected.record));
728 }
729 }
730 if let Some((id, (mismatched, related))) = failures.into_iter().next() {
731 return Err(version_mismatch(&id, "post", &mismatched, related));
732 }
733 Ok(())
734 }
735
736 pub(super) fn finish(mut self) -> BoundaryTransitionEvidence {
737 self.close_phase();
738 BoundaryTransitionEvidence {
739 registered: self.registered,
740 audits: self.audits,
741 pending: self.pending,
742 }
743 }
744}
745
746fn version_mismatch(
747 id: &TransitionManifestId,
748 stage: &str,
749 mismatched: &[String],
750 related: Vec<EntityRef>,
751) -> CanwuError {
752 let mut error = CanwuError::new(
753 ErrorCode::TransitionVersionMismatch,
754 format!(
755 "transition manifest {id} has mismatched expected_{stage} versions: {}",
756 version_list(mismatched)
757 ),
758 );
759 error.related_entities = related;
760 error
761}
762
763impl Simulation {
764 pub fn pending_transition_manifests(&self) -> impl Iterator<Item = &PendingTransitionManifest> {
767 self.state.scheduler.transition_manifests.values()
768 }
769}
770
771fn insert_snapshot_pending(
772 pending: &mut BTreeMap<TransitionManifestId, PendingTransitionManifest>,
773 registered: &PendingTransitionManifest,
774) -> Result<(), CanwuError> {
775 if pending_admission_error(
776 pending.values(),
777 ®istered.coordinator,
778 ®istered.manifest.lineage_id,
779 )
780 .is_some()
781 {
782 return invalid_snapshot(
783 "transition manifest registrations exceed the pending bounds or repeat a pending lineage",
784 );
785 }
786 pending.insert(registered.id(), registered.clone());
787 Ok(())
788}
789
790pub(super) fn validate_snapshot_transitions(
799 snapshot: &SimulationSnapshot,
800 plugins: &PluginRegistry,
801) -> Result<(), CanwuError> {
802 let key = StateKey::core_transitions();
803 let mut pending: BTreeMap<TransitionManifestId, PendingTransitionManifest> = BTreeMap::new();
804 for record in &snapshot.boundaries {
805 let mut previous_order = None;
806 let mut late = Vec::new();
807 for registered in &record.transition_manifests {
808 let Some(contract) = super::validation::snapshot_boundary_contract(
809 plugins,
810 ®istered.coordinator,
811 ®istered.system,
812 )
813 .filter(|contract| {
814 contract.writes.contains(&key)
815 && REGISTRATION_PHASES.contains(&contract.phase)
816 && super::boundary_system_due(
817 contract,
818 &record.cadences,
819 super::boundary_has_event_ingress(record),
820 )
821 }) else {
822 return invalid_snapshot("transition manifest registration has no declared writer");
823 };
824 let order = (
825 contract.phase,
826 registered.coordinator.as_str(),
827 registered.system.as_str(),
828 );
829 if registered.registered_at != record.id
830 || previous_order.is_some_and(|previous| previous > order)
831 || !ready_at_is_valid(contract.phase, record.id, registered.manifest.ready_at)
832 || validate_manifest(®istered.manifest, plugins).is_err()
833 {
834 return invalid_snapshot(
835 "transition manifest registration evidence is inconsistent",
836 );
837 }
838 previous_order = Some(order);
839 if contract.phase < BoundaryPhase::ConditionalTransitionCommit {
840 insert_snapshot_pending(&mut pending, registered)?;
841 } else {
842 late.push(registered);
843 }
844 }
845 let ready: Vec<_> = pending
846 .iter()
847 .filter(|(_, pending)| pending.manifest.ready_at == record.id)
848 .map(|(id, _)| id.clone())
849 .collect();
850 if ready.len() != record.transition_audits.len()
851 || ready
852 .iter()
853 .zip(&record.transition_audits)
854 .any(|(id, audit)| *id != audit.manifest_id)
855 {
856 return invalid_snapshot(
857 "transition audits do not settle exactly the manifests ready at their boundary",
858 );
859 }
860 for audit in &record.transition_audits {
861 let Some(settled) = pending.remove(&audit.manifest_id) else {
862 return invalid_snapshot("transition audit names no pending manifest");
863 };
864 if audit.ready_at != record.id
865 || audit.participants.len() != settled.manifest.participants.len()
866 || audit
867 .participants
868 .iter()
869 .zip(&settled.manifest.participants)
870 .any(|(audited, listed)| audited.plugin != listed.plugin)
871 || (audit.outcome == TransitionAuditOutcome::Expired
872 && audit
873 .participants
874 .iter()
875 .any(|participant| participant.staged_writes != 0))
876 {
877 return invalid_snapshot("transition audit evidence is inconsistent");
878 }
879 }
880 if pending
881 .values()
882 .any(|pending| pending.manifest.ready_at <= record.id)
883 {
884 return invalid_snapshot(
885 "a transition manifest passed its ready boundary without an audit",
886 );
887 }
888 for registered in late {
889 insert_snapshot_pending(&mut pending, registered)?;
890 }
891 }
892 if !pending
893 .values()
894 .eq(snapshot.pending_transition_manifests.iter())
895 {
896 return invalid_snapshot(
897 "boundary transition evidence does not reconstruct the persisted pending manifests",
898 );
899 }
900 Ok(())
901}