use std::time::Duration;
use super::continuity::ProjectedReadiness;
use super::identity::{
AudienceScopeCommitment, ConsumerLatencyBudget, GroupRef, ProviderSelector, ResultMode,
};
#[derive(Clone, Copy, Debug)]
pub struct CandidatePolicy {
pub initial_fanout: usize,
pub standby_count: usize,
pub maximum_fanout: usize,
pub each_mode_max_providers: usize,
}
impl Default for CandidatePolicy {
fn default() -> Self {
Self {
initial_fanout: 1,
standby_count: 1,
maximum_fanout: 3,
each_mode_max_providers: 32,
}
}
}
#[derive(Clone, Debug)]
pub struct TagAssertion {
pub key: String,
pub value: String,
pub asserted_by: AudienceScopeCommitment,
}
#[derive(Clone, Debug)]
pub struct CandidateProvider {
pub node_id: u64,
pub capability_generation: u64,
pub authorized: bool,
pub reachable: bool,
pub route_estimate: Duration,
pub tags: Vec<TagAssertion>,
pub groups: Vec<GroupRef>,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum ResolutionRefusal {
SelectorTooBroad {
matched: usize,
cap: usize,
},
AllBranchesRefused,
QuorumExceedsFanout {
required: usize,
cap: usize,
},
AuthorityMismatch,
AdmittedLegMismatch,
}
#[derive(Clone, PartialEq, Eq, Debug)]
pub struct ResolvedCandidates {
pub active: Vec<u64>,
pub standby: Vec<u64>,
}
pub const fn population_is_boundable(selector: &ProviderSelector) -> bool {
matches!(
selector,
ProviderSelector::Node(_) | ProviderSelector::Nodes(_) | ProviderSelector::Group(_)
)
}
pub fn resolve_candidates(
selector: &ProviderSelector,
result_mode: ResultMode,
snapshot: &[CandidateProvider],
owner_root: &AudienceScopeCommitment,
policy: &CandidatePolicy,
) -> Result<ResolvedCandidates, ResolutionRefusal> {
if let ProviderSelector::Node(id) = selector {
return Ok(ResolvedCandidates {
active: vec![*id],
standby: Vec::new(),
});
}
let mut eligible: Vec<&CandidateProvider> = snapshot
.iter()
.filter(|candidate| candidate.authorized && candidate.reachable)
.filter(|candidate| match selector {
ProviderSelector::AnyAuthorized => true,
ProviderSelector::Node(_) => unreachable!("handled above"),
ProviderSelector::Nodes(ids) => ids.contains(&candidate.node_id),
ProviderSelector::Group(group) => candidate.groups.contains(group),
ProviderSelector::Tags(matches) => matches.iter().all(|wanted| {
candidate.tags.iter().any(|assertion| {
assertion.key == wanted.key
&& assertion.value == wanted.value
&& assertion.asserted_by == *owner_root
})
}),
})
.collect();
eligible.sort_by_key(|candidate| candidate.route_estimate);
let matched = eligible.len();
let active_bound = match result_mode {
ResultMode::Any => policy.initial_fanout,
ResultMode::TopK(k) => (k as usize)
.max(policy.initial_fanout)
.min(policy.maximum_fanout),
ResultMode::Quorum(k) => {
let required = k as usize;
if required > policy.maximum_fanout {
return Err(ResolutionRefusal::QuorumExceedsFanout {
required,
cap: policy.maximum_fanout,
});
}
required
.max(policy.initial_fanout)
.min(policy.maximum_fanout)
}
ResultMode::Each => {
if matched > policy.each_mode_max_providers {
return Err(ResolutionRefusal::SelectorTooBroad {
matched,
cap: policy.each_mode_max_providers,
});
}
matched
}
};
let active: Vec<u64> = eligible
.iter()
.take(active_bound)
.map(|candidate| candidate.node_id)
.collect();
let standby: Vec<u64> = eligible
.iter()
.skip(active.len())
.take(policy.standby_count)
.map(|candidate| candidate.node_id)
.collect();
Ok(ResolvedCandidates { active, standby })
}
#[derive(Clone, Copy, Debug)]
pub struct BranchView {
pub provider: u64,
pub projection: ProjectedReadiness,
pub estimated_start: Option<Duration>,
pub route_estimate: Duration,
}
#[derive(Clone, PartialEq, Eq, Debug)]
pub enum AggregateView {
Scalar {
status: ProjectedReadiness,
supporting: Vec<u64>,
},
PerProvider(Vec<(u64, ProjectedReadiness)>),
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum BranchViability {
Viable(Duration),
Potential,
NonViable,
}
pub fn classify_branch(branch: &BranchView, budget: &ConsumerLatencyBudget) -> BranchViability {
match branch.projection {
ProjectedReadiness::Ready
if budget.admits(branch.route_estimate, branch.estimated_start) =>
{
BranchViability::Viable(
branch
.route_estimate
.saturating_add(branch.estimated_start.unwrap_or(Duration::ZERO)),
)
}
ProjectedReadiness::NotReady => BranchViability::NonViable,
_ => BranchViability::Potential,
}
}
pub fn project_aggregate(
selector: &ProviderSelector,
result_mode: ResultMode,
budget: &ConsumerLatencyBudget,
branches: &[BranchView],
search_complete: bool,
) -> AggregateView {
if result_mode == ResultMode::Each {
return AggregateView::PerProvider(
branches
.iter()
.map(|branch| {
let projection = if branch.projection == ProjectedReadiness::Ready
&& !matches!(classify_branch(branch, budget), BranchViability::Viable(_))
{
ProjectedReadiness::Unknown
} else {
branch.projection
};
(branch.provider, projection)
})
.collect(),
);
}
let required = match result_mode {
ResultMode::Quorum(k) => k as usize,
_ => 1,
};
let mut viable_ranked: Vec<(Duration, u64)> = branches
.iter()
.filter_map(|branch| match classify_branch(branch, budget) {
BranchViability::Viable(cost) => Some((cost, branch.provider)),
_ => None,
})
.collect();
viable_ranked.sort();
let viable: Vec<u64> = viable_ranked.into_iter().map(|(_, id)| id).collect();
let explicit_not_ready = branches
.iter()
.filter(|branch| classify_branch(branch, budget) == BranchViability::NonViable)
.count();
let complete = search_complete && population_is_boundable(selector);
let potential = branches.len().saturating_sub(explicit_not_ready);
let status = if viable.len() >= required {
ProjectedReadiness::Ready
} else if complete && potential < required {
ProjectedReadiness::NotReady
} else {
ProjectedReadiness::Unknown
};
let supporting = match result_mode {
ResultMode::TopK(k) => viable.into_iter().take(k as usize).collect(),
_ => viable,
};
AggregateView::Scalar { status, supporting }
}
#[cfg(test)]
mod tests {
use std::time::Instant;
use super::super::continuity::AttestedStatus;
use super::super::delivery::{Attestation, SensingConsumer, SensingRelay};
use super::super::identity::{
CanonicalConstraints, CapabilityId, DisclosureClass, InterestRegistration, InterestSpec,
ProviderInterestKey, ProviderObservationKey, TagMatch, WorkLatencyEnvelope,
};
use super::super::incarnation::Incarnation;
use super::super::table::DownstreamId;
use super::*;
fn root() -> AudienceScopeCommitment {
AudienceScopeCommitment::from_bytes([0xAA; 32])
}
fn ms(v: u64) -> Duration {
Duration::from_millis(v)
}
fn provider(id: u64, route_ms: u64) -> CandidateProvider {
CandidateProvider {
node_id: id,
capability_generation: 1,
authorized: true,
reachable: true,
route_estimate: ms(route_ms),
tags: Vec::new(),
groups: Vec::new(),
}
}
fn view(provider: u64, projection: ProjectedReadiness) -> BranchView {
BranchView {
provider,
projection,
estimated_start: Some(ms(50)),
route_estimate: ms(10),
}
}
fn costed(
provider: u64,
projection: ProjectedReadiness,
route_ms: u64,
start_ms: Option<u64>,
) -> BranchView {
BranchView {
provider,
projection,
estimated_start: start_ms.map(ms),
route_estimate: ms(route_ms),
}
}
#[test]
fn topk_ranks_viable_branches_by_local_economics() {
use ProjectedReadiness::Ready;
let branches = [
costed(7, Ready, 20, Some(30)),
costed(5, Ready, 15, None),
costed(9, Ready, 60, Some(40)),
costed(3, Ready, 10, Some(5)),
];
let view = project_aggregate(
&ProviderSelector::AnyAuthorized,
ResultMode::TopK(2),
&ConsumerLatencyBudget::default(),
&branches,
false,
);
match view {
AggregateView::Scalar { status, supporting } => {
assert_eq!(status, Ready);
assert_eq!(
supporting,
vec![3, 5],
"the two cheapest by local economics, id-tie-broken",
);
}
other => panic!("TopK aggregates to Scalar, got {other:?}"),
}
let view = project_aggregate(
&ProviderSelector::AnyAuthorized,
ResultMode::Any,
&ConsumerLatencyBudget::default(),
&branches,
false,
);
match view {
AggregateView::Scalar { supporting, .. } => {
assert_eq!(supporting, vec![3, 5, 7, 9], "fully ranked");
}
other => panic!("Any aggregates to Scalar, got {other:?}"),
}
}
#[test]
fn each_mode_budget_gates_ready_projections() {
use ProjectedReadiness::{NotReady, Ready, Unknown};
let budget = ConsumerLatencyBudget {
end_to_end_within: Some(ms(50)),
};
let branches = [
costed(1, Ready, 10, Some(10)), costed(2, Ready, 100, Some(10)), costed(3, NotReady, 10, None),
];
let view = project_aggregate(
&ProviderSelector::nodes(vec![1, 2, 3]),
ResultMode::Each,
&budget,
&branches,
true,
);
assert_eq!(
view,
AggregateView::PerProvider(vec![(1, Ready), (2, Unknown), (3, NotReady)]),
"over-budget Ready is locally non-viable, never Ready",
);
}
#[test]
fn any_mode_exploration_is_bounded_by_policy() {
let snapshot: Vec<CandidateProvider> =
(1..=10).map(|id| provider(id, 100 - id * 5)).collect();
let resolved = resolve_candidates(
&ProviderSelector::AnyAuthorized,
ResultMode::Any,
&snapshot,
&root(),
&CandidatePolicy::default(),
)
.unwrap();
assert_eq!(resolved.active, vec![10]);
assert_eq!(resolved.standby, vec![9]);
}
#[test]
fn node_selector_resolves_exactly_the_named_provider() {
let snapshot = vec![provider(2, 1), provider(1, 500)];
let resolved = resolve_candidates(
&ProviderSelector::Node(1),
ResultMode::Each,
&snapshot,
&root(),
&CandidatePolicy::default(),
)
.unwrap();
assert_eq!(resolved.active, vec![1]);
assert!(resolved.standby.is_empty());
}
#[test]
fn tags_require_authorized_assertions() {
let calibrated = |by: AudienceScopeCommitment| TagAssertion {
key: "calibrated".into(),
value: "true".into(),
asserted_by: by,
};
let mut legit = provider(1, 10);
legit.tags = vec![calibrated(root())];
let mut imposter = provider(2, 1);
imposter.tags = vec![calibrated(AudienceScopeCommitment::from_bytes([0xEE; 32]))];
let selector = ProviderSelector::tags(vec![TagMatch {
key: "calibrated".into(),
value: "true".into(),
}]);
let resolved = resolve_candidates(
&selector,
ResultMode::Any,
&[legit, imposter],
&root(),
&CandidatePolicy::default(),
)
.unwrap();
assert_eq!(
resolved.active,
vec![1],
"self-asserted authority tags must not admit a candidate",
);
}
#[test]
fn broad_each_selector_is_refused_before_activation() {
let snapshot: Vec<CandidateProvider> = (1..=40).map(|id| provider(id, id)).collect();
let result = resolve_candidates(
&ProviderSelector::AnyAuthorized,
ResultMode::Each,
&snapshot,
&root(),
&CandidatePolicy::default(),
);
assert_eq!(
result,
Err(ResolutionRefusal::SelectorTooBroad {
matched: 40,
cap: 32,
}),
);
}
#[test]
fn quorum_exceeding_max_fanout_is_refused_not_silently_unsatisfiable() {
let snapshot: Vec<CandidateProvider> = (1..=8).map(|id| provider(id, id)).collect();
let policy = CandidatePolicy::default(); let refusal = resolve_candidates(
&ProviderSelector::AnyAuthorized,
ResultMode::Quorum(5),
&snapshot,
&root(),
&policy,
);
assert_eq!(
refusal,
Err(ResolutionRefusal::QuorumExceedsFanout {
required: 5,
cap: policy.maximum_fanout,
}),
);
let resolved = resolve_candidates(
&ProviderSelector::AnyAuthorized,
ResultMode::Quorum(3),
&snapshot,
&root(),
&policy,
)
.expect("a quorum within the fanout budget resolves");
assert!(
resolved.active.len() >= 3,
"the active set must be able to hold the full quorum",
);
}
#[test]
fn group_each_yields_the_unflattened_map() {
let branches = [
view(1, ProjectedReadiness::Ready),
view(2, ProjectedReadiness::NotReady),
view(3, ProjectedReadiness::Unknown),
];
let group = ProviderSelector::Group(GroupRef::from_bytes([7; 32]));
let aggregate = project_aggregate(
&group,
ResultMode::Each,
&ConsumerLatencyBudget::default(),
&branches,
true,
);
assert_eq!(
aggregate,
AggregateView::PerProvider(vec![
(1, ProjectedReadiness::Ready),
(2, ProjectedReadiness::NotReady),
(3, ProjectedReadiness::Unknown),
]),
);
}
#[test]
fn quorum_flips_on_viable_count_and_respects_completeness() {
let quorum = ResultMode::Quorum(2);
let nodes = ProviderSelector::nodes(vec![1, 2, 3]);
let budget = ConsumerLatencyBudget::default();
let one_ready = [
view(1, ProjectedReadiness::Ready),
view(2, ProjectedReadiness::NotReady),
view(3, ProjectedReadiness::Unknown),
];
let aggregate = project_aggregate(&nodes, quorum, &budget, &one_ready, true);
assert_eq!(
aggregate,
AggregateView::Scalar {
status: ProjectedReadiness::Unknown,
supporting: vec![1],
},
);
let two_ready = [
view(1, ProjectedReadiness::Ready),
view(2, ProjectedReadiness::Ready),
view(3, ProjectedReadiness::NotReady),
];
let aggregate = project_aggregate(&nodes, quorum, &budget, &two_ready, true);
assert_eq!(
aggregate,
AggregateView::Scalar {
status: ProjectedReadiness::Ready,
supporting: vec![1, 2],
},
);
let starved = [
view(1, ProjectedReadiness::Ready),
view(2, ProjectedReadiness::NotReady),
view(3, ProjectedReadiness::NotReady),
];
let aggregate = project_aggregate(&nodes, quorum, &budget, &starved, true);
assert_eq!(
aggregate,
AggregateView::Scalar {
status: ProjectedReadiness::NotReady,
supporting: vec![1],
},
);
let aggregate = project_aggregate(&nodes, quorum, &budget, &starved, false);
assert_eq!(
aggregate,
AggregateView::Scalar {
status: ProjectedReadiness::Unknown,
supporting: vec![1],
},
);
}
#[test]
fn open_world_selectors_never_project_not_ready() {
let branches = [
view(1, ProjectedReadiness::NotReady),
view(2, ProjectedReadiness::NotReady),
];
let budget = ConsumerLatencyBudget::default();
for selector in [
ProviderSelector::AnyAuthorized,
ProviderSelector::tags(vec![TagMatch {
key: "site".into(),
value: "factory-7".into(),
}]),
] {
let aggregate = project_aggregate(&selector, ResultMode::Any, &budget, &branches, true);
assert_eq!(
aggregate,
AggregateView::Scalar {
status: ProjectedReadiness::Unknown,
supporting: vec![],
},
"open-world population projected NotReady",
);
}
let bounded = project_aggregate(
&ProviderSelector::nodes(vec![1, 2]),
ResultMode::Any,
&budget,
&branches,
true,
);
assert_eq!(
bounded,
AggregateView::Scalar {
status: ProjectedReadiness::NotReady,
supporting: vec![],
},
);
}
#[test]
fn flagship_coalescing_surfaces_and_the_honest_limitation() {
let any_color_a4_printer = || InterestSpec {
capability_id: CapabilityId::new("print.document"),
constraints: CanonicalConstraints::from_entries([
("color", "true"),
("duplex", "true"),
("media", "a4"),
])
.unwrap(),
work_latency: WorkLatencyEnvelope::start_within(Duration::from_secs(5)),
providers: ProviderSelector::AnyAuthorized,
result_mode: ResultMode::Any,
disclosure_class: DisclosureClass::Owner,
audience: root(),
};
let hermes = InterestRegistration {
spec: any_color_a4_printer(),
requested_sample_interval: ms(100),
soft_state_ttl: Duration::from_secs(30),
consumer_budget: ConsumerLatencyBudget::default(),
};
let desktop_ui = InterestRegistration {
spec: any_color_a4_printer(),
requested_sample_interval: Duration::from_secs(1),
soft_state_ttl: Duration::from_secs(300),
consumer_budget: ConsumerLatencyBudget {
end_to_end_within: Some(Duration::from_secs(2)),
},
};
let interest = hermes.spec.key();
assert_eq!(interest, desktop_ui.spec.key());
let p1 = 1u64;
let shared_snapshot = vec![provider(p1, 10), provider(2, 40)];
for _node in ["A", "C"] {
let resolved = resolve_candidates(
&ProviderSelector::AnyAuthorized,
ResultMode::Any,
&shared_snapshot,
&root(),
&CandidatePolicy::default(),
)
.unwrap();
assert_eq!(resolved.active, vec![p1]);
}
let t0 = Instant::now();
let branch = ProviderInterestKey::new(interest.clone(), p1);
let (a, c) = (DownstreamId::Peer(0xA), DownstreamId::Peer(0xC));
let mut relay = SensingRelay::new(3, 512);
let mut consumer_a = SensingConsumer::new(3);
let mut consumer_c = SensingConsumer::new(3);
consumer_a.register_interest(&branch, ms(100), t0);
consumer_c.register_interest(&branch, ms(100), t0);
relay.register_downstream(&branch, a, ms(100), Duration::from_secs(30), root(), t0);
relay.register_downstream(&branch, c, ms(100), Duration::from_secs(30), root(), t0);
assert_eq!(
relay.table.len(),
1,
"equivalent demand must share one entry"
);
let proof = Attestation::new(
ProviderObservationKey::new(interest.clone(), p1, 42),
Incarnation::new(1),
AttestedStatus::Ready,
Some(ms(800)),
1,
ms(100),
);
let out = relay.on_attestation(t0 + ms(100), &proof, true);
let to_a = out
.iter()
.find(|d| d.to == a)
.expect("A receives the proof");
let to_c = out
.iter()
.find(|d| d.to == c)
.expect("C receives the proof");
assert_eq!(
to_a.attestation.fingerprint, to_c.attestation.fingerprint,
"one provider stream serves both consumers with identical signed bytes",
);
consumer_a.on_delivery(t0 + ms(100), to_a);
consumer_c.on_delivery(t0 + ms(100), to_c);
assert_eq!(consumer_a.projected(&branch), ProjectedReadiness::Ready);
assert_eq!(consumer_c.projected(&branch), ProjectedReadiness::Ready);
let snapshot_a = vec![provider(1, 10), provider(2, 40)];
let snapshot_c = vec![provider(1, 40), provider(2, 10)];
let resolve = |snapshot: &[CandidateProvider]| {
resolve_candidates(
&ProviderSelector::AnyAuthorized,
ResultMode::Any,
snapshot,
&root(),
&CandidatePolicy::default(),
)
.unwrap()
.active
};
assert_eq!(resolve(&snapshot_a), vec![1]);
assert_eq!(resolve(&snapshot_c), vec![2]);
let mut divergent_relay = SensingRelay::new(3, 512);
divergent_relay.register_downstream(
&ProviderInterestKey::new(interest.clone(), 1),
a,
ms(100),
Duration::from_secs(30),
root(),
t0,
);
divergent_relay.register_downstream(
&ProviderInterestKey::new(interest, 2),
c,
ms(100),
Duration::from_secs(30),
root(),
t0,
);
assert_eq!(
divergent_relay.table.len(),
2,
"divergent provider resolution does not merge — the honest v1 cost",
);
}
#[test]
fn budget_makes_viability_consumer_relative() {
let proof = |route_ms: u64| BranchView {
provider: 1,
projection: ProjectedReadiness::Ready,
estimated_start: Some(ms(300)),
route_estimate: ms(route_ms),
};
let budget = ConsumerLatencyBudget {
end_to_end_within: Some(ms(500)),
};
let selector = ProviderSelector::AnyAuthorized;
let near = project_aggregate(&selector, ResultMode::Any, &budget, &[proof(150)], false);
assert_eq!(
near,
AggregateView::Scalar {
status: ProjectedReadiness::Ready,
supporting: vec![1],
},
);
let far = project_aggregate(&selector, ResultMode::Any, &budget, &[proof(250)], false);
assert_eq!(
far,
AggregateView::Scalar {
status: ProjectedReadiness::Unknown,
supporting: vec![],
},
"an over-budget Ready is not viable — and not NotReady either",
);
}
}