1#![forbid(unsafe_code)]
2#![warn(rustdoc::broken_intra_doc_links)]
3
4use std::{
67 borrow::Cow,
68 time::{SystemTime, UNIX_EPOCH},
69};
70
71use serde::{Deserialize, Serialize};
72use sha2::{Digest, Sha256};
73use thiserror::Error;
74
75#[cfg(feature = "asupersync")]
76pub mod asupersync;
77
78#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
79#[serde(rename_all = "snake_case")]
80pub enum RuntimeMode {
81 Strict,
82 Hardened,
83}
84
85#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
86#[serde(rename_all = "snake_case")]
87pub enum DecisionAction {
88 Allow,
89 Reject,
90 Repair,
91}
92
93#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
94#[serde(rename_all = "snake_case")]
95pub enum IssueKind {
96 UnknownFeature,
97 MalformedInput,
98 JoinCardinality,
99 PolicyOverride,
100}
101
102#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
103pub struct CompatibilityIssue {
104 pub kind: IssueKind,
105 pub subject: String,
106 pub detail: String,
107}
108
109#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
110pub struct EvidenceTerm {
111 pub name: Cow<'static, str>,
112 pub log_likelihood_if_compatible: f64,
113 pub log_likelihood_if_incompatible: f64,
114}
115
116#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)]
117pub struct LossMatrix {
118 pub allow_if_compatible: f64,
119 pub allow_if_incompatible: f64,
120 pub reject_if_compatible: f64,
121 pub reject_if_incompatible: f64,
122 pub repair_if_compatible: f64,
123 pub repair_if_incompatible: f64,
124}
125
126impl Default for LossMatrix {
127 fn default() -> Self {
128 Self {
129 allow_if_compatible: 0.0,
130 allow_if_incompatible: 100.0,
131 reject_if_compatible: 6.0,
132 reject_if_incompatible: 0.5,
133 repair_if_compatible: 2.0,
134 repair_if_incompatible: 3.0,
135 }
136 }
137}
138
139const UNKNOWN_FEATURE_PRIOR: f64 = 0.25;
140const JOIN_ADMISSION_PRIOR: f64 = 0.6;
141const PRIOR_COMPATIBLE_EPSILON: f64 = 1e-10;
142
143const UNKNOWN_FEATURE_EVIDENCE: [EvidenceTerm; 2] = [
144 EvidenceTerm {
145 name: Cow::Borrowed("compatibility_allowlist_miss"),
146 log_likelihood_if_compatible: -3.5,
147 log_likelihood_if_incompatible: -0.2,
148 },
149 EvidenceTerm {
150 name: Cow::Borrowed("unknown_protocol_field"),
151 log_likelihood_if_compatible: -2.0,
152 log_likelihood_if_incompatible: -0.1,
153 },
154];
155
156const JOIN_ADMISSION_EVIDENCE_WITHIN_CAP: [EvidenceTerm; 2] = [
157 EvidenceTerm {
158 name: Cow::Borrowed("estimator_overflow_risk"),
159 log_likelihood_if_compatible: -0.3,
160 log_likelihood_if_incompatible: -1.2,
161 },
162 EvidenceTerm {
163 name: Cow::Borrowed("memory_budget_signal"),
164 log_likelihood_if_compatible: -0.4,
165 log_likelihood_if_incompatible: -1.5,
166 },
167];
168
169const JOIN_ADMISSION_EVIDENCE_OVER_CAP: [EvidenceTerm; 2] = [
170 EvidenceTerm {
171 name: Cow::Borrowed("estimator_overflow_risk"),
172 log_likelihood_if_compatible: -2.8,
173 log_likelihood_if_incompatible: -0.1,
174 },
175 EvidenceTerm {
176 name: Cow::Borrowed("memory_budget_signal"),
177 log_likelihood_if_compatible: -2.2,
178 log_likelihood_if_incompatible: -0.2,
179 },
180];
181
182const JOIN_ADMISSION_LOSS: LossMatrix = LossMatrix {
183 allow_if_compatible: 0.0,
184 allow_if_incompatible: 130.0,
185 reject_if_compatible: 5.0,
186 reject_if_incompatible: 0.5,
187 repair_if_compatible: 1.5,
188 repair_if_incompatible: 3.0,
189};
190
191const DEFAULT_CONFORMAL_ALPHA: f64 = 0.1;
192const MIN_CONFORMAL_ALPHA: f64 = 0.01;
193const MAX_CONFORMAL_ALPHA: f64 = 0.5;
194
195#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
196pub struct DecisionMetrics {
197 pub posterior_compatible: f64,
198 pub bayes_factor_compatible_over_incompatible: f64,
199 pub expected_loss_allow: f64,
200 pub expected_loss_reject: f64,
201 pub expected_loss_repair: f64,
202}
203
204#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
205pub struct DecisionRecord {
206 pub ts_unix_ms: u64,
207 pub mode: RuntimeMode,
208 pub action: DecisionAction,
209 pub issue: CompatibilityIssue,
210 pub prior_compatible: f64,
211 pub metrics: DecisionMetrics,
212 pub evidence: Vec<EvidenceTerm>,
213}
214
215#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
216pub struct SemanticIndexIdentity {
217 pub role: String,
218 pub len: usize,
219 pub has_duplicates: bool,
220 pub fingerprint: String,
221}
222
223#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
224pub struct SemanticWitnessRecord {
225 pub ts_unix_ms: u64,
226 pub operation: String,
227 pub materialization_reason: String,
228 pub alignment_mode: String,
229 pub input_index_identity: Vec<SemanticIndexIdentity>,
230 pub output_index_identity: SemanticIndexIdentity,
231 pub null_nan_policy: String,
232 pub output_ordering_contract: String,
233}
234
235impl SemanticWitnessRecord {
236 #[must_use]
237 pub fn new(
238 operation: impl Into<String>,
239 materialization_reason: impl Into<String>,
240 alignment_mode: impl Into<String>,
241 input_index_identity: Vec<SemanticIndexIdentity>,
242 output_index_identity: SemanticIndexIdentity,
243 null_nan_policy: impl Into<String>,
244 output_ordering_contract: impl Into<String>,
245 ) -> Self {
246 Self {
247 ts_unix_ms: now_unix_ms().unwrap_or_default(),
248 operation: operation.into(),
249 materialization_reason: materialization_reason.into(),
250 alignment_mode: alignment_mode.into(),
251 input_index_identity,
252 output_index_identity,
253 null_nan_policy: null_nan_policy.into(),
254 output_ordering_contract: output_ordering_contract.into(),
255 }
256 }
257}
258
259#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
260pub struct GalaxyBrainCard {
261 pub title: String,
262 pub equation: String,
263 pub substitution: String,
264 pub intuition: String,
265}
266
267impl GalaxyBrainCard {
268 #[must_use]
269 pub fn render_plain(&self) -> String {
270 format!(
271 "[{}]\n{}\n{}\n{}",
272 self.title, self.equation, self.substitution, self.intuition
273 )
274 }
275}
276
277#[must_use]
278pub fn decision_to_card(record: &DecisionRecord) -> GalaxyBrainCard {
279 GalaxyBrainCard {
280 title: format!("{}::{:?}", record.issue.subject, record.action),
281 equation: "argmin_a Σ_s L(a,s) P(s|evidence)".to_owned(),
282 substitution: format!(
283 "P(compatible|e)={:.4}, E[allow]={:.4}, E[reject]={:.4}, E[repair]={:.4}",
284 record.metrics.posterior_compatible,
285 record.metrics.expected_loss_allow,
286 record.metrics.expected_loss_reject,
287 record.metrics.expected_loss_repair
288 ),
289 intuition: "Lower expected loss wins; strict mode may still force fail-closed.".to_owned(),
290 }
291}
292
293#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
294pub struct EvidenceLedger {
295 records: Vec<DecisionRecord>,
296 #[serde(default)]
297 semantic_witnesses: Vec<SemanticWitnessRecord>,
298 #[serde(skip)]
304 record_semantic_witnesses: bool,
305}
306
307impl Default for EvidenceLedger {
308 fn default() -> Self {
309 Self::new()
310 }
311}
312
313impl EvidenceLedger {
314 #[must_use]
315 pub fn new() -> Self {
316 Self {
317 records: Vec::new(),
318 semantic_witnesses: Vec::new(),
319 record_semantic_witnesses: true,
320 }
321 }
322
323 #[must_use]
328 pub fn without_semantic_witnesses(mut self) -> Self {
329 self.record_semantic_witnesses = false;
330 self
331 }
332
333 #[must_use]
335 pub fn records_semantic_witnesses(&self) -> bool {
336 self.record_semantic_witnesses
337 }
338
339 pub fn push(&mut self, record: DecisionRecord) {
340 self.records.push(record);
341 }
342
343 pub fn push_semantic_witness(&mut self, record: SemanticWitnessRecord) {
344 self.semantic_witnesses.push(record);
345 }
346
347 #[must_use]
348 pub fn records(&self) -> &[DecisionRecord] {
349 &self.records
350 }
351
352 #[must_use]
353 pub fn semantic_witnesses(&self) -> &[SemanticWitnessRecord] {
354 &self.semantic_witnesses
355 }
356}
357
358#[derive(Debug, Clone, PartialEq, Eq)]
359pub struct RuntimePolicy {
360 pub mode: RuntimeMode,
361 pub fail_closed_unknown_features: bool,
362 pub hardened_join_row_cap: Option<usize>,
363}
364
365impl RuntimePolicy {
366 #[must_use]
367 pub fn strict() -> Self {
368 Self {
369 mode: RuntimeMode::Strict,
370 fail_closed_unknown_features: true,
371 hardened_join_row_cap: None,
372 }
373 }
374
375 #[must_use]
376 pub fn hardened(join_row_cap: Option<usize>) -> Self {
377 Self {
378 mode: RuntimeMode::Hardened,
379 fail_closed_unknown_features: false,
380 hardened_join_row_cap: join_row_cap,
381 }
382 }
383
384 pub fn decide_unknown_feature(
385 &self,
386 subject: impl Into<String>,
387 detail: impl Into<String>,
388 ledger: &mut EvidenceLedger,
389 ) -> DecisionAction {
390 let issue = CompatibilityIssue {
391 kind: IssueKind::UnknownFeature,
392 subject: subject.into(),
393 detail: detail.into(),
394 };
395
396 let mut record = decide(
397 self.mode,
398 issue,
399 UNKNOWN_FEATURE_PRIOR,
400 LossMatrix::default(),
401 UNKNOWN_FEATURE_EVIDENCE.to_vec(),
402 );
403 if self.fail_closed_unknown_features {
404 record.action = DecisionAction::Reject;
405 }
406 let action = record.action;
407 ledger.push(record);
408 action
409 }
410
411 pub fn decide_join_admission(
412 &self,
413 estimated_rows: usize,
414 ledger: &mut EvidenceLedger,
415 ) -> DecisionAction {
416 let issue = CompatibilityIssue {
417 kind: IssueKind::JoinCardinality,
418 subject: "join_estimator".to_owned(),
419 detail: format!("estimated_rows={estimated_rows}"),
420 };
421
422 let cap = self.hardened_join_row_cap.unwrap_or(usize::MAX);
423 let evidence = if estimated_rows <= cap {
424 JOIN_ADMISSION_EVIDENCE_WITHIN_CAP.to_vec()
425 } else {
426 JOIN_ADMISSION_EVIDENCE_OVER_CAP.to_vec()
427 };
428 let mut record = decide(
429 self.mode,
430 issue,
431 JOIN_ADMISSION_PRIOR,
432 JOIN_ADMISSION_LOSS,
433 evidence,
434 );
435
436 if matches!(self.mode, RuntimeMode::Hardened) && estimated_rows > cap {
437 record.action = DecisionAction::Repair;
438 }
439
440 let action = record.action;
441 ledger.push(record);
442 action
443 }
444}
445
446impl Default for RuntimePolicy {
447 fn default() -> Self {
448 Self::strict()
449 }
450}
451
452#[derive(Debug, Error)]
453pub enum RuntimeError {
454 #[error("system clock is before UNIX_EPOCH")]
455 ClockSkew,
456}
457
458fn now_unix_ms() -> Result<u64, RuntimeError> {
459 let ms = SystemTime::now()
460 .duration_since(UNIX_EPOCH)
461 .map_err(|_| RuntimeError::ClockSkew)?
462 .as_millis();
463 Ok(ms as u64)
464}
465
466fn normalize_prior_compatible(prior_compatible: f64) -> f64 {
467 if !prior_compatible.is_finite() {
468 return 0.5;
469 }
470 prior_compatible.clamp(PRIOR_COMPATIBLE_EPSILON, 1.0 - PRIOR_COMPATIBLE_EPSILON)
471}
472
473fn decide(
474 mode: RuntimeMode,
475 issue: CompatibilityIssue,
476 prior_compatible: f64,
477 loss: LossMatrix,
478 evidence: Vec<EvidenceTerm>,
479) -> DecisionRecord {
480 let prior_compatible = normalize_prior_compatible(prior_compatible);
481 let log_odds_prior = (prior_compatible / (1.0 - prior_compatible)).ln();
482 let llr_sum: f64 = evidence
483 .iter()
484 .map(|term| term.log_likelihood_if_compatible - term.log_likelihood_if_incompatible)
485 .sum();
486 let log_odds_post = log_odds_prior + llr_sum;
487
488 let posterior_compatible = 1.0 / (1.0 + (-log_odds_post).exp());
489 let posterior_incompatible = 1.0 - posterior_compatible;
490
491 let expected_loss_allow = loss.allow_if_compatible * posterior_compatible
492 + loss.allow_if_incompatible * posterior_incompatible;
493 let expected_loss_reject = loss.reject_if_compatible * posterior_compatible
494 + loss.reject_if_incompatible * posterior_incompatible;
495 let expected_loss_repair = loss.repair_if_compatible * posterior_compatible
496 + loss.repair_if_incompatible * posterior_incompatible;
497
498 let mut best_action = DecisionAction::Allow;
499 let mut best_loss = expected_loss_allow;
500
501 if expected_loss_repair < best_loss {
502 best_action = DecisionAction::Repair;
503 best_loss = expected_loss_repair;
504 }
505 if expected_loss_reject < best_loss {
506 best_action = DecisionAction::Reject;
507 }
508
509 DecisionRecord {
510 ts_unix_ms: now_unix_ms().unwrap_or_default(),
511 mode,
512 action: best_action,
513 issue,
514 prior_compatible,
515 metrics: DecisionMetrics {
516 posterior_compatible,
517 bayes_factor_compatible_over_incompatible: llr_sum.exp(),
518 expected_loss_allow,
519 expected_loss_reject,
520 expected_loss_repair,
521 },
522 evidence,
523 }
524}
525
526#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
527pub struct RaptorQEnvelope {
528 pub artifact_id: String,
529 pub artifact_type: String,
530 pub source_hash: String,
531 pub raptorq: RaptorQMetadata,
532 pub scrub: ScrubStatus,
533 pub decode_proofs: Vec<DecodeProof>,
534}
535
536pub const MAX_DECODE_PROOFS: usize = 1_000;
537pub const DEFAULT_RAPTORQ_SYMBOL_BYTES: usize = 1_024;
538
539#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
540pub struct RaptorQMetadata {
541 pub k: u32,
542 pub repair_symbols: u32,
543 pub overhead_ratio: f64,
544 pub symbol_hashes: Vec<String>,
545}
546
547#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
548pub struct ScrubStatus {
549 pub last_ok_unix_ms: u64,
550 pub status: String,
551}
552
553#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
554pub struct DecodeProof {
555 pub ts_unix_ms: u64,
556 pub reason: String,
557 pub recovered_blocks: u32,
558 pub proof_hash: String,
559}
560
561impl RaptorQEnvelope {
562 #[must_use]
563 pub fn from_source_bytes(
564 artifact_id: impl Into<String>,
565 artifact_type: impl Into<String>,
566 source_bytes: &[u8],
567 repair_symbols: u32,
568 ) -> Self {
569 let symbol_hashes: Vec<String> = source_bytes
570 .chunks(DEFAULT_RAPTORQ_SYMBOL_BYTES)
571 .map(|chunk| format!("sha256:{}", sha256_hex(chunk)))
572 .collect();
573 let k = u32::try_from(symbol_hashes.len()).unwrap_or(u32::MAX);
574 let overhead_ratio = if k == 0 {
575 0.0
576 } else {
577 f64::from(repair_symbols) / f64::from(k)
578 };
579
580 Self {
581 artifact_id: artifact_id.into(),
582 artifact_type: artifact_type.into(),
583 source_hash: format!("sha256:{}", sha256_hex(source_bytes)),
584 raptorq: RaptorQMetadata {
585 k,
586 repair_symbols,
587 overhead_ratio,
588 symbol_hashes,
589 },
590 scrub: ScrubStatus {
591 last_ok_unix_ms: now_unix_ms().unwrap_or_default(),
592 status: "ok".to_owned(),
593 },
594 decode_proofs: Vec::new(),
595 }
596 }
597
598 pub fn push_decode_proof_capped(&mut self, proof: DecodeProof) {
602 if self.decode_proofs.len() >= MAX_DECODE_PROOFS {
603 let overflow = self.decode_proofs.len() + 1 - MAX_DECODE_PROOFS;
604 self.decode_proofs.drain(0..overflow);
605 }
606 self.decode_proofs.push(proof);
607 }
608}
609
610#[must_use]
611pub fn semantic_fingerprint_bytes(bytes: &[u8]) -> String {
612 format!("sha256:{}", sha256_hex(bytes))
613}
614
615#[derive(Debug)]
616pub struct SemanticFingerprintBuilder {
617 hasher: Sha256,
618}
619
620impl Default for SemanticFingerprintBuilder {
621 fn default() -> Self {
622 Self::new()
623 }
624}
625
626impl SemanticFingerprintBuilder {
627 #[must_use]
628 pub fn new() -> Self {
629 Self {
630 hasher: Sha256::new(),
631 }
632 }
633
634 pub fn update(&mut self, bytes: &[u8]) {
635 self.hasher.update(bytes);
636 }
637
638 #[must_use]
639 pub fn finish(self) -> String {
640 format!("sha256:{}", sha256_digest_hex(self.hasher.finalize()))
641 }
642}
643
644fn sha256_hex(bytes: &[u8]) -> String {
645 let digest = Sha256::digest(bytes);
646 sha256_digest_hex(digest)
647}
648
649fn sha256_digest_hex(digest: impl IntoIterator<Item = u8>) -> String {
650 const HEX: &[u8; 16] = b"0123456789abcdef";
651 let mut hex = String::with_capacity(64);
652 for byte in digest {
653 hex.push(char::from(HEX[usize::from(byte >> 4)]));
654 hex.push(char::from(HEX[usize::from(byte & 0x0f)]));
655 }
656 hex
657}
658
659fn nonconformity_score(record: &DecisionRecord) -> f64 {
664 let p = record
666 .metrics
667 .posterior_compatible
668 .clamp(1e-15, 1.0 - 1e-15);
669 (p / (1.0 - p)).ln().abs()
670}
671
672fn normalize_conformal_alpha(alpha: f64) -> f64 {
673 if alpha.is_finite() {
674 alpha.clamp(MIN_CONFORMAL_ALPHA, MAX_CONFORMAL_ALPHA)
675 } else {
676 DEFAULT_CONFORMAL_ALPHA
677 }
678}
679
680#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
682pub struct ConformalPredictionSet {
683 pub quantile_threshold: f64,
685 pub current_score: f64,
687 pub bayesian_action_in_set: bool,
689 pub admissible_actions: Vec<DecisionAction>,
691 pub empirical_coverage: f64,
693}
694
695#[derive(Debug, Clone, Serialize, Deserialize)]
697pub struct ConformalGuard {
698 scores: Vec<f64>,
700 window_size: usize,
702 alpha: f64,
704 in_set_count: usize,
706 total_count: usize,
708}
709
710impl ConformalGuard {
711 #[must_use]
713 pub fn new(window_size: usize, alpha: f64) -> Self {
714 let window_size = window_size.max(1);
715 Self {
716 scores: Vec::with_capacity(window_size),
717 window_size,
718 alpha: normalize_conformal_alpha(alpha),
719 in_set_count: 0,
720 total_count: 0,
721 }
722 }
723
724 #[must_use]
726 pub fn default_config() -> Self {
727 Self::new(1000, 0.1)
728 }
729
730 #[must_use]
733 pub fn conformal_quantile(&self) -> Option<f64> {
734 let mut sorted: Vec<f64> = self
735 .scores
736 .iter()
737 .copied()
738 .filter(|score| score.is_finite())
739 .collect();
740 if sorted.len() < 2 {
741 return None;
742 }
743 sorted.sort_by(f64::total_cmp);
744 let n = sorted.len() as f64;
746 let level = (1.0 - normalize_conformal_alpha(self.alpha)) * (1.0 + 1.0 / n);
747 let idx = (level * n).ceil() as usize;
748 let idx = idx.min(sorted.len()).saturating_sub(1);
749 Some(sorted[idx])
750 }
751
752 pub fn evaluate(&mut self, record: &DecisionRecord) -> ConformalPredictionSet {
755 self.normalize_runtime_config();
756 let score = nonconformity_score(record);
757
758 let quantile = self.conformal_quantile();
759
760 if self.scores.len() >= self.window_size {
762 self.scores.remove(0);
763 }
764 self.scores.push(score);
765
766 let threshold = match quantile {
767 Some(q) => q,
768 None => {
769 self.total_count += 1;
771 self.in_set_count += 1;
772 return ConformalPredictionSet {
773 quantile_threshold: f64::INFINITY,
774 current_score: score,
775 bayesian_action_in_set: true,
776 admissible_actions: vec![
777 DecisionAction::Allow,
778 DecisionAction::Reject,
779 DecisionAction::Repair,
780 ],
781 empirical_coverage: 1.0,
782 };
783 }
784 };
785
786 let bayesian_in_set = score <= threshold;
787
788 let admissible = if bayesian_in_set {
792 vec![record.action]
793 } else {
794 vec![
795 DecisionAction::Allow,
796 DecisionAction::Reject,
797 DecisionAction::Repair,
798 ]
799 };
800
801 self.total_count += 1;
802 if bayesian_in_set {
803 self.in_set_count += 1;
804 }
805
806 let empirical_coverage = if self.total_count > 0 {
807 self.in_set_count as f64 / self.total_count as f64
808 } else {
809 1.0
810 };
811
812 ConformalPredictionSet {
813 quantile_threshold: threshold,
814 current_score: score,
815 bayesian_action_in_set: bayesian_in_set,
816 admissible_actions: admissible,
817 empirical_coverage,
818 }
819 }
820
821 #[must_use]
823 pub fn empirical_coverage(&self) -> f64 {
824 if self.total_count == 0 {
825 return 1.0;
826 }
827 self.in_set_count.min(self.total_count) as f64 / self.total_count as f64
828 }
829
830 #[must_use]
832 pub fn calibration_count(&self) -> usize {
833 self.scores.len()
834 }
835
836 #[must_use]
838 pub fn is_calibrated(&self) -> bool {
839 self.scores.iter().filter(|score| score.is_finite()).count() >= 2
840 }
841
842 #[must_use]
844 pub fn coverage_alert(&self) -> bool {
845 self.total_count >= 100
846 && self.empirical_coverage() < (1.0 - normalize_conformal_alpha(self.alpha))
847 }
848
849 fn normalize_runtime_config(&mut self) {
850 self.window_size = self.window_size.max(1);
851 self.alpha = normalize_conformal_alpha(self.alpha);
852 self.scores.retain(|score| score.is_finite());
853 if self.scores.len() > self.window_size {
854 let overflow = self.scores.len() - self.window_size;
855 self.scores.drain(0..overflow);
856 }
857 self.in_set_count = self.in_set_count.min(self.total_count);
858 }
859}
860
861#[cfg(feature = "asupersync")]
862#[must_use]
863pub fn outcome_to_action<T, E>(outcome: &::asupersync::Outcome<T, E>) -> DecisionAction {
864 match outcome {
865 ::asupersync::Outcome::Ok(_) => DecisionAction::Allow,
866 ::asupersync::Outcome::Err(_) => DecisionAction::Repair,
867 ::asupersync::Outcome::Cancelled(_) | ::asupersync::Outcome::Panicked(_) => {
868 DecisionAction::Reject
869 }
870 }
871}
872
873#[cfg(test)]
874mod tests {
875 use std::{borrow::Cow, hint::black_box, time::Instant};
876
877 use serde::Serialize;
878
879 use super::{
880 ConformalGuard, DecisionAction, EvidenceLedger, RaptorQEnvelope, RuntimeMode,
881 RuntimePolicy, SemanticIndexIdentity, SemanticWitnessRecord, decision_to_card,
882 };
883
884 const ASUPERSYNC_PACKET_ID: &str = "ASUPERSYNC-E";
885 const REPLAY_PREFIX: &str = "cargo test -p fp-runtime --";
886
887 #[test]
891 fn semantic_fingerprint_sha256_known_answers_01gdm() {
892 assert_eq!(
894 super::semantic_fingerprint_bytes(b""),
895 "sha256:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
896 );
897 assert_eq!(
898 super::semantic_fingerprint_bytes(b"abc"),
899 "sha256:ba7816bf8f01cfea414140de5dae2223b00361a396177a9cb410ff61f20015ad"
900 );
901
902 let f = super::semantic_fingerprint_bytes(b"frankenpandas");
904 assert!(f.starts_with("sha256:"), "prefix");
905 assert_eq!(f.len(), 7 + 64, "length");
906 assert!(
907 f[7..]
908 .bytes()
909 .all(|c| c.is_ascii_digit() || (b'a'..=b'f').contains(&c)),
910 "lowercase hex"
911 );
912 assert_eq!(
913 f,
914 super::semantic_fingerprint_bytes(b"frankenpandas"),
915 "deterministic"
916 );
917
918 let inputs: [&[u8]; 5] = [b"a", b"b", b"ab", b"ba", b""];
920 for (i, x) in inputs.iter().enumerate() {
921 for (j, y) in inputs.iter().enumerate() {
922 if i != j {
923 assert_ne!(
924 super::semantic_fingerprint_bytes(x),
925 super::semantic_fingerprint_bytes(y),
926 "distinct {i} vs {j}"
927 );
928 }
929 }
930 }
931
932 let mut b = super::SemanticFingerprintBuilder::new();
934 b.update(b"hello");
935 b.update(b" ");
936 b.update(b"world");
937 assert_eq!(
938 b.finish(),
939 super::semantic_fingerprint_bytes(b"hello world")
940 );
941 }
942
943 #[test]
944 fn semantic_fingerprint_streaming_equals_oneshot_h2i8m() {
945 let mut st: u64 = 0x4f1e_0b1c_2d3e_4f50;
949 let mut next = || {
950 st = st
951 .wrapping_mul(6_364_136_223_846_793_005)
952 .wrapping_add(1_442_695_040_888_963_407);
953 (st >> 33) as u32
954 };
955 for iter in 0..400u32 {
956 let total = (next() % 40) as usize;
957 let bytes: Vec<u8> = (0..total).map(|_| (next() % 256) as u8).collect();
958 let mut builder = super::SemanticFingerprintBuilder::new();
960 let mut pos = 0usize;
961 while pos < bytes.len() {
962 let remaining = bytes.len() - pos;
963 let take = (next() as usize % (remaining + 1)).min(remaining);
964 builder.update(&bytes[pos..pos + take]);
965 pos += take;
966 if take == 0 {
967 builder.update(&bytes[pos..pos + 1.min(bytes.len() - pos)]);
969 pos += 1;
970 }
971 }
972 assert_eq!(
973 builder.finish(),
974 super::semantic_fingerprint_bytes(&bytes),
975 "streaming==one-shot iter={iter} total={total}"
976 );
977 }
978 }
979
980 #[test]
983 fn raptorq_envelope_from_source_bytes_invariants_b9vvk() {
984 let sym = super::DEFAULT_RAPTORQ_SYMBOL_BYTES;
985 let repair = 3u32;
986 for &len in &[0usize, 1, sym, sym + 1, 3 * sym - 72] {
987 let source: Vec<u8> = (0..len).map(|i| (i % 251) as u8).collect();
988 let env = RaptorQEnvelope::from_source_bytes("pkt-1", "conformance", &source, repair);
989
990 let expected_k = len.div_ceil(sym) as u32; assert_eq!(env.raptorq.k, expected_k, "k for len={len}");
992 assert_eq!(
993 env.raptorq.symbol_hashes.len() as u32,
994 expected_k,
995 "one symbol hash per source symbol, len={len}"
996 );
997 assert_eq!(
998 env.raptorq.repair_symbols, repair,
999 "repair_symbols len={len}"
1000 );
1001 assert_eq!(
1002 env.source_hash,
1003 super::semantic_fingerprint_bytes(&source),
1004 "source_hash == fingerprint, len={len}"
1005 );
1006 let expected_overhead = if expected_k == 0 {
1007 0.0
1008 } else {
1009 f64::from(repair) / f64::from(expected_k)
1010 };
1011 assert_eq!(
1012 env.raptorq.overhead_ratio, expected_overhead,
1013 "overhead len={len}"
1014 );
1015 assert_eq!(env.scrub.status, "ok", "scrub ok len={len}");
1016 assert!(
1017 env.decode_proofs.is_empty(),
1018 "no decode proofs yet len={len}"
1019 );
1020 assert!(
1022 env.raptorq
1023 .symbol_hashes
1024 .iter()
1025 .all(|h| h.starts_with("sha256:") && h.len() == 7 + 64),
1026 "symbol hash format len={len}"
1027 );
1028 }
1029 }
1030
1031 #[test]
1034 fn raptorq_decode_proof_cap_fifo_bhlwt() {
1035 let mut env = RaptorQEnvelope::from_source_bytes("p", "t", b"x", 1);
1036 let cap = super::MAX_DECODE_PROOFS;
1037 let total = cap + 5;
1038 for i in 0..total {
1039 env.push_decode_proof_capped(super::DecodeProof {
1040 ts_unix_ms: i as u64,
1041 reason: "scrub".to_owned(),
1042 recovered_blocks: i as u32,
1043 proof_hash: "sha256:deadbeef".to_owned(),
1044 });
1045 }
1046 assert_eq!(
1048 env.decode_proofs.len(),
1049 cap,
1050 "history capped at MAX_DECODE_PROOFS"
1051 );
1052 assert_eq!(
1054 env.decode_proofs.first().unwrap().recovered_blocks,
1055 (total - cap) as u32,
1056 "oldest evicted; first retained is seq=overflow"
1057 );
1058 assert_eq!(
1059 env.decode_proofs.last().unwrap().recovered_blocks,
1060 (total - 1) as u32,
1061 "newest retained"
1062 );
1063 assert!(
1065 env.decode_proofs
1066 .windows(2)
1067 .all(|w| w[1].recovered_blocks == w[0].recovered_blocks + 1),
1068 "retained window is contiguous"
1069 );
1070 }
1071
1072 #[test]
1075 fn runtime_policy_failclosed_and_join_cap_mbjpj() {
1076 let s = RuntimePolicy::strict();
1078 assert_eq!(s.mode, RuntimeMode::Strict);
1079 assert!(s.fail_closed_unknown_features);
1080 assert_eq!(s.hardened_join_row_cap, None);
1081 let h = RuntimePolicy::hardened(Some(10));
1082 assert_eq!(h.mode, RuntimeMode::Hardened);
1083 assert!(!h.fail_closed_unknown_features);
1084 assert_eq!(h.hardened_join_row_cap, Some(10));
1085 assert_eq!(
1086 RuntimePolicy::default().mode,
1087 RuntimeMode::Strict,
1088 "default is strict"
1089 );
1090
1091 let mut led = EvidenceLedger::new();
1093 let action =
1094 RuntimePolicy::strict().decide_unknown_feature("widget", "no handler", &mut led);
1095 assert_eq!(
1096 action,
1097 DecisionAction::Reject,
1098 "strict fail-closes unknown features"
1099 );
1100 assert_eq!(led.records().len(), 1, "decision recorded");
1101
1102 let mut led2 = EvidenceLedger::new();
1104 let over = RuntimePolicy::hardened(Some(10)).decide_join_admission(1_000, &mut led2);
1105 assert_eq!(over, DecisionAction::Repair, "hardened caps over-cap joins");
1106 assert_eq!(led2.records().len(), 1, "join decision recorded");
1107 }
1108
1109 #[derive(Debug, Clone, PartialEq, Eq, Serialize)]
1110 struct StructuredTestLog {
1111 packet_id: String,
1112 case_id: String,
1113 mode: RuntimeMode,
1114 seed: u64,
1115 trace_id: String,
1116 assertion_path: String,
1117 result: String,
1118 replay_cmd: String,
1119 }
1120
1121 fn make_structured_log(
1122 case_id: &str,
1123 mode: RuntimeMode,
1124 seed: u64,
1125 assertion_path: &str,
1126 result: &str,
1127 ) -> StructuredTestLog {
1128 StructuredTestLog {
1129 packet_id: ASUPERSYNC_PACKET_ID.to_owned(),
1130 case_id: case_id.to_owned(),
1131 mode,
1132 seed,
1133 trace_id: format!("{ASUPERSYNC_PACKET_ID}:{case_id}:{seed:016x}"),
1134 assertion_path: assertion_path.to_owned(),
1135 result: result.to_owned(),
1136 replay_cmd: format!("{REPLAY_PREFIX} {case_id} --nocapture"),
1137 }
1138 }
1139
1140 fn assert_required_log_fields(log: &serde_json::Value) {
1141 for field in [
1142 "packet_id",
1143 "case_id",
1144 "mode",
1145 "seed",
1146 "trace_id",
1147 "assertion_path",
1148 "result",
1149 "replay_cmd",
1150 ] {
1151 assert!(
1152 log.get(field).is_some(),
1153 "structured log missing field: {field}"
1154 );
1155 }
1156 }
1157
1158 #[test]
1159 fn evidence_ledger_records_semantic_witnesses_tn6qb3() {
1160 let mut ledger = EvidenceLedger::new();
1161 let witness = SemanticWitnessRecord::new(
1162 "series.add",
1163 "series_binary_arithmetic_materialization",
1164 "outer",
1165 vec![
1166 SemanticIndexIdentity {
1167 role: "left".to_owned(),
1168 len: 2,
1169 has_duplicates: false,
1170 fingerprint: super::semantic_fingerprint_bytes(b"left"),
1171 },
1172 SemanticIndexIdentity {
1173 role: "right".to_owned(),
1174 len: 2,
1175 has_duplicates: false,
1176 fingerprint: super::semantic_fingerprint_bytes(b"right"),
1177 },
1178 ],
1179 SemanticIndexIdentity {
1180 role: "output".to_owned(),
1181 len: 3,
1182 has_duplicates: false,
1183 fingerprint: super::semantic_fingerprint_bytes(b"output"),
1184 },
1185 "missing aligned operands materialize as NaN/null before arithmetic",
1186 "outer union preserves left order then right-only labels",
1187 );
1188
1189 ledger.push_semantic_witness(witness);
1190
1191 let witnesses = ledger.semantic_witnesses();
1192 assert_eq!(witnesses.len(), 1);
1193 assert_eq!(witnesses[0].operation, "series.add");
1194 assert_eq!(witnesses[0].alignment_mode, "outer");
1195 assert_eq!(witnesses[0].output_index_identity.len, 3);
1196 assert_eq!(witnesses[0].input_index_identity[0].role, "left");
1197 assert!(
1198 witnesses[0].input_index_identity[0]
1199 .fingerprint
1200 .starts_with("sha256:")
1201 );
1202 }
1203
1204 fn decide_join_admission_baseline(
1205 policy: &RuntimePolicy,
1206 estimated_rows: usize,
1207 ledger: &mut EvidenceLedger,
1208 ) -> DecisionAction {
1209 let issue = super::CompatibilityIssue {
1210 kind: super::IssueKind::JoinCardinality,
1211 subject: "join_estimator".to_owned(),
1212 detail: format!("estimated_rows={estimated_rows}"),
1213 };
1214 let cap = policy.hardened_join_row_cap.unwrap_or(usize::MAX);
1215 let evidence = vec![
1216 super::EvidenceTerm {
1217 name: Cow::Owned("estimator_overflow_risk".to_owned()),
1218 log_likelihood_if_compatible: if estimated_rows <= cap { -0.3 } else { -2.8 },
1219 log_likelihood_if_incompatible: if estimated_rows <= cap { -1.2 } else { -0.1 },
1220 },
1221 super::EvidenceTerm {
1222 name: Cow::Owned("memory_budget_signal".to_owned()),
1223 log_likelihood_if_compatible: if estimated_rows <= cap { -0.4 } else { -2.2 },
1224 log_likelihood_if_incompatible: if estimated_rows <= cap { -1.5 } else { -0.2 },
1225 },
1226 ];
1227 let loss = super::LossMatrix {
1228 allow_if_compatible: 0.0,
1229 allow_if_incompatible: 130.0,
1230 reject_if_compatible: 5.0,
1231 reject_if_incompatible: 0.5,
1232 repair_if_compatible: 1.5,
1233 repair_if_incompatible: 3.0,
1234 };
1235 let mut record = super::decide(policy.mode, issue, 0.6, loss, evidence);
1236 if matches!(policy.mode, RuntimeMode::Hardened) && estimated_rows > cap {
1237 record.action = DecisionAction::Repair;
1238 }
1239 let action = record.action;
1240 ledger.push(record);
1241 action
1242 }
1243
1244 fn assert_join_record_equivalent(
1245 optimized: &super::DecisionRecord,
1246 baseline: &super::DecisionRecord,
1247 ) {
1248 assert_eq!(optimized.mode, baseline.mode);
1249 assert_eq!(optimized.action, baseline.action);
1250 assert_eq!(optimized.issue.kind, baseline.issue.kind);
1251 assert_eq!(optimized.issue.subject, baseline.issue.subject);
1252 assert_eq!(optimized.issue.detail, baseline.issue.detail);
1253 assert_eq!(optimized.prior_compatible, baseline.prior_compatible);
1254 assert_eq!(optimized.metrics, baseline.metrics);
1255 assert_eq!(optimized.evidence.len(), baseline.evidence.len());
1256 for (left, right) in optimized.evidence.iter().zip(&baseline.evidence) {
1257 assert_eq!(left.name.as_ref(), right.name.as_ref());
1258 assert_eq!(
1259 left.log_likelihood_if_compatible,
1260 right.log_likelihood_if_compatible
1261 );
1262 assert_eq!(
1263 left.log_likelihood_if_incompatible,
1264 right.log_likelihood_if_incompatible
1265 );
1266 }
1267 }
1268
1269 fn quantile_from_sorted(samples: &[u128], pct: usize) -> u128 {
1270 let len = samples.len();
1271 assert!(len > 0);
1272 let idx = (len.saturating_sub(1) * pct) / 100;
1273 samples[idx]
1274 }
1275
1276 fn latency_quantiles(mut samples_ns: Vec<u128>) -> (u128, u128, u128) {
1277 samples_ns.sort_unstable();
1278 (
1279 quantile_from_sorted(&samples_ns, 50),
1280 quantile_from_sorted(&samples_ns, 95),
1281 quantile_from_sorted(&samples_ns, 99),
1282 )
1283 }
1284
1285 #[test]
1286 fn asupersync_join_admission_optimized_path_is_isomorphic_to_baseline() {
1287 let policy = RuntimePolicy::hardened(Some(1024));
1288 let mut optimized = EvidenceLedger::new();
1289 let mut baseline = EvidenceLedger::new();
1290
1291 for seed in 0_usize..256 {
1292 let rows = if seed % 2 == 0 {
1293 512 + seed
1294 } else {
1295 4096 + seed
1296 };
1297 let optimized_action = policy.decide_join_admission(rows, &mut optimized);
1298 let baseline_action = decide_join_admission_baseline(&policy, rows, &mut baseline);
1299 assert_eq!(optimized_action, baseline_action);
1300
1301 let optimized_record = optimized.records().last().expect("optimized record");
1302 let baseline_record = baseline.records().last().expect("baseline record");
1303 assert_join_record_equivalent(optimized_record, baseline_record);
1304 }
1305 }
1306
1307 #[test]
1308 fn asupersync_join_admission_profile_snapshot_reports_allocation_delta() {
1309 const ITERATIONS: usize = 256;
1310 let policy = RuntimePolicy::hardened(Some(2048));
1311 let mut optimized = EvidenceLedger::new();
1312 let mut baseline = EvidenceLedger::new();
1313 let mut optimized_ns = Vec::with_capacity(ITERATIONS);
1314 let mut baseline_ns = Vec::with_capacity(ITERATIONS);
1315
1316 for seed in 0_usize..ITERATIONS {
1317 let rows = if seed % 3 == 0 {
1318 1024 + seed
1319 } else {
1320 8192 + seed
1321 };
1322
1323 let baseline_start = Instant::now();
1324 let baseline_action = decide_join_admission_baseline(&policy, rows, &mut baseline);
1325 baseline_ns.push(baseline_start.elapsed().as_nanos());
1326 black_box(baseline_action);
1327
1328 let optimized_start = Instant::now();
1329 let optimized_action = policy.decide_join_admission(rows, &mut optimized);
1330 optimized_ns.push(optimized_start.elapsed().as_nanos());
1331 black_box(optimized_action);
1332 }
1333
1334 for (optimized_record, baseline_record) in
1335 optimized.records().iter().zip(baseline.records())
1336 {
1337 assert_join_record_equivalent(optimized_record, baseline_record);
1338 }
1339
1340 let (baseline_p50_ns, baseline_p95_ns, baseline_p99_ns) = latency_quantiles(baseline_ns);
1341 let (optimized_p50_ns, optimized_p95_ns, optimized_p99_ns) =
1342 latency_quantiles(optimized_ns);
1343 let baseline_name_bytes_per_call =
1344 "estimator_overflow_risk".len() + "memory_budget_signal".len();
1345 let baseline_name_bytes_total = baseline_name_bytes_per_call * ITERATIONS;
1346 let optimized_name_bytes_total = 0_usize;
1347 assert!(baseline_name_bytes_total > optimized_name_bytes_total);
1348
1349 println!(
1350 "asupersync_join_admission_profile_snapshot baseline_ns[p50={baseline_p50_ns},p95={baseline_p95_ns},p99={baseline_p99_ns}] optimized_ns[p50={optimized_p50_ns},p95={optimized_p95_ns},p99={optimized_p99_ns}] name_alloc_bytes_baseline={baseline_name_bytes_total} name_alloc_bytes_optimized={optimized_name_bytes_total}"
1351 );
1352 }
1353
1354 #[test]
1355 fn asupersync_structured_log_contains_required_fields() {
1356 let log = make_structured_log(
1357 "asupersync_structured_log_contains_required_fields",
1358 RuntimeMode::Strict,
1359 42,
1360 "ASUPERSYNC-E/log_schema",
1361 "pass",
1362 );
1363 let value = serde_json::to_value(log).expect("serialize log");
1364 assert_required_log_fields(&value);
1365 }
1366
1367 #[test]
1368 fn asupersync_structured_log_is_deterministic_for_same_inputs() {
1369 let left = make_structured_log(
1370 "asupersync_structured_log_is_deterministic_for_same_inputs",
1371 RuntimeMode::Hardened,
1372 1337,
1373 "ASUPERSYNC-E/log_determinism",
1374 "pass",
1375 );
1376 let right = make_structured_log(
1377 "asupersync_structured_log_is_deterministic_for_same_inputs",
1378 RuntimeMode::Hardened,
1379 1337,
1380 "ASUPERSYNC-E/log_determinism",
1381 "pass",
1382 );
1383 assert_eq!(left, right);
1384 let left_json = serde_json::to_string(&left).expect("left json");
1385 let right_json = serde_json::to_string(&right).expect("right json");
1386 assert_eq!(left_json, right_json);
1387 }
1388
1389 #[test]
1390 fn asupersync_property_strict_unknown_feature_always_rejects() {
1391 let policy = RuntimePolicy::strict();
1392 let mut ledger = EvidenceLedger::new();
1393 let case_id = "asupersync_property_strict_unknown_feature_always_rejects";
1394
1395 for seed in 0_u64..128 {
1396 let action = policy.decide_unknown_feature(
1397 format!("unknown_subject_{seed}"),
1398 format!("unknown_detail_{:08x}", seed.wrapping_mul(37)),
1399 &mut ledger,
1400 );
1401 let log = make_structured_log(
1402 case_id,
1403 RuntimeMode::Strict,
1404 seed,
1405 "ASUPERSYNC-E/strict_unknown_feature_reject",
1406 if action == DecisionAction::Reject {
1407 "pass"
1408 } else {
1409 "fail"
1410 },
1411 );
1412 let log_json = serde_json::to_value(log).expect("serialize log");
1413 assert_required_log_fields(&log_json);
1414 assert_eq!(
1415 action,
1416 DecisionAction::Reject,
1417 "strict mode must reject unknown feature; log={}",
1418 serde_json::to_string(&log_json).expect("json")
1419 );
1420 }
1421
1422 assert_eq!(ledger.records().len(), 128);
1423 }
1424
1425 #[test]
1426 fn asupersync_property_hardened_over_cap_forces_repair() {
1427 let cap = 1024_usize;
1428 let policy = RuntimePolicy::hardened(Some(cap));
1429 let mut ledger = EvidenceLedger::new();
1430 let case_id = "asupersync_property_hardened_over_cap_forces_repair";
1431
1432 for seed in 0_u64..256 {
1433 let rows = if seed % 2 == 0 {
1434 cap + 1 + (seed as usize % 10_000)
1435 } else {
1436 cap.saturating_sub(seed as usize % cap)
1437 };
1438 let action = policy.decide_join_admission(rows, &mut ledger);
1439 let log = make_structured_log(
1440 case_id,
1441 RuntimeMode::Hardened,
1442 seed,
1443 "ASUPERSYNC-E/hardened_join_cap_boundary",
1444 if rows > cap && action == DecisionAction::Repair {
1445 "pass"
1446 } else {
1447 "check"
1448 },
1449 );
1450 let log_json = serde_json::to_value(log).expect("serialize log");
1451 assert_required_log_fields(&log_json);
1452 if rows > cap {
1453 assert_eq!(
1454 action,
1455 DecisionAction::Repair,
1456 "rows over cap must force repair; rows={rows}; log={}",
1457 serde_json::to_string(&log_json).expect("json")
1458 );
1459 }
1460 }
1461 }
1462
1463 #[test]
1464 fn asupersync_property_decision_metrics_are_finite_and_bounded() {
1465 let policy = RuntimePolicy::hardened(Some(2048));
1466 let mut ledger = EvidenceLedger::new();
1467 let case_id = "asupersync_property_decision_metrics_are_finite_and_bounded";
1468
1469 for seed in 0_u64..128 {
1470 let rows = 1 + (seed as usize * 97 % 500_000);
1471 policy.decide_join_admission(rows, &mut ledger);
1472 let record = ledger.records().last().expect("record");
1473 let metrics = &record.metrics;
1474 let posterior = metrics.posterior_compatible;
1475 let bounded = (0.0..=1.0).contains(&posterior);
1476 let finite = metrics
1477 .bayes_factor_compatible_over_incompatible
1478 .is_finite()
1479 && metrics.expected_loss_allow.is_finite()
1480 && metrics.expected_loss_reject.is_finite()
1481 && metrics.expected_loss_repair.is_finite();
1482
1483 let log = make_structured_log(
1484 case_id,
1485 RuntimeMode::Hardened,
1486 seed,
1487 "ASUPERSYNC-E/decision_metrics_finite",
1488 if bounded && finite { "pass" } else { "fail" },
1489 );
1490 let log_json = serde_json::to_value(log).expect("serialize log");
1491 assert_required_log_fields(&log_json);
1492 assert!(bounded, "posterior out of range; log={log_json}");
1493 assert!(finite, "non-finite metrics; log={log_json}");
1494 }
1495 }
1496
1497 #[test]
1498 fn decide_clamps_boundary_priors_to_finite_range() {
1499 for (input_prior, expected_prior) in [
1500 (0.0, super::PRIOR_COMPATIBLE_EPSILON),
1501 (1.0, 1.0 - super::PRIOR_COMPATIBLE_EPSILON),
1502 ] {
1503 let record = super::decide(
1504 RuntimeMode::Strict,
1505 super::CompatibilityIssue {
1506 kind: super::IssueKind::MalformedInput,
1507 subject: "prior_clamp_test".to_owned(),
1508 detail: "boundary prior".to_owned(),
1509 },
1510 input_prior,
1511 super::LossMatrix::default(),
1512 Vec::new(),
1513 );
1514
1515 assert_eq!(
1516 record.prior_compatible, expected_prior,
1517 "prior should be clamped into open interval (0,1)"
1518 );
1519 assert!(
1520 record.metrics.posterior_compatible.is_finite(),
1521 "posterior must remain finite for boundary priors"
1522 );
1523 assert!(
1524 record.metrics.expected_loss_allow.is_finite()
1525 && record.metrics.expected_loss_reject.is_finite()
1526 && record.metrics.expected_loss_repair.is_finite(),
1527 "expected-loss metrics must remain finite for boundary priors"
1528 );
1529 }
1530 }
1531
1532 #[test]
1533 fn decide_normalizes_non_finite_priors_to_neutral() {
1534 for input_prior in [f64::NAN, f64::INFINITY, f64::NEG_INFINITY] {
1535 let record = super::decide(
1536 RuntimeMode::Strict,
1537 super::CompatibilityIssue {
1538 kind: super::IssueKind::MalformedInput,
1539 subject: "prior_clamp_test".to_owned(),
1540 detail: "non-finite prior".to_owned(),
1541 },
1542 input_prior,
1543 super::LossMatrix::default(),
1544 Vec::new(),
1545 );
1546
1547 assert_eq!(
1548 record.prior_compatible, 0.5,
1549 "non-finite priors should normalize to neutral prior"
1550 );
1551 assert!(
1552 record.metrics.posterior_compatible.is_finite(),
1553 "posterior must remain finite for non-finite priors"
1554 );
1555 }
1556 }
1557
1558 #[test]
1559 fn asupersync_adversarial_extreme_join_estimate_remains_repair_and_loggable() {
1560 let policy = RuntimePolicy::hardened(Some(8));
1561 let mut ledger = EvidenceLedger::new();
1562 let action = policy.decide_join_admission(usize::MAX, &mut ledger);
1563 assert_eq!(action, DecisionAction::Repair);
1564 let record = ledger.records().last().expect("record");
1565 assert_eq!(record.mode, RuntimeMode::Hardened);
1566 assert!(
1567 record.issue.detail.contains("estimated_rows="),
1568 "issue detail should include estimated_rows"
1569 );
1570
1571 let log = make_structured_log(
1572 "asupersync_adversarial_extreme_join_estimate_remains_repair_and_loggable",
1573 RuntimeMode::Hardened,
1574 u64::MAX,
1575 "ASUPERSYNC-E/adversarial_extreme_rows",
1576 "pass",
1577 );
1578 let log_json = serde_json::to_value(log).expect("serialize log");
1579 assert_required_log_fields(&log_json);
1580 }
1581
1582 #[test]
1583 fn strict_mode_fails_closed_for_unknown_features() {
1584 let mut ledger = EvidenceLedger::new();
1585 let policy = RuntimePolicy::strict();
1586
1587 let action = policy.decide_unknown_feature("csv", "field=experimental", &mut ledger);
1588 assert_eq!(action, DecisionAction::Reject);
1589 assert_eq!(ledger.records()[0].mode, RuntimeMode::Strict);
1590 }
1591
1592 #[test]
1593 fn hardened_mode_repairs_large_join_estimates() {
1594 let mut ledger = EvidenceLedger::new();
1595 let policy = RuntimePolicy::hardened(Some(10_000));
1596
1597 let action = policy.decide_join_admission(100_000, &mut ledger);
1598 assert_eq!(action, DecisionAction::Repair);
1599 assert_eq!(ledger.records().len(), 1);
1600 }
1601
1602 #[test]
1603 fn source_backed_raptorq_envelope_records_manifest_fields() {
1604 let mut source = vec![7_u8; super::DEFAULT_RAPTORQ_SYMBOL_BYTES];
1605 source.extend_from_slice(b"tail");
1606
1607 let envelope = RaptorQEnvelope::from_source_bytes("packet-001", "conformance", &source, 3);
1608
1609 assert_eq!(envelope.artifact_id, "packet-001");
1610 assert_eq!(envelope.artifact_type, "conformance");
1611 assert!(envelope.source_hash.starts_with("sha256:"));
1612 assert_eq!(envelope.source_hash.len(), "sha256:".len() + 64);
1613 assert_eq!(envelope.raptorq.k, 2);
1614 assert_eq!(envelope.raptorq.repair_symbols, 3);
1615 assert_eq!(envelope.raptorq.overhead_ratio, 1.5);
1616 assert_eq!(envelope.raptorq.symbol_hashes.len(), 2);
1617 assert!(
1618 envelope
1619 .raptorq
1620 .symbol_hashes
1621 .iter()
1622 .all(|hash| hash.starts_with("sha256:") && hash.len() == "sha256:".len() + 64)
1623 );
1624 assert_eq!(envelope.scrub.status, "ok");
1625 assert!(envelope.scrub.last_ok_unix_ms > 0);
1626 }
1627
1628 #[test]
1629 fn decode_proof_append_is_capped_and_evicts_oldest() {
1630 let mut envelope =
1631 RaptorQEnvelope::from_source_bytes("packet-001", "conformance", b"source", 1);
1632 let total = super::MAX_DECODE_PROOFS + 5;
1633
1634 for idx in 0..total {
1635 envelope.push_decode_proof_capped(super::DecodeProof {
1636 ts_unix_ms: u64::try_from(idx).expect("idx within u64 range"),
1637 reason: format!("proof-{idx}"),
1638 recovered_blocks: u32::try_from(idx).expect("idx within u32 range"),
1639 proof_hash: format!("sha256:{idx:08x}"),
1640 });
1641 }
1642
1643 assert_eq!(envelope.decode_proofs.len(), super::MAX_DECODE_PROOFS);
1644 assert_eq!(
1645 envelope.decode_proofs[0].proof_hash,
1646 format!("sha256:{:08x}", total - super::MAX_DECODE_PROOFS)
1647 );
1648 assert_eq!(
1649 envelope
1650 .decode_proofs
1651 .last()
1652 .expect("decode proof should exist")
1653 .proof_hash,
1654 format!("sha256:{:08x}", total - 1)
1655 );
1656 }
1657
1658 #[test]
1659 fn decision_card_is_renderable_for_ftui_consumers() {
1660 let mut ledger = EvidenceLedger::new();
1661 let policy = RuntimePolicy::strict();
1662 policy.decide_unknown_feature("csv", "field=experimental", &mut ledger);
1663
1664 let card = decision_to_card(&ledger.records()[0]);
1665 let rendered = card.render_plain();
1666 assert!(rendered.contains("argmin_a"));
1667 assert!(rendered.contains("P(compatible|e)"));
1668 }
1669
1670 #[test]
1673 fn conformal_guard_uncalibrated_accepts_all() {
1674 let mut guard = ConformalGuard::new(100, 0.1);
1675 assert!(!guard.is_calibrated());
1676
1677 let mut ledger = EvidenceLedger::new();
1678 let policy = RuntimePolicy::strict();
1679 policy.decide_unknown_feature("test", "detail", &mut ledger);
1680
1681 let ps = guard.evaluate(&ledger.records()[0]);
1682 assert!(ps.bayesian_action_in_set);
1683 assert_eq!(ps.admissible_actions.len(), 3); assert_eq!(ps.quantile_threshold, f64::INFINITY);
1685 }
1686
1687 #[test]
1688 fn conformal_guard_calibrates_after_sufficient_data() {
1689 let mut guard = ConformalGuard::new(100, 0.1);
1690 let mut ledger = EvidenceLedger::new();
1691 let policy = RuntimePolicy::hardened(Some(100_000));
1692
1693 for _ in 0..10 {
1695 policy.decide_join_admission(50_000, &mut ledger);
1696 }
1697
1698 for record in ledger.records() {
1699 guard.evaluate(record);
1700 }
1701
1702 assert!(guard.is_calibrated());
1703 assert!(guard.conformal_quantile().is_some());
1704 assert_eq!(guard.calibration_count(), 10);
1705 }
1706
1707 #[test]
1708 fn conformal_guard_rolling_window_evicts_old_scores() {
1709 let mut guard = ConformalGuard::new(5, 0.1);
1710 let mut ledger = EvidenceLedger::new();
1711 let policy = RuntimePolicy::hardened(Some(100_000));
1712
1713 for _ in 0..10 {
1714 policy.decide_join_admission(1000, &mut ledger);
1715 }
1716
1717 for record in ledger.records() {
1718 guard.evaluate(record);
1719 }
1720
1721 assert_eq!(guard.calibration_count(), 5);
1723 }
1724
1725 #[test]
1726 fn conformal_guard_coverage_tracking() {
1727 let mut guard = ConformalGuard::new(50, 0.1);
1728 let mut ledger = EvidenceLedger::new();
1729 let policy = RuntimePolicy::hardened(Some(100_000));
1730
1731 for _ in 0..20 {
1733 policy.decide_join_admission(1000, &mut ledger);
1734 }
1735
1736 for record in ledger.records() {
1737 guard.evaluate(record);
1738 }
1739
1740 let coverage = guard.empirical_coverage();
1742 assert!(coverage > 0.5, "coverage should be reasonable: {coverage}");
1743 }
1744
1745 #[test]
1746 fn conformal_guard_no_coverage_alert_under_100_decisions() {
1747 let mut guard = ConformalGuard::new(100, 0.1);
1748 let mut ledger = EvidenceLedger::new();
1749 let policy = RuntimePolicy::hardened(Some(100_000));
1750
1751 for _ in 0..10 {
1752 policy.decide_join_admission(1000, &mut ledger);
1753 }
1754 for record in ledger.records() {
1755 guard.evaluate(record);
1756 }
1757
1758 assert!(!guard.coverage_alert());
1760 }
1761
1762 #[test]
1763 fn conformal_guard_zero_window_size_is_clamped() {
1764 let mut guard = ConformalGuard::new(0, 0.1);
1765 let mut ledger = EvidenceLedger::new();
1766 let policy = RuntimePolicy::hardened(Some(100_000));
1767
1768 policy.decide_join_admission(1000, &mut ledger);
1769 let set = guard.evaluate(&ledger.records()[0]);
1770
1771 assert!(set.bayesian_action_in_set);
1772 assert_eq!(guard.calibration_count(), 1);
1773 }
1774
1775 #[test]
1776 fn conformal_guard_non_finite_alpha_uses_default() {
1777 let guard = ConformalGuard::new(100, f64::NAN);
1778 assert_eq!(guard.alpha, super::DEFAULT_CONFORMAL_ALPHA);
1779 assert!(!guard.coverage_alert());
1780 }
1781
1782 #[test]
1783 fn conformal_guard_repairs_deserialized_zero_window_before_evaluate() {
1784 let mut guard: ConformalGuard = serde_json::from_str(
1785 r#"{"scores":[],"window_size":0,"alpha":0.1,"in_set_count":0,"total_count":0}"#,
1786 )
1787 .expect("deserialize guard");
1788 let mut ledger = EvidenceLedger::new();
1789 let policy = RuntimePolicy::hardened(Some(100_000));
1790
1791 policy.decide_join_admission(1000, &mut ledger);
1792 let set = guard.evaluate(&ledger.records()[0]);
1793
1794 assert!(set.bayesian_action_in_set);
1795 assert_eq!(guard.window_size, 1);
1796 assert_eq!(guard.calibration_count(), 1);
1797 }
1798
1799 #[test]
1800 fn conformal_quantile_ignores_non_finite_persisted_scores() {
1801 let guard = ConformalGuard {
1802 scores: vec![f64::NAN, f64::INFINITY, 1.0, 2.0],
1803 window_size: 10,
1804 alpha: f64::NAN,
1805 in_set_count: 5,
1806 total_count: 3,
1807 };
1808
1809 assert!(guard.is_calibrated());
1810 assert_eq!(guard.conformal_quantile(), Some(2.0));
1811 assert_eq!(guard.empirical_coverage(), 1.0);
1812 }
1813
1814 #[test]
1815 fn conformal_guard_quantile_is_deterministic() {
1816 let mut guard = ConformalGuard::new(100, 0.1);
1817 let mut ledger = EvidenceLedger::new();
1818 let policy = RuntimePolicy::hardened(Some(100_000));
1819
1820 for _ in 0..5 {
1821 policy.decide_join_admission(1000, &mut ledger);
1822 }
1823 for record in ledger.records() {
1824 guard.evaluate(record);
1825 }
1826
1827 let q1 = guard.conformal_quantile();
1828 let q2 = guard.conformal_quantile();
1829 assert_eq!(q1, q2);
1830 }
1831
1832 #[test]
1836 fn conformal_quantile_basic() {
1837 let mut guard = ConformalGuard::new(100, 0.1);
1838 let mut ledger = EvidenceLedger::new();
1840 let policy = RuntimePolicy::hardened(Some(100_000));
1841
1842 for _ in 0..5 {
1844 policy.decide_join_admission(1000, &mut ledger);
1845 }
1846 for record in ledger.records() {
1847 guard.evaluate(record);
1848 }
1849
1850 let q = guard.conformal_quantile();
1851 assert!(q.is_some());
1852 let quantile = q.unwrap();
1853 assert!(quantile.is_finite(), "quantile must be finite: {quantile}");
1854 assert!(quantile >= 0.0, "quantile must be non-negative: {quantile}");
1855 }
1856
1857 #[test]
1859 fn conformal_quantile_trivial() {
1860 let mut guard = ConformalGuard::new(100, 0.1);
1861 let mut ledger = EvidenceLedger::new();
1862 let policy = RuntimePolicy::hardened(Some(100_000));
1863
1864 policy.decide_join_admission(1000, &mut ledger);
1866 policy.decide_join_admission(1000, &mut ledger);
1867 guard.evaluate(&ledger.records()[0]);
1868 guard.evaluate(&ledger.records()[1]);
1869
1870 let q = guard.conformal_quantile();
1871 assert!(q.is_some());
1872 }
1873
1874 #[test]
1876 fn conformal_quantile_empty() {
1877 let guard = ConformalGuard::new(100, 0.1);
1878 assert!(guard.conformal_quantile().is_none());
1879 assert!(!guard.is_calibrated());
1880 }
1881
1882 #[test]
1885 fn conformal_guard_agrees_with_bayesian() {
1886 let mut guard = ConformalGuard::new(100, 0.1);
1887 let mut ledger = EvidenceLedger::new();
1888 let policy = RuntimePolicy::hardened(Some(100_000));
1889
1890 for _ in 0..20 {
1892 policy.decide_join_admission(1000, &mut ledger);
1893 }
1894
1895 let mut bayesian_agreed = 0;
1896 let mut total = 0;
1897
1898 for record in ledger.records() {
1899 let ps = guard.evaluate(record);
1900 total += 1;
1901 if ps.bayesian_action_in_set && ps.admissible_actions.len() == 1 {
1902 assert_eq!(ps.admissible_actions[0], record.action);
1904 bayesian_agreed += 1;
1905 }
1906 }
1907
1908 assert!(total > 0, "should have evaluated at least one decision");
1910 assert!(
1912 bayesian_agreed > 0 || total < 3,
1913 "at least some decisions should agree with Bayesian"
1914 );
1915 }
1916
1917 #[test]
1920 fn conformal_guard_widens_on_uncertainty() {
1921 let mut guard = ConformalGuard::new(10, 0.1);
1922 let mut ledger = EvidenceLedger::new();
1923 let policy = RuntimePolicy::hardened(Some(100_000));
1924
1925 for _ in 0..10 {
1927 policy.decide_join_admission(100, &mut ledger);
1928 }
1929 for record in ledger.records() {
1930 guard.evaluate(record);
1931 }
1932
1933 let mut outlier_ledger = EvidenceLedger::new();
1935 let extreme_policy = RuntimePolicy::hardened(Some(10));
1936 extreme_policy.decide_join_admission(1_000_000, &mut outlier_ledger);
1937
1938 let ps = guard.evaluate(&outlier_ledger.records()[0]);
1939 if !ps.bayesian_action_in_set {
1941 assert_eq!(
1942 ps.admissible_actions.len(),
1943 3,
1944 "widened set should admit all actions"
1945 );
1946 }
1947 }
1948
1949 #[test]
1951 fn conformal_coverage_guarantee_1000_decisions() {
1952 let mut guard = ConformalGuard::new(1000, 0.1);
1953 let mut ledger = EvidenceLedger::new();
1954 let policy = RuntimePolicy::hardened(Some(100_000));
1955
1956 for i in 0..1000 {
1958 let rows = 1000 + (i * 7) % 500;
1960 policy.decide_join_admission(rows, &mut ledger);
1961 }
1962
1963 for record in ledger.records() {
1964 guard.evaluate(record);
1965 }
1966
1967 let coverage = guard.empirical_coverage();
1970 assert!(
1971 coverage >= 0.7,
1972 "coverage {coverage} should be >= 0.7 (relaxed bound for finite sample)"
1973 );
1974 }
1975
1976 #[test]
1978 fn conformal_rolling_window_exact_eviction() {
1979 let window_size = 5;
1980 let mut guard = ConformalGuard::new(window_size, 0.1);
1981 let mut ledger = EvidenceLedger::new();
1982 let policy = RuntimePolicy::hardened(Some(100_000));
1983
1984 for _ in 0..15 {
1986 policy.decide_join_admission(1000, &mut ledger);
1987 }
1988
1989 for record in ledger.records() {
1990 guard.evaluate(record);
1991 }
1992
1993 assert_eq!(
1994 guard.calibration_count(),
1995 window_size,
1996 "window should be exactly {window_size}"
1997 );
1998 }
1999
2000 #[test]
2003 fn conformal_galaxy_brain_card_content() {
2004 let mut ledger = EvidenceLedger::new();
2005 let policy = RuntimePolicy::hardened(Some(100_000));
2006 policy.decide_join_admission(50_000, &mut ledger);
2007
2008 let card = decision_to_card(&ledger.records()[0]);
2009 assert!(card.equation.contains("argmin_a"));
2010 assert!(card.substitution.contains("P(compatible|e)"));
2011 assert!(card.substitution.contains("E[allow]"));
2012 assert!(card.substitution.contains("E[reject]"));
2013 assert!(card.substitution.contains("E[repair]"));
2014 }
2015
2016 #[test]
2017 fn conformal_prediction_set_serializes() {
2018 let mut guard = ConformalGuard::new(100, 0.1);
2019 let mut ledger = EvidenceLedger::new();
2020 let policy = RuntimePolicy::hardened(Some(100_000));
2021 policy.decide_join_admission(1000, &mut ledger);
2022
2023 let ps = guard.evaluate(&ledger.records()[0]);
2024 let json = serde_json::to_string(&ps).expect("serialize");
2025 let _: serde_json::Value = serde_json::from_str(&json).expect("valid JSON");
2026 assert!(json.contains("quantile_threshold"));
2027 assert!(json.contains("empirical_coverage"));
2028 }
2029}