use std::collections::BTreeMap;
use std::num::NonZeroUsize;
use std::ops::Bound::Excluded;
use std::ops::Bound::Unbounded;
use serde::Deserialize;
use serde::Serialize;
use crate::ApplyOutcome;
use crate::Authentication;
use crate::CreditPolicy;
use crate::CreditRecord;
use crate::CreditScore;
use crate::MeasureError;
use crate::MeasurementBatch;
use crate::MeasurementEvent;
use crate::MeasurementSnapshot;
use crate::ReliabilityClass;
use crate::ReliabilityEvidence;
use crate::ReliabilityPolicy;
use crate::ReliabilityWindow;
use crate::SnapshotRecord;
use crate::UnixTime;
use crate::MEASUREMENT_SNAPSHOT_VERSION;
#[allow(
clippy::unwrap_used,
reason = "the non-zero integer literal is validated during const evaluation"
)]
pub const DEFAULT_MAX_RETAINED_PEERS: NonZeroUsize = NonZeroUsize::new(16_384).unwrap();
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub struct PeerRecord {
credit: CreditRecord,
reliability: ReliabilityWindow,
}
impl PeerRecord {
pub const fn new(credit: CreditRecord, reliability: ReliabilityWindow) -> Self {
Self {
credit,
reliability,
}
}
pub const fn empty(first_seen: UnixTime) -> Self {
Self::new(
CreditRecord::empty(first_seen),
ReliabilityWindow::new(None, None, ReliabilityEvidence::new(0, 0, 0, 0, 0, 0)),
)
}
pub const fn credit(self) -> CreditRecord {
self.credit
}
pub const fn reliability(self) -> ReliabilityWindow {
self.reliability
}
fn transition_batch(
self,
authentication: Authentication,
batch: MeasurementBatch,
at: UnixTime,
reliability_policy: ReliabilityPolicy,
) -> Result<Self, MeasureError> {
let mut next = self;
next.reliability
.observe_batch(batch, at, reliability_policy)?;
match (authentication.refreshes_peer_observation(), batch.event()) {
(true, MeasurementEvent::Sent { useful_bytes }) => {
next.credit.record_sent(useful_bytes, at)?;
}
(true, MeasurementEvent::Received { useful_bytes }) => {
next.credit.record_received(useful_bytes, at)?;
}
(
true,
MeasurementEvent::Connected
| MeasurementEvent::Disconnected
| MeasurementEvent::FailedToSend
| MeasurementEvent::FailedToReceive,
) => next.credit.touch(at)?,
(false, _) => next.credit.ensure_not_before_last_seen(at)?,
}
Ok(next)
}
fn validate_snapshot(self) -> Result<(), MeasureError> {
match (
self.reliability.epoch_start(),
self.reliability.window_seconds(),
self.reliability.stored_evidence().is_unobserved(),
) {
(None, _, false) => Err(MeasureError::SnapshotEvidenceWithoutEpoch),
(Some(_), _, true) => Err(MeasureError::SnapshotEpochWithoutEvidence),
(Some(_), None, false) => Err(MeasureError::SnapshotReliabilityWindowMissing),
(None, Some(_), true) => Err(MeasureError::SnapshotReliabilityWindowWithoutEpoch),
(Some(_), Some(0), false) => Err(MeasureError::SnapshotReliabilityWindowMissing),
(Some(epoch_start), Some(window), false) if epoch_start.as_secs() % window != 0 => {
Err(MeasureError::SnapshotReliabilityEpochMisaligned)
}
_ => Ok(()),
}
}
fn reconcile_runtime(
&mut self,
now: UnixTime,
policy: ReliabilityPolicy,
) -> RecordReconciliation {
let reliability_reset = self.reliability.reset_for_policy_change(policy);
let credit_adjusted = self.credit.reconcile_clock(now);
let reliability_adjusted = self.reliability.reconcile_clock(now, policy);
RecordReconciliation {
clock_adjusted: credit_adjusted || reliability_adjusted,
reliability_reset,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
struct RecordReconciliation {
clock_adjusted: bool,
reliability_reset: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ApplyReport<P> {
outcome: ApplyOutcome,
evicted_peer: Option<P>,
}
impl<P> ApplyReport<P> {
const fn applied(evicted_peer: Option<P>) -> Self {
Self {
outcome: ApplyOutcome::Applied,
evicted_peer,
}
}
const fn ignored_unattributable() -> Self {
Self {
outcome: ApplyOutcome::IgnoredUnattributable,
evicted_peer: None,
}
}
const fn ignored_unknown_peer() -> Self {
Self {
outcome: ApplyOutcome::IgnoredUnknownPeer,
evicted_peer: None,
}
}
pub const fn outcome(&self) -> ApplyOutcome {
self.outcome
}
pub const fn evicted_peer(&self) -> Option<&P> {
self.evicted_peer.as_ref()
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct PeerMeasurement<P> {
pub peer: P,
pub credit: CreditRecord,
pub credit_score: CreditScore,
pub reliability: ReliabilityEvidence,
pub reliability_class: ReliabilityClass,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PeerMeasurementFailure<P> {
peer: P,
error: MeasureError,
}
impl<P> PeerMeasurementFailure<P> {
fn new(peer: P, error: MeasureError) -> Self {
Self { peer, error }
}
pub const fn peer(&self) -> &P {
&self.peer
}
pub const fn error(&self) -> MeasureError {
self.error
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct MeasurementProjection<P> {
measurements: Vec<PeerMeasurement<P>>,
failures: Vec<PeerMeasurementFailure<P>>,
}
#[derive(Debug, Clone, PartialEq)]
pub struct MeasurementPage<P> {
measurements: Vec<PeerMeasurement<P>>,
failures: Vec<PeerMeasurementFailure<P>>,
next_cursor: Option<P>,
}
impl<P> MeasurementPage<P> {
pub fn measurements(&self) -> &[PeerMeasurement<P>] {
&self.measurements
}
pub fn failures(&self) -> &[PeerMeasurementFailure<P>] {
&self.failures
}
pub const fn next_cursor(&self) -> Option<&P> {
self.next_cursor.as_ref()
}
pub fn into_parts(
self,
) -> (
Vec<PeerMeasurement<P>>,
Vec<PeerMeasurementFailure<P>>,
Option<P>,
) {
(self.measurements, self.failures, self.next_cursor)
}
}
impl<P> MeasurementProjection<P> {
pub fn measurements(&self) -> &[PeerMeasurement<P>] {
&self.measurements
}
pub fn failures(&self) -> &[PeerMeasurementFailure<P>] {
&self.failures
}
pub fn into_parts(self) -> (Vec<PeerMeasurement<P>>, Vec<PeerMeasurementFailure<P>>) {
(self.measurements, self.failures)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PruneReport<P> {
removed: Vec<P>,
failures: Vec<PeerMeasurementFailure<P>>,
}
impl<P> PruneReport<P> {
pub fn removed(&self) -> &[P] {
&self.removed
}
pub fn failures(&self) -> &[PeerMeasurementFailure<P>] {
&self.failures
}
pub fn removed_count(&self) -> usize {
self.removed.len()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct LedgerReconciliation {
clock_adjusted_records: usize,
reliability_reset_records: usize,
}
impl LedgerReconciliation {
pub const fn clock_adjusted_records(self) -> usize {
self.clock_adjusted_records
}
pub const fn reliability_reset_records(self) -> usize {
self.reliability_reset_records
}
pub const fn is_adjusted(self) -> bool {
self.clock_adjusted_records > 0 || self.reliability_reset_records > 0
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MeasurementLedger<P> {
records: BTreeMap<P, PeerRecord>,
max_records: NonZeroUsize,
}
impl<P> Default for MeasurementLedger<P>
where P: Ord
{
fn default() -> Self {
Self::new()
}
}
impl<P> MeasurementLedger<P>
where P: Ord
{
pub const fn new() -> Self {
Self {
records: BTreeMap::new(),
max_records: DEFAULT_MAX_RETAINED_PEERS,
}
}
pub const fn with_max_records(max_records: NonZeroUsize) -> Self {
Self {
records: BTreeMap::new(),
max_records,
}
}
pub const fn max_records(&self) -> usize {
self.max_records.get()
}
pub fn len(&self) -> usize {
self.records.len()
}
pub fn is_empty(&self) -> bool {
self.records.is_empty()
}
pub fn apply(
&mut self,
peer: P,
authentication: Authentication,
event: MeasurementEvent,
at: UnixTime,
reliability_policy: ReliabilityPolicy,
) -> Result<ApplyReport<P>, MeasureError>
where
P: Clone,
{
self.apply_batch(
peer,
authentication,
MeasurementBatch::single(event),
at,
reliability_policy,
)
}
pub fn apply_batch(
&mut self,
peer: P,
authentication: Authentication,
batch: MeasurementBatch,
at: UnixTime,
reliability_policy: ReliabilityPolicy,
) -> Result<ApplyReport<P>, MeasureError>
where
P: Clone,
{
if !authentication.permits(batch.event()) {
return Ok(ApplyReport::ignored_unattributable());
}
let is_new_peer = !self.records.contains_key(&peer);
if is_new_peer && !authentication.establishes_peer() {
return Ok(ApplyReport::ignored_unknown_peer());
}
let current = self
.records
.get(&peer)
.copied()
.unwrap_or_else(|| PeerRecord::empty(at));
let next = current.transition_batch(authentication, batch, at, reliability_policy)?;
let evicted_peer = if is_new_peer && self.records.len() >= self.max_records.get() {
let stalest = self
.records
.iter()
.min_by(|(left_peer, left_record), (right_peer, right_record)| {
left_record
.credit
.last_seen()
.cmp(&right_record.credit.last_seen())
.then_with(|| left_peer.cmp(right_peer))
})
.map(|(oldest_peer, _)| oldest_peer.clone());
if let Some(stalest_peer) = &stalest {
self.records.remove(stalest_peer);
}
stalest
} else {
None
};
self.records.insert(peer, next);
Ok(ApplyReport::applied(evicted_peer))
}
pub fn record(&self, peer: &P) -> Option<PeerRecord> {
self.records.get(peer).copied()
}
pub fn next_retention_boundary(&self, credit_policy: CreditPolicy) -> Option<UnixTime> {
self.records
.values()
.map(|record| {
UnixTime::from_secs(
record
.credit
.last_seen()
.as_secs()
.saturating_add(credit_policy.retention_seconds()),
)
})
.min()
}
pub fn from_snapshot(snapshot: MeasurementSnapshot<P>) -> Result<Self, MeasureError> {
Self::from_snapshot_with_max_records(snapshot, DEFAULT_MAX_RETAINED_PEERS)
}
pub fn from_snapshot_with_max_records(
snapshot: MeasurementSnapshot<P>,
max_records: NonZeroUsize,
) -> Result<Self, MeasureError> {
if snapshot.schema_version != MEASUREMENT_SNAPSHOT_VERSION {
return Err(MeasureError::UnsupportedSnapshotVersion {
found: snapshot.schema_version,
});
}
if snapshot.records.len() > max_records.get() {
return Err(MeasureError::SnapshotPeerLimitExceeded {
found: snapshot.records.len(),
max: max_records.get(),
});
}
let mut records = BTreeMap::new();
for entry in snapshot.records {
entry.record.validate_snapshot()?;
if records.insert(entry.peer, entry.record).is_some() {
return Err(MeasureError::DuplicatePeerInSnapshot);
}
}
Ok(Self {
records,
max_records,
})
}
}
impl<P> MeasurementLedger<P>
where P: Clone + Ord
{
pub fn measurement(
&self,
peer: &P,
now: UnixTime,
credit_policy: CreditPolicy,
reliability_policy: ReliabilityPolicy,
) -> Result<Option<PeerMeasurement<P>>, MeasureError> {
self.records
.get(peer)
.copied()
.map(|record| {
project_record(peer.clone(), record, now, credit_policy, reliability_policy)
})
.transpose()
}
pub fn measurements(
&self,
now: UnixTime,
credit_policy: CreditPolicy,
reliability_policy: ReliabilityPolicy,
) -> MeasurementProjection<P> {
let mut measurements = Vec::with_capacity(self.records.len());
let mut failures = Vec::new();
for (peer, record) in &self.records {
match project_record(
peer.clone(),
*record,
now,
credit_policy,
reliability_policy,
) {
Ok(measurement) => measurements.push(measurement),
Err(error) => {
failures.push(PeerMeasurementFailure::new(peer.clone(), error));
}
}
}
MeasurementProjection {
measurements,
failures,
}
}
pub fn measurements_page(
&self,
after: Option<&P>,
limit: NonZeroUsize,
now: UnixTime,
credit_policy: CreditPolicy,
reliability_policy: ReliabilityPolicy,
) -> MeasurementPage<P> {
let mut records = match after {
Some(peer) => self.records.range((Excluded(peer), Unbounded)),
None => self.records.range(..),
};
let mut measurements = Vec::with_capacity(self.records.len().min(limit.get()));
let mut failures = Vec::new();
let mut last_scanned = None;
for _ in 0..limit.get() {
let Some((peer, record)) = records.next() else {
break;
};
last_scanned = Some(peer.clone());
match project_record(
peer.clone(),
*record,
now,
credit_policy,
reliability_policy,
) {
Ok(measurement) => measurements.push(measurement),
Err(error) => failures.push(PeerMeasurementFailure::new(peer.clone(), error)),
}
}
let next_cursor = if records.next().is_some() {
last_scanned
} else {
None
};
MeasurementPage {
measurements,
failures,
next_cursor,
}
}
pub fn prune(&mut self, now: UnixTime, credit_policy: CreditPolicy) -> PruneReport<P> {
let mut expired = Vec::new();
let mut failures = Vec::new();
for (peer, record) in &self.records {
match record.credit.is_expired(now, credit_policy) {
Ok(true) => expired.push(peer.clone()),
Ok(false) => {}
Err(error) => {
failures.push(PeerMeasurementFailure::new(peer.clone(), error));
}
}
}
for peer in &expired {
self.records.remove(peer);
}
PruneReport {
removed: expired,
failures,
}
}
pub fn reconcile_runtime(
&mut self,
now: UnixTime,
reliability_policy: ReliabilityPolicy,
) -> LedgerReconciliation {
let mut reconciliation = LedgerReconciliation::default();
for record in self.records.values_mut() {
let adjusted = record.reconcile_runtime(now, reliability_policy);
reconciliation.clock_adjusted_records += usize::from(adjusted.clock_adjusted);
reconciliation.reliability_reset_records += usize::from(adjusted.reliability_reset);
}
reconciliation
}
pub fn snapshot(&self) -> MeasurementSnapshot<P> {
MeasurementSnapshot {
schema_version: MEASUREMENT_SNAPSHOT_VERSION,
records: self
.records
.iter()
.map(|(peer, record)| SnapshotRecord {
peer: peer.clone(),
record: *record,
})
.collect(),
}
}
}
fn project_record<P>(
peer: P,
record: PeerRecord,
now: UnixTime,
credit_policy: CreditPolicy,
reliability_policy: ReliabilityPolicy,
) -> Result<PeerMeasurement<P>, MeasureError> {
record.credit.ensure_not_before_last_seen(now)?;
let reliability = record.reliability.evidence_at(now, reliability_policy)?;
Ok(PeerMeasurement {
peer,
credit: record.credit,
credit_score: record.credit.score(credit_policy),
reliability,
reliability_class: reliability.classify_with_policy(reliability_policy),
})
}
#[cfg(test)]
#[allow(clippy::panic)]
mod tests;