use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Duration;
use super::continuity::AttestedStatus;
use super::identity::{
CanonicalConstraints, CapabilityId, ConstraintError, Digest256, WorkLatencyEnvelope,
};
pub const DEFAULT_ATTESTATION_CADENCE_FLOOR: Duration = Duration::from_millis(50);
#[derive(Clone, Copy, Debug)]
pub struct EvaluationRequest<'a> {
pub capability_id: &'a CapabilityId,
pub constraints: &'a CanonicalConstraints,
pub work_latency: &'a WorkLatencyEnvelope,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum ReadinessEvaluation {
Ready {
estimated_start: Option<Duration>,
},
NotReady {
reason: u16,
},
UnsupportedPredicate,
TemporarilyUnevaluable,
InvalidConstraints,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, serde::Serialize, serde::Deserialize)]
pub enum StatusReason {
None,
Provider(u16),
UnsupportedPredicate,
TemporarilyUnevaluable,
InvalidConstraints,
SamplingIntervalUnsupported,
}
pub const fn project_evaluation(
evaluation: &ReadinessEvaluation,
) -> (AttestedStatus, StatusReason) {
match evaluation {
ReadinessEvaluation::Ready { .. } => (AttestedStatus::Ready, StatusReason::None),
ReadinessEvaluation::NotReady { reason } => {
(AttestedStatus::NotReady, StatusReason::Provider(*reason))
}
ReadinessEvaluation::UnsupportedPredicate => (
AttestedStatus::ProviderUnknown,
StatusReason::UnsupportedPredicate,
),
ReadinessEvaluation::TemporarilyUnevaluable => (
AttestedStatus::ProviderUnknown,
StatusReason::TemporarilyUnevaluable,
),
ReadinessEvaluation::InvalidConstraints => (
AttestedStatus::ProviderUnknown,
StatusReason::InvalidConstraints,
),
}
}
pub trait ReadinessEvaluator {
fn evaluate(&self, request: &EvaluationRequest<'_>) -> ReadinessEvaluation;
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub struct CadenceRefusal {
pub minimum_supported: Duration,
}
impl CadenceRefusal {
pub const fn as_status(&self) -> (AttestedStatus, StatusReason) {
(
AttestedStatus::ProviderUnknown,
StatusReason::SamplingIntervalUnsupported,
)
}
}
pub const fn check_cadence(
requested_strictest: Duration,
floor: Duration,
) -> Result<(), CadenceRefusal> {
if requested_strictest.as_nanos() < floor.as_nanos() {
Err(CadenceRefusal {
minimum_supported: floor,
})
} else {
Ok(())
}
}
#[derive(Default, Debug)]
pub struct SensingCounters {
pub invalid_constraints: AtomicU64,
pub protocol_invalid: AtomicU64,
pub cadence_refusals: AtomicU64,
pub scope_refusals: AtomicU64,
pub broad_selector_refusals: AtomicU64,
pub org_cert_invalid: AtomicU64,
pub org_below_floor: AtomicU64,
pub org_foreign_org: AtomicU64,
pub org_sender_member_mismatch: AtomicU64,
pub org_audience_mismatch: AtomicU64,
pub org_authority_unavailable: AtomicU64,
pub org_stale_stamp: AtomicU64,
pub org_store_poisoned: AtomicU64,
pub interests_registered: AtomicU64,
pub interests_coalesced: AtomicU64,
pub candidate_fanout_total: AtomicU64,
pub attestations_emitted: AtomicU64,
pub attestations_forwarded: AtomicU64,
pub attestations_gated: AtomicU64,
pub attestations_superseded: AtomicU64,
pub provider_free_registrations: AtomicU64,
pub divergent_resolution_merge_miss: AtomicU64,
}
impl SensingCounters {
pub fn get(counter: &AtomicU64) -> u64 {
counter.load(Ordering::Relaxed)
}
}
pub fn validate_interest_constraints(
bytes: &[u8],
claimed: &Digest256,
counters: &SensingCounters,
) -> Result<CanonicalConstraints, ConstraintError> {
match CanonicalConstraints::validate_inline(bytes, claimed) {
Ok(constraints) => Ok(constraints),
Err(error) => {
counters.invalid_constraints.fetch_add(1, Ordering::Relaxed);
if error.is_security_relevant() {
counters.protocol_invalid.fetch_add(1, Ordering::Relaxed);
}
Err(error)
}
}
}
#[cfg(test)]
mod tests {
use super::*;
struct LoadEvaluator {
current_load: u16,
}
impl ReadinessEvaluator for LoadEvaluator {
fn evaluate(&self, request: &EvaluationRequest<'_>) -> ReadinessEvaluation {
let Some(max_load) = request.constraints.get("max_load") else {
return ReadinessEvaluation::UnsupportedPredicate;
};
let Ok(max_load) = max_load.parse::<u16>() else {
return ReadinessEvaluation::InvalidConstraints;
};
if self.current_load <= max_load {
ReadinessEvaluation::Ready {
estimated_start: Some(Duration::from_millis(5)),
}
} else {
ReadinessEvaluation::NotReady { reason: 42 }
}
}
}
fn request<'a>(
capability_id: &'a CapabilityId,
constraints: &'a CanonicalConstraints,
work_latency: &'a WorkLatencyEnvelope,
) -> EvaluationRequest<'a> {
EvaluationRequest {
capability_id,
constraints,
work_latency,
}
}
#[test]
fn evaluator_contract_round_trips_through_a_real_impl() {
let id = CapabilityId::new("job.run");
let latency = WorkLatencyEnvelope::start_within(Duration::from_millis(100));
let ok = CanonicalConstraints::from_entries([("max_load", "50")]).unwrap();
let alien = CanonicalConstraints::from_entries([("gpu_class", "h100")]).unwrap();
let idle = LoadEvaluator { current_load: 10 };
let busy = LoadEvaluator { current_load: 90 };
assert_eq!(
idle.evaluate(&request(&id, &ok, &latency)),
ReadinessEvaluation::Ready {
estimated_start: Some(Duration::from_millis(5)),
},
);
assert_eq!(
busy.evaluate(&request(&id, &ok, &latency)),
ReadinessEvaluation::NotReady { reason: 42 },
);
assert_eq!(
idle.evaluate(&request(&id, &alien, &latency)),
ReadinessEvaluation::UnsupportedPredicate,
);
}
#[test]
fn projection_collapses_to_provider_unknown_with_distinct_reasons() {
use AttestedStatus as S;
assert_eq!(
project_evaluation(&ReadinessEvaluation::Ready {
estimated_start: None,
}),
(S::Ready, StatusReason::None),
);
assert_eq!(
project_evaluation(&ReadinessEvaluation::NotReady { reason: 7 }),
(S::NotReady, StatusReason::Provider(7)),
);
let unknowns = [
(
ReadinessEvaluation::UnsupportedPredicate,
StatusReason::UnsupportedPredicate,
),
(
ReadinessEvaluation::TemporarilyUnevaluable,
StatusReason::TemporarilyUnevaluable,
),
(
ReadinessEvaluation::InvalidConstraints,
StatusReason::InvalidConstraints,
),
];
for (evaluation, expected_reason) in unknowns {
assert_eq!(
project_evaluation(&evaluation),
(S::ProviderUnknown, expected_reason),
);
}
}
#[test]
fn cadence_below_floor_is_refused_with_the_floor_attached() {
let floor = DEFAULT_ATTESTATION_CADENCE_FLOOR;
assert_eq!(check_cadence(Duration::from_millis(50), floor), Ok(()));
assert_eq!(check_cadence(Duration::from_secs(1), floor), Ok(()));
let refusal = check_cadence(Duration::from_millis(20), floor).unwrap_err();
assert_eq!(refusal.minimum_supported, floor);
assert_eq!(
refusal.as_status(),
(
AttestedStatus::ProviderUnknown,
StatusReason::SamplingIntervalUnsupported,
),
);
}
#[test]
fn digest_mismatch_counts_as_security_plain_decode_failures_do_not() {
let counters = SensingCounters::default();
let constraints = CanonicalConstraints::from_entries([("a", "1")]).unwrap();
let bytes = constraints.canonical_bytes();
let right = constraints.constraints_digest();
let wrong = Digest256::from_bytes([0u8; 32]);
assert!(validate_interest_constraints(&bytes, &right, &counters).is_ok());
assert_eq!(SensingCounters::get(&counters.invalid_constraints), 0);
assert_eq!(SensingCounters::get(&counters.protocol_invalid), 0);
assert_eq!(
validate_interest_constraints(&bytes, &wrong, &counters),
Err(ConstraintError::DigestMismatch),
);
assert_eq!(SensingCounters::get(&counters.invalid_constraints), 1);
assert_eq!(SensingCounters::get(&counters.protocol_invalid), 1);
assert!(validate_interest_constraints(&bytes[..3], &right, &counters).is_err());
assert_eq!(SensingCounters::get(&counters.invalid_constraints), 2);
assert_eq!(SensingCounters::get(&counters.protocol_invalid), 1);
}
}